use crate::transports::ice::IceCandidate;
use anyhow::{Result, anyhow};
use igd::AddPortError;
use igd::PortMappingProtocol;
use igd::aio::Gateway;
use std::collections::HashMap;
use std::net::{IpAddr, Ipv4Addr, SocketAddr, SocketAddrV4};
use std::sync::Arc;
use std::time::Duration;
use tokio::sync::Mutex;
use tokio::time::timeout;
use tracing::{debug, trace, warn};
pub const DEFAULT_LEASE_DURATION: u32 = 3600;
pub const MIN_LEASE_DURATION: u32 = 300;
pub const MAX_LEASE_DURATION: u32 = 86400;
pub const DEFAULT_UPNP_DISCOVERY_TIMEOUT: Duration = Duration::from_secs(2);
#[derive(Debug, Clone)]
pub struct PortMapping {
pub external_port: u16,
pub internal_addr: SocketAddr,
pub lease_duration: u32,
pub description: String,
pub created_at: std::time::Instant,
}
impl PortMapping {
pub fn is_expired_or_stale(&self) -> bool {
let elapsed = self.created_at.elapsed().as_secs() as u32;
elapsed + 60 >= self.lease_duration
}
pub fn remaining_lifetime(&self) -> u32 {
let elapsed = self.created_at.elapsed().as_secs() as u32;
self.lease_duration.saturating_sub(elapsed)
}
}
#[derive(Debug, Clone)]
pub struct UpnpPortMapper {
gateway: Option<Gateway>,
mappings: Arc<Mutex<HashMap<u16, PortMapping>>>,
pub local_addr: SocketAddr,
pub default_lease_duration: u32,
enabled: bool,
}
impl UpnpPortMapper {
pub fn new(local_addr: SocketAddr) -> Self {
Self {
gateway: None,
mappings: Arc::new(Mutex::new(HashMap::new())),
local_addr,
default_lease_duration: DEFAULT_LEASE_DURATION,
enabled: true,
}
}
pub fn with_lease_duration(local_addr: SocketAddr, lease_duration: u32) -> Self {
let lease_duration = lease_duration.clamp(MIN_LEASE_DURATION, MAX_LEASE_DURATION);
Self {
gateway: None,
mappings: Arc::new(Mutex::new(HashMap::new())),
local_addr,
default_lease_duration: lease_duration,
enabled: true,
}
}
pub fn disable(&mut self) {
self.enabled = false;
self.gateway = None;
}
pub fn enable(&mut self) {
self.enabled = true;
}
#[cfg(test)]
pub fn set_gateway_for_test(&mut self, gateway: Gateway) {
self.gateway = Some(gateway);
}
#[cfg(test)]
pub async fn insert_mapping_for_test(&self, external_port: u16, age: Duration) {
let mut mappings = self.mappings.lock().await;
mappings.insert(
external_port,
PortMapping {
external_port,
internal_addr: self.local_addr,
lease_duration: self.default_lease_duration,
description: format!("rustrtc-{}", self.local_addr.port()),
created_at: std::time::Instant::now().checked_sub(age).unwrap(),
},
);
}
pub fn is_enabled(&self) -> bool {
self.enabled
}
pub async fn discover(&mut self) -> Result<()> {
self.discover_with_timeout(DEFAULT_UPNP_DISCOVERY_TIMEOUT)
.await
}
pub async fn discover_with_timeout(&mut self, timeout_duration: Duration) -> Result<()> {
if !self.enabled {
return Err(anyhow!("UPnP is disabled"));
}
if self.local_addr.ip().is_loopback() {
return Err(anyhow!("Cannot map loopback address"));
}
trace!(
"Starting UPnP gateway discovery (timeout: {:?})",
timeout_duration
);
let gateway = timeout(
timeout_duration,
igd::aio::search_gateway(Default::default()),
)
.await
.map_err(|_| {
anyhow!(
"UPnP gateway discovery timed out after {:?}",
timeout_duration
)
})?
.map_err(|e| anyhow!("UPnP gateway discovery failed: {}", e))?;
debug!("Found UPnP gateway");
self.gateway = Some(gateway);
Ok(())
}
pub fn has_gateway(&self) -> bool {
self.gateway.is_some()
}
pub async fn get_external_ip(&self) -> Result<Ipv4Addr> {
let gateway = self
.gateway
.as_ref()
.ok_or_else(|| anyhow!("No UPnP gateway available"))?;
let ip = gateway
.get_external_ip()
.await
.map_err(|e| anyhow!("Failed to get external IP: {}", e))?;
Ok(ip)
}
pub async fn add_mapping(&self, external_port: u16) -> Result<SocketAddr> {
if !self.enabled {
return Err(anyhow!("UPnP is disabled"));
}
let gateway = self
.gateway
.as_ref()
.ok_or_else(|| anyhow!("No UPnP gateway available, call discover() first"))?;
let external_ip = self.get_external_ip().await?;
let requested_port = if external_port == 0 {
self.local_addr.port()
} else {
external_port
};
let description = format!("rustrtc-{}", self.local_addr.port());
let local_ip = match self.local_addr.ip() {
IpAddr::V4(ip) => ip,
IpAddr::V6(_) => return Err(anyhow!("IPv6 not supported for UPnP IGD")),
};
trace!(
"Adding UPnP port mapping: {}:{} -> {}:{}",
external_ip,
requested_port,
local_ip,
self.local_addr.port()
);
let internal_sock_addr = SocketAddrV4::new(local_ip, self.local_addr.port());
match gateway
.add_port(
PortMappingProtocol::UDP,
requested_port,
internal_sock_addr,
self.default_lease_duration,
&description,
)
.await
{
Ok(()) => {
let external_addr = SocketAddr::new(IpAddr::V4(external_ip), requested_port);
let mapping = PortMapping {
external_port: requested_port,
internal_addr: self.local_addr,
lease_duration: self.default_lease_duration,
description,
created_at: std::time::Instant::now(),
};
self.mappings.lock().await.insert(requested_port, mapping);
debug!(
"UPnP port mapping added: {} -> {}",
external_addr, self.local_addr
);
Ok(external_addr)
}
Err(e) => {
if external_port != 0 && requested_port != 0 {
warn!(
"Port {} is taken, trying random port: {}",
requested_port, e
);
self.add_mapping_random_port(gateway, external_ip, local_ip)
.await
} else {
Err(anyhow!("Failed to add UPnP port mapping: {}", e))
}
}
}
}
async fn add_mapping_random_port(
&self,
gateway: &Gateway,
external_ip: Ipv4Addr,
local_ip: Ipv4Addr,
) -> Result<SocketAddr> {
for port in 10000..=65535u16 {
let description = format!("rustrtc-{}", self.local_addr.port());
let internal_sock_addr = SocketAddrV4::new(local_ip, self.local_addr.port());
match gateway
.add_port(
PortMappingProtocol::UDP,
port,
internal_sock_addr,
self.default_lease_duration,
&description,
)
.await
{
Ok(()) => {
let external_addr = SocketAddr::new(IpAddr::V4(external_ip), port);
let mapping = PortMapping {
external_port: port,
internal_addr: self.local_addr,
lease_duration: self.default_lease_duration,
description,
created_at: std::time::Instant::now(),
};
self.mappings.lock().await.insert(port, mapping);
debug!(
"UPnP port mapping added (random port): {} -> {}",
external_addr, self.local_addr
);
return Ok(external_addr);
}
Err(_) => continue,
}
}
Err(anyhow!("Failed to find available port for UPnP mapping"))
}
pub async fn remove_mapping(&self, external_port: u16) -> Result<()> {
let gateway = match &self.gateway {
Some(g) => g,
None => {
self.mappings.lock().await.remove(&external_port);
return Ok(());
}
};
gateway
.remove_port(PortMappingProtocol::UDP, external_port)
.await
.map_err(|e| anyhow!("Failed to remove UPnP mapping: {}", e))?;
self.mappings.lock().await.remove(&external_port);
debug!("UPnP port mapping removed: {}", external_port);
Ok(())
}
pub async fn cleanup(&self) -> Result<()> {
let mappings = self.mappings.lock().await.clone();
let mut last_error = None;
for (port, _) in mappings {
if let Err(e) = self.remove_mapping(port).await {
warn!("Failed to remove UPnP mapping for port {}: {}", port, e);
last_error = Some(e);
}
}
match last_error {
Some(e) => Err(e),
None => Ok(()),
}
}
pub async fn mapping_count(&self) -> usize {
self.mappings.lock().await.len()
}
pub async fn has_mapping(&self, external_port: u16) -> bool {
self.mappings.lock().await.contains_key(&external_port)
}
pub async fn get_mappings(&self) -> HashMap<u16, PortMapping> {
self.mappings.lock().await.clone()
}
pub async fn refresh_mapping(&self, external_port: u16) -> Result<bool> {
let gateway = match &self.gateway {
Some(g) => g,
None => return Ok(false), };
let mapping = {
let mappings = self.mappings.lock().await;
match mappings.get(&external_port) {
Some(m) => m.clone(),
None => return Ok(false),
}
};
let local_ip = match mapping.internal_addr.ip() {
IpAddr::V4(ip) => ip,
IpAddr::V6(_) => return Err(anyhow!("IPv6 not supported for UPnP IGD")),
};
let internal_sock_addr = SocketAddrV4::new(local_ip, mapping.internal_addr.port());
let refreshed = match gateway
.add_port(
PortMappingProtocol::UDP,
external_port,
internal_sock_addr,
self.default_lease_duration,
&mapping.description,
)
.await
{
Ok(()) => true,
Err(AddPortError::PortInUse) => {
let _ = self.remove_mapping(external_port).await;
gateway
.add_port(
PortMappingProtocol::UDP,
external_port,
internal_sock_addr,
self.default_lease_duration,
&mapping.description,
)
.await
.map_err(|e| anyhow!("Failed to re-add UPnP mapping after conflict: {}", e))?;
true
}
Err(AddPortError::OnlyPermanentLeasesSupported) => {
true
}
Err(e) => {
return Err(anyhow!("Failed to refresh UPnP mapping: {}", e));
}
};
if refreshed {
{
let mut mappings = self.mappings.lock().await;
if let Some(m) = mappings.get_mut(&external_port) {
m.created_at = std::time::Instant::now();
}
}
debug!("Refreshed UPnP mapping for port {}", external_port);
}
Ok(refreshed)
}
pub async fn renew_mapping(&self, external_port: u16) -> Result<bool> {
let needs_renewal = {
let mappings = self.mappings.lock().await;
match mappings.get(&external_port) {
Some(mapping) if mapping.is_expired_or_stale() => true,
Some(_) => return Ok(false), None => return Ok(false), }
};
if !needs_renewal {
return Ok(false);
}
self.refresh_mapping(external_port).await
}
pub async fn renew_all_stale(&self) -> Result<usize> {
let ports_to_renew: Vec<u16> = {
let mappings = self.mappings.lock().await;
mappings
.values()
.filter(|m| m.is_expired_or_stale())
.map(|m| m.external_port)
.collect()
};
let mut renewed = 0;
for port in ports_to_renew {
match self.renew_mapping(port).await {
Ok(true) => renewed += 1,
Ok(false) => {}
Err(e) => {
warn!("Failed to renew UPnP mapping for port {}: {}", port, e);
}
}
}
Ok(renewed)
}
pub async fn create_candidate(&self) -> Result<IceCandidate> {
let mappings = self.mappings.lock().await;
let mapping = mappings
.values()
.next()
.ok_or_else(|| anyhow!("No UPnP mappings available"))?;
let external_addr = SocketAddr::new(
IpAddr::V4(self.get_external_ip().await?),
mapping.external_port,
);
Ok(IceCandidate::server_reflexive(
mapping.internal_addr,
external_addr,
1, ))
}
}
pub async fn try_create_upnp_candidate(local_addr: SocketAddr) -> Option<IceCandidate> {
if local_addr.ip().is_loopback() {
return None;
}
let mut mapper = UpnpPortMapper::new(local_addr);
if let Err(e) = mapper.discover().await {
trace!("UPnP discovery failed for {}: {}", local_addr, e);
return None;
}
let external_addr = match mapper.add_mapping(0).await {
Ok(addr) => addr,
Err(e) => {
debug!("UPnP mapping failed for {}: {}", local_addr, e);
return None;
}
};
let candidate = IceCandidate::server_reflexive(local_addr, external_addr, 1);
debug!(
"Created UPnP candidate: {} -> {}",
local_addr, external_addr
);
Some(candidate)
}
#[cfg(test)]
pub(crate) mod test_mock_igd {
use super::*;
use igd::aio::Gateway;
use std::collections::VecDeque;
use std::net::SocketAddrV4;
use std::sync::atomic::{AtomicUsize, Ordering};
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::TcpListener;
use tokio::sync::Mutex as TokioMutex;
#[derive(Debug, Clone)]
pub enum MockAddResponse {
Ok,
PortInUse,
PermanentOnly,
Unauthorized,
}
#[derive(Debug, Clone)]
pub struct MockIgd {
addr: SocketAddrV4,
pub add_calls: Arc<AtomicUsize>,
pub delete_calls: Arc<AtomicUsize>,
add_queue: Arc<TokioMutex<VecDeque<MockAddResponse>>>,
}
impl MockIgd {
pub async fn start() -> Self {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let addr_v4 = match addr {
SocketAddr::V4(v4) => v4,
_ => unreachable!("bound to IPv4 loopback"),
};
let this = Self {
addr: addr_v4,
add_calls: Arc::new(AtomicUsize::new(0)),
delete_calls: Arc::new(AtomicUsize::new(0)),
add_queue: Arc::new(TokioMutex::new(VecDeque::new())),
};
let server = this.clone();
tokio::spawn(async move {
loop {
let Ok((mut sock, _)) = listener.accept().await else {
break;
};
let server = server.clone();
tokio::spawn(async move {
let mut buf = [0u8; 8192];
let n = match sock.read(&mut buf).await {
Ok(n) => n,
Err(_) => return,
};
let req = String::from_utf8_lossy(&buf[..n]);
let soap_action = req
.lines()
.find_map(|l| {
let lower = l.to_ascii_lowercase();
lower
.starts_with("soapaction:")
.then(|| l["soapaction:".len()..].trim().to_string())
})
.unwrap_or_default();
let action = soap_action
.rsplit('#')
.next()
.unwrap_or("")
.trim_matches('"')
.to_string();
let body = server.handle_action(&action).await;
let response = format!(
"HTTP/1.1 200 OK\r\nContent-Type: text/xml\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
body.len(),
body
);
let _ = sock.write_all(response.as_bytes()).await;
});
}
});
this
}
pub async fn handle_action(&self, action: &str) -> String {
if action.contains("GetExternalIPAddress") {
return ok_response(
"GetExternalIPAddressResponse",
"<NewExternalIPAddress>203.0.113.1</NewExternalIPAddress>",
);
}
if action.contains("AddPortMapping") {
let _ = self.add_calls.fetch_add(1, Ordering::SeqCst);
let resp = {
let mut q = self.add_queue.lock().await;
q.pop_front().unwrap_or(MockAddResponse::Ok)
};
return match resp {
MockAddResponse::Ok => ok_response("AddPortMappingResponse", ""),
MockAddResponse::PortInUse => error_response(718, "ConflictInMappingEntry"),
MockAddResponse::PermanentOnly => {
error_response(725, "OnlyPermanentLeasesSupported")
}
MockAddResponse::Unauthorized => error_response(606, "ActionNotAuthorized"),
};
}
if action.contains("DeletePortMapping") {
let _ = self.delete_calls.fetch_add(1, Ordering::SeqCst);
return ok_response("DeletePortMappingResponse", "");
}
error_response(401, "InvalidAction")
}
pub fn gateway(&self) -> Gateway {
let mut schema = HashMap::new();
schema.insert(
"AddPortMapping".to_string(),
vec![
"NewEnabled".to_string(),
"NewExternalPort".to_string(),
"NewInternalClient".to_string(),
"NewInternalPort".to_string(),
"NewLeaseDuration".to_string(),
"NewPortMappingDescription".to_string(),
"NewProtocol".to_string(),
"NewRemoteHost".to_string(),
],
);
schema.insert(
"DeletePortMapping".to_string(),
vec![
"NewExternalPort".to_string(),
"NewProtocol".to_string(),
"NewRemoteHost".to_string(),
],
);
Gateway {
addr: self.addr,
root_url: format!("http://{}/", self.addr),
control_url: "/upnp/control".to_string(),
control_schema_url: "/upnp/control/scpd.xml".to_string(),
control_schema: schema,
}
}
pub async fn enqueue_add(&self, resp: MockAddResponse) {
self.add_queue.lock().await.push_back(resp);
}
}
pub fn ok_response(name: &str, inner: &str) -> String {
soap_envelope(&format!(
r#"<u:{name} xmlns:u="urn:schemas-upnp-org:service:WANIPConnection:1">{inner}</u:{name}>"#
))
}
pub fn error_response(code: u16, desc: &str) -> String {
soap_envelope(&format!(
r#"<s:Fault><faultcode>s:Client</faultcode><faultstring>UPnPError</faultstring><detail>
<UPnPError><errorCode>{code}</errorCode><errorDescription>{desc}</errorDescription></UPnPError>
</detail></s:Fault>"#
))
}
pub fn soap_envelope(body: &str) -> String {
format!(
r#"<?xml version="1.0"?>
<s:Envelope xmlns:s="http://schemas.xmlsoap.org/soap/envelope/" s:encodingStyle="http://schemas.xmlsoap.org/soap/encoding/">
<s:Body>{body}</s:Body>
</s:Envelope>"#
)
}
pub const STALE_AGE: Duration = Duration::from_secs(3540); pub const FRESH_AGE: Duration = Duration::from_secs(60);
}
#[cfg(test)]
mod tests {
use super::test_mock_igd::{FRESH_AGE, MockAddResponse, MockIgd, STALE_AGE};
use super::*;
use std::sync::atomic::Ordering;
#[test]
fn test_port_mapping_expiry() {
let mapping = PortMapping {
external_port: 12345,
internal_addr: "192.168.1.100:5000".parse().unwrap(),
lease_duration: 70,
description: "test".to_string(),
created_at: std::time::Instant::now(),
};
assert!(!mapping.is_expired_or_stale());
let remaining = mapping.remaining_lifetime();
assert!((69..=70).contains(&remaining));
}
#[test]
fn test_port_mapping_remaining_lifetime() {
let mapping = PortMapping {
external_port: 12345,
internal_addr: "192.168.1.100:5000".parse().unwrap(),
lease_duration: 60,
description: "test".to_string(),
created_at: std::time::Instant::now(),
};
let remaining = mapping.remaining_lifetime();
assert!(remaining > 55 && remaining <= 60);
std::thread::sleep(std::time::Duration::from_millis(100));
let new_remaining = mapping.remaining_lifetime();
assert!(
new_remaining <= remaining,
"remaining={}, new_remaining={}",
remaining,
new_remaining
);
}
#[test]
fn test_upnp_mapper_creation() {
let addr: SocketAddr = "192.168.1.100:5000".parse().unwrap();
let mapper = UpnpPortMapper::new(addr);
assert!(mapper.is_enabled());
assert!(!mapper.has_gateway());
assert_eq!(mapper.local_addr, addr);
}
#[test]
fn test_upnp_mapper_disable_enable() {
let addr: SocketAddr = "192.168.1.100:5000".parse().unwrap();
let mut mapper = UpnpPortMapper::new(addr);
assert!(mapper.is_enabled());
mapper.disable();
assert!(!mapper.is_enabled());
assert!(mapper.gateway.is_none());
mapper.enable();
assert!(mapper.is_enabled());
}
#[test]
fn test_upnp_mapper_custom_lease() {
let addr: SocketAddr = "192.168.1.100:5000".parse().unwrap();
let mapper = UpnpPortMapper::with_lease_duration(addr, 100);
assert_eq!(mapper.default_lease_duration, MIN_LEASE_DURATION);
let mapper = UpnpPortMapper::with_lease_duration(addr, 100000);
assert_eq!(mapper.default_lease_duration, MAX_LEASE_DURATION);
let mapper = UpnpPortMapper::with_lease_duration(addr, 1800);
assert_eq!(mapper.default_lease_duration, 1800);
}
#[tokio::test]
async fn test_upnp_mapper_loopback_rejection() {
let addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
let mut mapper = UpnpPortMapper::new(addr);
let result = mapper.discover().await;
assert!(result.is_err());
assert!(result.unwrap_err().to_string().contains("loopback"));
}
#[tokio::test]
async fn test_upnp_mapper_disabled() {
let addr: SocketAddr = "192.168.1.100:5000".parse().unwrap();
let mut mapper = UpnpPortMapper::new(addr);
mapper.disable();
let result = mapper.discover().await;
assert!(result.is_err());
assert!(result.unwrap_err().to_string().contains("disabled"));
}
#[tokio::test]
async fn test_upnp_mapper_no_gateway() {
let addr: SocketAddr = "192.168.1.100:5000".parse().unwrap();
let mapper = UpnpPortMapper::new(addr);
let result = mapper.add_mapping(12345).await;
assert!(result.is_err());
assert!(result.unwrap_err().to_string().contains("No UPnP gateway"));
}
#[test]
fn test_try_create_upnp_candidate_loopback() {
let rt = tokio::runtime::Runtime::new().unwrap();
let result = rt.block_on(async {
let addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
try_create_upnp_candidate(addr).await
});
assert!(result.is_none());
}
#[tokio::test]
async fn test_upnp_mapper_clone() {
let addr: SocketAddr = "192.168.1.100:5000".parse().unwrap();
let mapper = UpnpPortMapper::new(addr);
let cloned = mapper.clone();
assert_eq!(cloned.local_addr, addr);
assert!(cloned.is_enabled());
assert!(!cloned.has_gateway());
}
#[test]
#[allow(clippy::assertions_on_constants)] fn test_mapping_constants() {
assert!(MIN_LEASE_DURATION > 0);
assert!(MAX_LEASE_DURATION > MIN_LEASE_DURATION);
assert!(DEFAULT_LEASE_DURATION >= MIN_LEASE_DURATION);
assert!(DEFAULT_LEASE_DURATION <= MAX_LEASE_DURATION);
}
#[tokio::test]
async fn test_refresh_mapping_reissues_add_without_delete() {
let mock = MockIgd::start().await;
let mut mapper = UpnpPortMapper::new("192.168.1.100:5000".parse().unwrap());
mapper.set_gateway_for_test(mock.gateway());
mapper.insert_mapping_for_test(10001, STALE_AGE).await;
let renewed = mapper.renew_mapping(10001).await.unwrap();
assert!(renewed, "stale mapping must be renewed");
assert_eq!(mock.add_calls.load(Ordering::SeqCst), 1);
assert_eq!(mock.delete_calls.load(Ordering::SeqCst), 0);
assert!(
!mapper.mappings.lock().await[&10001].is_expired_or_stale(),
"created_at must be reset after refresh"
);
}
#[tokio::test]
async fn test_refresh_mapping_router_lost_mapping_readds() {
let mock = MockIgd::start().await;
let mut mapper = UpnpPortMapper::new("192.168.1.100:5000".parse().unwrap());
mapper.set_gateway_for_test(mock.gateway());
mapper.insert_mapping_for_test(10002, STALE_AGE).await;
let renewed = mapper.renew_mapping(10002).await.unwrap();
assert!(renewed);
assert_eq!(mock.add_calls.load(Ordering::SeqCst), 1);
assert_eq!(mock.delete_calls.load(Ordering::SeqCst), 0);
}
#[tokio::test]
async fn test_refresh_mapping_port_in_use_falls_back_to_delete_then_add() {
let mock = MockIgd::start().await;
mock.enqueue_add(MockAddResponse::PortInUse).await; let mut mapper = UpnpPortMapper::new("192.168.1.100:5000".parse().unwrap());
mapper.set_gateway_for_test(mock.gateway());
mapper.insert_mapping_for_test(10003, STALE_AGE).await;
let renewed = mapper.renew_mapping(10003).await.unwrap();
assert!(renewed, "PortInUse fallback should still renew");
assert_eq!(
mock.add_calls.load(Ordering::SeqCst),
2,
"expect delete-then-add: two AddPortMapping calls"
);
assert_eq!(
mock.delete_calls.load(Ordering::SeqCst),
1,
"expect one DeletePortMapping before re-add"
);
}
#[tokio::test]
async fn test_refresh_mapping_never_switches_port_on_conflict() {
let mock = MockIgd::start().await;
mock.enqueue_add(MockAddResponse::PortInUse).await;
let mut mapper = UpnpPortMapper::new("192.168.1.100:5000".parse().unwrap());
mapper.set_gateway_for_test(mock.gateway());
mapper.insert_mapping_for_test(10004, STALE_AGE).await;
let renewed = mapper.renew_mapping(10004).await.unwrap();
assert!(renewed);
assert_eq!(mock.add_calls.load(Ordering::SeqCst), 2);
assert_eq!(mock.delete_calls.load(Ordering::SeqCst), 1);
}
#[tokio::test]
async fn test_refresh_mapping_only_permanent_lease_supported() {
let mock = MockIgd::start().await;
mock.enqueue_add(MockAddResponse::PermanentOnly).await;
let mut mapper = UpnpPortMapper::new("192.168.1.100:5000".parse().unwrap());
mapper.set_gateway_for_test(mock.gateway());
mapper.insert_mapping_for_test(10005, STALE_AGE).await;
let renewed = mapper.renew_mapping(10005).await.unwrap();
assert!(renewed);
assert_eq!(mock.add_calls.load(Ordering::SeqCst), 1);
assert_eq!(mock.delete_calls.load(Ordering::SeqCst), 0);
assert!(!mapper.mappings.lock().await[&10005].is_expired_or_stale());
}
#[tokio::test]
async fn test_refresh_mapping_unauthorized_returns_error() {
let mock = MockIgd::start().await;
mock.enqueue_add(MockAddResponse::Unauthorized).await;
let mut mapper = UpnpPortMapper::new("192.168.1.100:5000".parse().unwrap());
mapper.set_gateway_for_test(mock.gateway());
mapper.insert_mapping_for_test(10006, STALE_AGE).await;
let result = mapper.renew_mapping(10006).await;
assert!(result.is_err(), "non-recoverable router error must surface");
}
#[tokio::test]
async fn test_renew_mapping_missing_locally_returns_false() {
let mock = MockIgd::start().await;
let mut mapper = UpnpPortMapper::new("192.168.1.100:5000".parse().unwrap());
mapper.set_gateway_for_test(mock.gateway());
let renewed = mapper.renew_mapping(10007).await.unwrap();
assert!(!renewed);
assert_eq!(mock.add_calls.load(Ordering::SeqCst), 0);
assert_eq!(mock.delete_calls.load(Ordering::SeqCst), 0);
}
#[tokio::test]
async fn test_renew_mapping_no_gateway_returns_false() {
let mapper = UpnpPortMapper::new("192.168.1.100:5000".parse().unwrap());
mapper.insert_mapping_for_test(10008, STALE_AGE).await;
let renewed = mapper.renew_mapping(10008).await.unwrap();
assert!(!renewed, "no gateway -> nothing to refresh, must not error");
}
#[tokio::test]
async fn test_renew_all_stale_skips_fresh_mappings() {
let mock = MockIgd::start().await;
let mut mapper = UpnpPortMapper::new("192.168.1.100:5000".parse().unwrap());
mapper.set_gateway_for_test(mock.gateway());
mapper.insert_mapping_for_test(10009, FRESH_AGE).await;
let renewed = mapper.renew_all_stale().await.unwrap();
assert_eq!(renewed, 0);
assert_eq!(mock.add_calls.load(Ordering::SeqCst), 0);
}
#[tokio::test]
async fn test_renew_all_stale_continues_on_error() {
let mock = MockIgd::start().await;
mock.enqueue_add(MockAddResponse::Unauthorized).await;
let mut mapper = UpnpPortMapper::new("192.168.1.100:5000".parse().unwrap());
mapper.set_gateway_for_test(mock.gateway());
mapper.insert_mapping_for_test(10010, STALE_AGE).await;
mapper.insert_mapping_for_test(10011, STALE_AGE).await;
let renewed = mapper.renew_all_stale().await.unwrap();
assert_eq!(renewed, 1, "second mapping must still be renewed");
assert_eq!(mock.add_calls.load(Ordering::SeqCst), 2);
}
#[tokio::test]
async fn test_renew_all_stale_no_gateway_is_noop() {
let mapper = UpnpPortMapper::new("192.168.1.100:5000".parse().unwrap());
mapper.insert_mapping_for_test(10012, STALE_AGE).await;
let renewed = mapper.renew_all_stale().await.unwrap();
assert_eq!(renewed, 0, "no gateway -> no mappings renewed, no error");
}
}