Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
81 changes: 81 additions & 0 deletions src/transaction/tests/test_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -473,3 +473,84 @@ async fn test_invite_2xx_upstream_via_delivery_and_ack() -> Result<()> {
}
Ok(())
}

/// RFC 3261 §17.1.2.2: Timer F still runs after a provisional response. A
/// BYE answered with one 1xx and then nothing must end with a 408 at 64*T1,
/// on unreliable and reliable transports alike.
#[tokio::test]
async fn test_non_invite_timer_f_after_provisional() -> Result<()> {
use crate::transaction::{endpoint::EndpointOption, EndpointBuilder};
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::{TcpListener, UdpSocket};

// The request with its start line swapped for a provisional status line.
fn provisional(code: u16, request: &[u8]) -> Vec<u8> {
let text = String::from_utf8_lossy(request);
let (_, headers) = text.split_once("\r\n").unwrap();
format!("SIP/2.0 {code} Provisional\r\n{headers}").into_bytes()
}

for (code, tcp) in [(100, false), (180, false), (183, true)] {
let tl = crate::transport::TransportLayer::new(Default::default());
let udp = UdpConnection::create_connection("127.0.0.1:0".parse()?, None, None).await?;
tl.add_transport(udp.into());
let t1 = Duration::from_millis(20);
let option = EndpointOption {
t1,
t1x64: t1 * 64,
..Default::default()
};
let endpoint = EndpointBuilder::new()
.with_transport_layer(tl)
.with_option(option)
.build();

// The peer answers the BYE with one provisional, then stays silent.
let socket = UdpSocket::bind("127.0.0.1:0").await?;
let listener = TcpListener::bind("127.0.0.1:0").await?;
let uri = match tcp {
true => format!("sip:bob@{};transport=tcp", listener.local_addr()?),
false => format!("sip:bob@{}", socket.local_addr()?),
};
let peer = tokio::spawn(async move {
let mut buf = vec![0u8; 4096];
if tcp {
let (mut stream, _) = listener.accept().await?;
let mut len = 0;
while !buf[..len].windows(4).any(|w| w == b"\r\n\r\n") {
len += stream.read(&mut buf[len..]).await?;
}
stream.write_all(&provisional(code, &buf[..len])).await?;
std::future::pending::<()>().await; // keep the connection open
} else {
let (len, src) = socket.recv_from(&mut buf).await?;
socket.send_to(&provisional(code, &buf[..len]), src).await?;
}
std::future::pending::<std::io::Result<()>>().await
});

let mut bye = make_invite_request(&uri)?;
bye.method = crate::sip::Method::Bye;
bye.headers.unique_push(CSeq::new("2 BYE").into());
let key = TransactionKey::from_request(&bye, TransactionRole::Client)?;
let mut tx = Transaction::new_client(key, bye, endpoint.inner.clone(), None);
let mut codes = vec![];
let run = async {
tx.send().await?;
while let Some(SipMessage::Response(resp)) = tx.receive().await {
codes.push(resp.status_code.code());
}
Ok::<_, crate::Error>(())
};
// `receive` returning None means the transaction terminated.
let terminated = select! {
r = run => r.map(|_| true)?,
_ = endpoint.serve() => panic!("endpoint stopped"),
_ = sleep(t1 * 64 * 3) => false,
};
peer.abort();
assert_eq!(codes, [code, 408], "{code} over tcp={tcp}");
assert!(terminated, "{code} over tcp={tcp}: not terminated");
}
Ok(())
}
7 changes: 6 additions & 1 deletion src/transaction/transaction.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1161,7 +1161,12 @@ impl Transaction {
}
}
TransactionState::Proceeding => {
if let TransactionTimer::TimerC(_) = timer {
// Timer C (client INVITE), or Timer F, run as Timer B, for a
// non-INVITE client (RFC 3261 §17.1.2.2).
if matches!(timer, TransactionTimer::TimerC(_))
|| (matches!(timer, TransactionTimer::TimerB(_))
&& self.transaction_type == TransactionType::ClientNonInvite)
{
// Inform TU about timeout
let timeout_response = self.endpoint_inner.make_response(
&self.original,
Expand Down
Loading