oxedyne/fe2o3/fe2o3_net/src/ssdp.rs
53.2 KiB, 149 runs
created by r1870400018:17852, which is this file's identity for as long as the history lasts, whatever it is later renamed to
download · who wrote it · its history
| 1 | //! SSDP: the discovery half of UPnP (UPnP Device Architecture 2.0, §1). |
| 2 | //! |
| 3 | //! A television on a home network finds what it can play from by shouting into a |
| 4 | //! multicast group and listening for whoever answers. That is all SSDP is: three |
| 5 | //! HTTP-shaped messages sent over UDP to 239.255.255.250:1900, with a start line |
| 6 | //! and a block of fields and no body at all. |
| 7 | //! |
| 8 | //! - `M-SEARCH * HTTP/1.1` -- a searcher asking who is out there. |
| 9 | //! - `HTTP/1.1 200 OK` -- a unicast answer to a search, sent back to the asker. |
| 10 | //! - `NOTIFY * HTTP/1.1` -- an announcement, multicast: `ssdp:alive` when a |
| 11 | //! device appears or renews, `ssdp:byebye` when it goes. |
| 12 | //! |
| 13 | //! Every one of them names a service in two ways: a *target* (`ST` on a search |
| 14 | //! and its answer, `NT` on a notification) saying what kind of thing it is, and a |
| 15 | //! `USN` saying which particular thing it is. A responder that gets those two out |
| 16 | //! of step is discovered and then cannot be reached, which is the failure mode |
| 17 | //! that eats an afternoon. |
| 18 | //! |
| 19 | //! # What is here |
| 20 | //! |
| 21 | //! The messages, their parsing and their serialisation, and a responder that |
| 22 | //! binds the group and reads and writes them. The types and the parsing are |
| 23 | //! tested; the responder's live behaviour on a real network is not, and wants a |
| 24 | //! second machine to test against rather than a unit test. |
| 25 | //! |
| 26 | //! # What is not |
| 27 | //! |
| 28 | //! The description document a `LOCATION` points at, and the SOAP services behind |
| 29 | //! it, are UPnP rather than SSDP and live in [`crate::upnp`]. |
| 30 | //! |
| 31 | //! # Two responders |
| 32 | //! |
| 33 | //! [`Responder`] is async over tokio. [`SyncResponder`] is the same protocol over |
| 34 | //! `std::net`, for a binary that wants no runtime: discovery is one socket, three |
| 35 | //! message shapes and a thread that blocks on a read, and pulling in an executor |
| 36 | //! to run it is a poor trade. Neither is a wrapper around the other; they share |
| 37 | //! the messages above, which is where the protocol actually is. |
| 38 | //! |
| 39 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 40 | //! Anthropic Claude |
| 41 | |
| 42 | use crate::time::Time; |
| 43 | |
| 44 | use oxedyne_fe2o3_core::prelude::*; |
| 45 | |
| 46 | use std::{ |
| 47 | collections::BTreeMap, |
| 48 | fmt, |
| 49 | net::{ |
| 50 | Ipv4Addr, |
| 51 | SocketAddr, |
| 52 | SocketAddrV4, |
| 53 | }, |
| 54 | str::FromStr, |
| 55 | time::{ |
| 56 | Duration, |
| 57 | SystemTime, |
| 58 | UNIX_EPOCH, |
| 59 | }, |
| 60 | }; |
| 61 | |
| 62 | use tokio::net::UdpSocket; |
| 63 | |
| 64 | |
| 65 | //// The group, and the port that goes with it (UPnP DA 2.0 §1.1.1). |
| 66 | pub const MULTICAST_ADDR: Ipv4Addr = Ipv4Addr::new(239, 255, 255, 250); |
| 67 | pub const PORT: u16 = 1900; |
| 68 | |
| 69 | // The value of `MAN` on a search. The quotes are part of the field, and a |
| 70 | // searcher that leaves them off is ignored by conforming devices. |
| 71 | pub const DISCOVER: &str = "\"ssdp:discover\""; |
| 72 | |
| 73 | // The largest datagram this reader will accept. A conforming SSDP message is a |
| 74 | // few hundred bytes; the rest of the space is for the long `SERVER` and |
| 75 | // vendor-extension fields that real devices send. |
| 76 | pub const MAX_DATAGRAM: usize = 2048; |
| 77 | |
| 78 | // How long an announcement stands before it must be renewed. The specification |
| 79 | // requires at least 1800. |
| 80 | pub const DEFAULT_MAX_AGE: u32 = 1800; // seconds |
| 81 | |
| 82 | |
| 83 | /// What a message is about: a device type, a service type, a particular device by |
| 84 | /// its UUID, or everything at once. |
| 85 | /// |
| 86 | /// Held as an enum because a responder answers `All` and `RootDevice` and its own |
| 87 | /// `Uuid` differently, and a string comparison spread across the call sites is how |
| 88 | /// one of those cases quietly stops being answered. |
| 89 | #[derive(Clone, Debug, Eq, Ord, PartialEq, PartialOrd)] |
| 90 | pub enum Target { |
| 91 | All, // `ssdp:all`, answered once per thing |
| 92 | RootDevice, // `upnp:rootdevice` |
| 93 | Uuid(String), // `uuid:...`, one particular device |
| 94 | Urn(String), // `urn:...`, a device or service type |
| 95 | Other(String), // anything else a device chose to name itself by |
| 96 | } |
| 97 | |
| 98 | impl fmt::Display for Target { |
| 99 | fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { |
| 100 | match self { |
| 101 | Self::All => write!(f, "ssdp:all"), |
| 102 | Self::RootDevice => write!(f, "upnp:rootdevice"), |
| 103 | Self::Uuid(s) => write!(f, "uuid:{}", s), |
| 104 | Self::Urn(s) => write!(f, "urn:{}", s), |
| 105 | Self::Other(s) => write!(f, "{}", s), |
| 106 | } |
| 107 | } |
| 108 | } |
| 109 | |
| 110 | impl FromStr for Target { |
| 111 | type Err = Error<ErrTag>; |
| 112 | |
| 113 | fn from_str(s: &str) -> std::result::Result<Self, Self::Err> { |
| 114 | let s = s.trim(); |
| 115 | Ok(match s { |
| 116 | "ssdp:all" => Self::All, |
| 117 | "upnp:rootdevice" => Self::RootDevice, |
| 118 | _ => match s.split_once(':') { |
| 119 | Some(("uuid", rest)) => Self::Uuid(rest.to_string()), |
| 120 | Some(("urn", rest)) => Self::Urn(rest.to_string()), |
| 121 | _ => Self::Other(s.to_string()), |
| 122 | }, |
| 123 | }) |
| 124 | } |
| 125 | } |
| 126 | |
| 127 | impl Target { |
| 128 | /// Does an announcement of `self` answer a search for `wanted`? `ssdp:all` |
| 129 | /// matches everything, and everything else matches only itself. |
| 130 | pub fn answers(&self, wanted: &Target) -> bool { |
| 131 | match wanted { |
| 132 | Target::All => true, |
| 133 | other => self == other, |
| 134 | } |
| 135 | } |
| 136 | } |
| 137 | |
| 138 | /// What a `NOTIFY` is saying, carried in its `NTS` field. |
| 139 | #[derive(Clone, Copy, Debug, Eq, PartialEq)] |
| 140 | pub enum Nts { |
| 141 | Alive, // here, and stays for its `CACHE-CONTROL` lifetime |
| 142 | ByeBye, // going, now |
| 143 | Update, // still here, on a new boot identifier, UPnP DA 2.0 §1.2.4 |
| 144 | } |
| 145 | |
| 146 | impl fmt::Display for Nts { |
| 147 | fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { |
| 148 | write!(f, "{}", match self { |
| 149 | Self::Alive => "ssdp:alive", |
| 150 | Self::ByeBye => "ssdp:byebye", |
| 151 | Self::Update => "ssdp:update", |
| 152 | }) |
| 153 | } |
| 154 | } |
| 155 | |
| 156 | impl FromStr for Nts { |
| 157 | type Err = Error<ErrTag>; |
| 158 | |
| 159 | fn from_str(s: &str) -> std::result::Result<Self, Self::Err> { |
| 160 | Ok(match s.trim() { |
| 161 | "ssdp:alive" => Self::Alive, |
| 162 | "ssdp:byebye" => Self::ByeBye, |
| 163 | "ssdp:update" => Self::Update, |
| 164 | _ => return Err(err!( |
| 165 | "Unrecognised SSDP NTS '{}'.", s; |
| 166 | IO, Network, Unknown, Input)), |
| 167 | }) |
| 168 | } |
| 169 | } |
| 170 | |
| 171 | /// A search: `M-SEARCH * HTTP/1.1`. |
| 172 | #[derive(Clone, Debug, Eq, PartialEq)] |
| 173 | pub struct Search { |
| 174 | pub target: Target, |
| 175 | // `MX`: the largest number of seconds a responder should wait before |
| 176 | // answering, so a hundred devices do not answer at the same instant. |
| 177 | pub mx: u8, |
| 178 | pub user_agent: Option<String>, // who asked, and what they call themselves |
| 179 | pub extra: BTreeMap<String, String>, // fields this crate does not model |
| 180 | } |
| 181 | |
| 182 | impl Search { |
| 183 | /// The customary two second spread. |
| 184 | pub fn new(target: Target) -> Self { |
| 185 | Self { |
| 186 | target, |
| 187 | mx: 2, |
| 188 | user_agent: None, |
| 189 | extra: BTreeMap::new(), |
| 190 | } |
| 191 | } |
| 192 | } |
| 193 | |
| 194 | /// A unicast answer to a search: `HTTP/1.1 200 OK`. |
| 195 | #[derive(Clone, Debug, Eq, PartialEq)] |
| 196 | pub struct SearchResponse { |
| 197 | pub max_age: u32, // seconds the answer stands |
| 198 | pub date: Option<String>, // HTTP date, generated by Responder when absent |
| 199 | pub location: String, // the device description document |
| 200 | pub server: String, // OS, UPnP version, product |
| 201 | pub target: Target, // echoed back from the search |
| 202 | pub usn: String, // which particular thing is answering |
| 203 | // `BOOTID.UPNP.ORG` changes when the device restarts, `CONFIGID.UPNP.ORG` |
| 204 | // when its description does. |
| 205 | pub boot_id: Option<u32>, |
| 206 | pub config_id: Option<u32>, |
| 207 | pub extra: BTreeMap<String, String>, // fields this crate does not model |
| 208 | } |
| 209 | |
| 210 | /// An announcement: `NOTIFY * HTTP/1.1`. |
| 211 | #[derive(Clone, Debug, Eq, PartialEq)] |
| 212 | pub struct Notify { |
| 213 | pub target: Target, |
| 214 | pub nts: Nts, // arriving, leaving, or renewing |
| 215 | pub usn: String, // which particular thing |
| 216 | // A `byebye` carries neither of these: there is nothing left to fetch, and |
| 217 | // the announcement stands until contradicted. |
| 218 | pub max_age: Option<u32>, // seconds |
| 219 | pub location: Option<String>, // the description document |
| 220 | pub server: Option<String>, |
| 221 | pub boot_id: Option<u32>, // `BOOTID.UPNP.ORG` |
| 222 | pub config_id: Option<u32>, // `CONFIGID.UPNP.ORG` |
| 223 | pub extra: BTreeMap<String, String>, // fields this crate does not model |
| 224 | } |
| 225 | |
| 226 | impl Notify { |
| 227 | /// An `ssdp:alive` for a thing that has just appeared, or is renewing. |
| 228 | pub fn alive(target: Target, usn: String, location: String, server: String) -> Self { |
| 229 | Self { |
| 230 | target, |
| 231 | nts: Nts::Alive, |
| 232 | usn, |
| 233 | max_age: Some(DEFAULT_MAX_AGE), |
| 234 | location: Some(location), |
| 235 | server: Some(server), |
| 236 | boot_id: None, |
| 237 | config_id: None, |
| 238 | extra: BTreeMap::new(), |
| 239 | } |
| 240 | } |
| 241 | |
| 242 | /// Carries neither a lifetime nor a location: there is nothing left to fetch, |
| 243 | /// and saying otherwise leaves a control point holding a URL that has stopped |
| 244 | /// answering. |
| 245 | pub fn byebye(target: Target, usn: String) -> Self { |
| 246 | Self { |
| 247 | target, |
| 248 | nts: Nts::ByeBye, |
| 249 | usn, |
| 250 | max_age: None, |
| 251 | location: None, |
| 252 | server: None, |
| 253 | boot_id: None, |
| 254 | config_id: None, |
| 255 | extra: BTreeMap::new(), |
| 256 | } |
| 257 | } |
| 258 | } |
| 259 | |
| 260 | /// One SSDP datagram. |
| 261 | #[derive(Clone, Debug, Eq, PartialEq)] |
| 262 | pub enum SsdpMessage { |
| 263 | Search(Search), |
| 264 | Response(SearchResponse), |
| 265 | Notify(Notify), |
| 266 | } |
| 267 | |
| 268 | impl SsdpMessage { |
| 269 | |
| 270 | /// Real networks carry SSDP from devices that get the details wrong, so this |
| 271 | /// is forgiving about spacing, field name case and line endings, and strict |
| 272 | /// only about the things that decide what the message means: the start line, |
| 273 | /// and the fields the message cannot be acted on without. |
| 274 | pub fn parse(bytes: &[u8]) -> Outcome<Self> { |
| 275 | let text = String::from_utf8_lossy(bytes); |
| 276 | let mut lines = text.split('\n'); |
| 277 | let start = match lines.next() { |
| 278 | Some(line) => line.trim_end_matches('\r').trim(), |
| 279 | None => return Err(err!( |
| 280 | "An SSDP datagram was empty."; IO, Network, Input, Missing)), |
| 281 | }; |
| 282 | |
| 283 | let mut fields: BTreeMap<String, String> = BTreeMap::new(); |
| 284 | for line in lines { |
| 285 | let line = line.trim_end_matches('\r'); |
| 286 | if line.trim().is_empty() { |
| 287 | continue; |
| 288 | } |
| 289 | match line.split_once(':') { |
| 290 | Some((name, value)) => { |
| 291 | fields.insert(name.trim().to_uppercase(), value.trim().to_string()); |
| 292 | } |
| 293 | // A line with no colon is not a field. Devices send them; they |
| 294 | // are of no use and are not worth refusing the message over. |
| 295 | None => continue, |
| 296 | } |
| 297 | } |
| 298 | |
| 299 | let upper = start.to_uppercase(); |
| 300 | if upper.starts_with("M-SEARCH") { |
| 301 | return Ok(Self::Search(res!(Search::from_fields(&mut fields)))); |
| 302 | } |
| 303 | if upper.starts_with("NOTIFY") { |
| 304 | return Ok(Self::Notify(res!(Notify::from_fields(&mut fields)))); |
| 305 | } |
| 306 | if upper.starts_with("HTTP/") { |
| 307 | // Only a 200 is an answer. Anything else is a device refusing, and |
| 308 | // there is nothing to be discovered from it. |
| 309 | let code = start.split_whitespace().nth(1).unwrap_or(""); |
| 310 | if code != "200" { |
| 311 | return Err(err!( |
| 312 | "An SSDP response answered '{}', which is not a discovery.", start; |
| 313 | IO, Network, Input, Invalid)); |
| 314 | } |
| 315 | return Ok(Self::Response(res!(SearchResponse::from_fields(&mut fields)))); |
| 316 | } |
| 317 | |
| 318 | Err(err!( |
| 319 | "'{}' is not the start line of any SSDP message.", start; |
| 320 | IO, Network, Input, Invalid)) |
| 321 | } |
| 322 | |
| 323 | pub fn as_bytes(&self) -> Vec<u8> { |
| 324 | self.to_string().into_bytes() |
| 325 | } |
| 326 | } |
| 327 | |
| 328 | impl fmt::Display for SsdpMessage { |
| 329 | fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { |
| 330 | match self { |
| 331 | Self::Search(m) => write!(f, "{}", m.as_text()), |
| 332 | Self::Response(m) => write!(f, "{}", m.as_text()), |
| 333 | Self::Notify(m) => write!(f, "{}", m.as_text()), |
| 334 | } |
| 335 | } |
| 336 | } |
| 337 | |
| 338 | |
| 339 | // ┌───────────────────────────────────────────────────────────────────────────┐ |
| 340 | // │ FIELDS TO MESSAGES, AND BACK │ |
| 341 | // └───────────────────────────────────────────────────────────────────────────┘ |
| 342 | |
| 343 | /// Take a field out of the map, so what is left over is the extras. |
| 344 | fn take(fields: &mut BTreeMap<String, String>, name: &str) -> Option<String> { |
| 345 | fields.remove(name) |
| 346 | } |
| 347 | |
| 348 | fn need(fields: &mut BTreeMap<String, String>, name: &str, what: &str) -> Outcome<String> { |
| 349 | match fields.remove(name) { |
| 350 | Some(v) => Ok(v), |
| 351 | None => Err(err!( |
| 352 | "An SSDP {} carried no {} field.", what, name; |
| 353 | IO, Network, Input, Missing)), |
| 354 | } |
| 355 | } |
| 356 | |
| 357 | /// Read a `CACHE-CONTROL: max-age=1800` into its seconds. |
| 358 | fn max_age_of(value: &str) -> Option<u32> { |
| 359 | for part in value.split(',') { |
| 360 | if let Some((k, v)) = part.split_once('=') { |
| 361 | if k.trim().eq_ignore_ascii_case("max-age") { |
| 362 | return v.trim().parse().ok(); |
| 363 | } |
| 364 | } |
| 365 | } |
| 366 | None |
| 367 | } |
| 368 | |
| 369 | /// In a stable order, so a message written twice is the same bytes twice. |
| 370 | fn extra_lines(extra: &BTreeMap<String, String>) -> String { |
| 371 | let mut out = String::new(); |
| 372 | for (name, value) in extra { |
| 373 | out.push_str(&fmt!("{}: {}\r\n", name, value)); |
| 374 | } |
| 375 | out |
| 376 | } |
| 377 | |
| 378 | impl Search { |
| 379 | |
| 380 | fn from_fields(fields: &mut BTreeMap<String, String>) -> Outcome<Self> { |
| 381 | // `HOST` and `MAN` are required of a sender and useless to a receiver |
| 382 | // that already knows where the datagram arrived and what kind it is, so |
| 383 | // they are dropped rather than checked: a device that sends a slightly |
| 384 | // wrong `MAN` is still searching. |
| 385 | let _ = take(fields, "HOST"); |
| 386 | let _ = take(fields, "MAN"); |
| 387 | let st = res!(need(fields, "ST", "search")); |
| 388 | let target = res!(Target::from_str(&st)); |
| 389 | let mx = take(fields, "MX") |
| 390 | .and_then(|v| v.trim().parse::<u8>().ok()) |
| 391 | .unwrap_or(0); |
| 392 | let user_agent = take(fields, "USER-AGENT"); |
| 393 | Ok(Self { |
| 394 | target, |
| 395 | mx, |
| 396 | user_agent, |
| 397 | extra: std::mem::take(fields), |
| 398 | }) |
| 399 | } |
| 400 | |
| 401 | pub fn as_text(&self) -> String { |
| 402 | let mut out = String::new(); |
| 403 | out.push_str("M-SEARCH * HTTP/1.1\r\n"); |
| 404 | out.push_str(&fmt!("HOST: {}:{}\r\n", MULTICAST_ADDR, PORT)); |
| 405 | out.push_str(&fmt!("MAN: {}\r\n", DISCOVER)); |
| 406 | out.push_str(&fmt!("MX: {}\r\n", self.mx)); |
| 407 | out.push_str(&fmt!("ST: {}\r\n", self.target)); |
| 408 | if let Some(ua) = &self.user_agent { |
| 409 | out.push_str(&fmt!("USER-AGENT: {}\r\n", ua)); |
| 410 | } |
| 411 | out.push_str(&extra_lines(&self.extra)); |
| 412 | out.push_str("\r\n"); |
| 413 | out |
| 414 | } |
| 415 | } |
| 416 | |
| 417 | impl SearchResponse { |
| 418 | |
| 419 | fn from_fields(fields: &mut BTreeMap<String, String>) -> Outcome<Self> { |
| 420 | let _ = take(fields, "EXT"); |
| 421 | let max_age = take(fields, "CACHE-CONTROL") |
| 422 | .and_then(|v| max_age_of(&v)) |
| 423 | .unwrap_or(DEFAULT_MAX_AGE); |
| 424 | let date = take(fields, "DATE"); |
| 425 | let location = res!(need(fields, "LOCATION", "response")); |
| 426 | let server = take(fields, "SERVER").unwrap_or_default(); |
| 427 | let st = res!(need(fields, "ST", "response")); |
| 428 | let target = res!(Target::from_str(&st)); |
| 429 | let usn = res!(need(fields, "USN", "response")); |
| 430 | let boot_id = take(fields, "BOOTID.UPNP.ORG").and_then(|v| v.trim().parse().ok()); |
| 431 | let config_id = take(fields, "CONFIGID.UPNP.ORG").and_then(|v| v.trim().parse().ok()); |
| 432 | Ok(Self { |
| 433 | max_age, |
| 434 | date, |
| 435 | location, |
| 436 | server, |
| 437 | target, |
| 438 | usn, |
| 439 | boot_id, |
| 440 | config_id, |
| 441 | extra: std::mem::take(fields), |
| 442 | }) |
| 443 | } |
| 444 | |
| 445 | /// `EXT:` is an empty field that means nothing and is required anyway |
| 446 | /// (UPnP DA 2.0 §1.3.3): a control point that does not see it discards the |
| 447 | /// answer. |
| 448 | pub fn as_text(&self) -> String { |
| 449 | let mut out = String::new(); |
| 450 | out.push_str("HTTP/1.1 200 OK\r\n"); |
| 451 | out.push_str(&fmt!("CACHE-CONTROL: max-age={}\r\n", self.max_age)); |
| 452 | if let Some(date) = &self.date { |
| 453 | out.push_str(&fmt!("DATE: {}\r\n", date)); |
| 454 | } |
| 455 | out.push_str("EXT:\r\n"); |
| 456 | out.push_str(&fmt!("LOCATION: {}\r\n", self.location)); |
| 457 | out.push_str(&fmt!("SERVER: {}\r\n", self.server)); |
| 458 | out.push_str(&fmt!("ST: {}\r\n", self.target)); |
| 459 | out.push_str(&fmt!("USN: {}\r\n", self.usn)); |
| 460 | if let Some(id) = self.boot_id { |
| 461 | out.push_str(&fmt!("BOOTID.UPNP.ORG: {}\r\n", id)); |
| 462 | } |
| 463 | if let Some(id) = self.config_id { |
| 464 | out.push_str(&fmt!("CONFIGID.UPNP.ORG: {}\r\n", id)); |
| 465 | } |
| 466 | out.push_str(&extra_lines(&self.extra)); |
| 467 | out.push_str("\r\n"); |
| 468 | out |
| 469 | } |
| 470 | } |
| 471 | |
| 472 | impl Notify { |
| 473 | |
| 474 | fn from_fields(fields: &mut BTreeMap<String, String>) -> Outcome<Self> { |
| 475 | let _ = take(fields, "HOST"); |
| 476 | let nt = res!(need(fields, "NT", "notification")); |
| 477 | let target = res!(Target::from_str(&nt)); |
| 478 | let nts_txt = res!(need(fields, "NTS", "notification")); |
| 479 | let nts = res!(Nts::from_str(&nts_txt)); |
| 480 | let usn = res!(need(fields, "USN", "notification")); |
| 481 | let max_age = take(fields, "CACHE-CONTROL").and_then(|v| max_age_of(&v)); |
| 482 | let location = take(fields, "LOCATION"); |
| 483 | let server = take(fields, "SERVER"); |
| 484 | let boot_id = take(fields, "BOOTID.UPNP.ORG").and_then(|v| v.trim().parse().ok()); |
| 485 | let config_id = take(fields, "CONFIGID.UPNP.ORG").and_then(|v| v.trim().parse().ok()); |
| 486 | Ok(Self { |
| 487 | target, |
| 488 | nts, |
| 489 | usn, |
| 490 | max_age, |
| 491 | location, |
| 492 | server, |
| 493 | boot_id, |
| 494 | config_id, |
| 495 | extra: std::mem::take(fields), |
| 496 | }) |
| 497 | } |
| 498 | |
| 499 | pub fn as_text(&self) -> String { |
| 500 | let mut out = String::new(); |
| 501 | out.push_str("NOTIFY * HTTP/1.1\r\n"); |
| 502 | out.push_str(&fmt!("HOST: {}:{}\r\n", MULTICAST_ADDR, PORT)); |
| 503 | if let Some(age) = self.max_age { |
| 504 | out.push_str(&fmt!("CACHE-CONTROL: max-age={}\r\n", age)); |
| 505 | } |
| 506 | if let Some(loc) = &self.location { |
| 507 | out.push_str(&fmt!("LOCATION: {}\r\n", loc)); |
| 508 | } |
| 509 | out.push_str(&fmt!("NT: {}\r\n", self.target)); |
| 510 | out.push_str(&fmt!("NTS: {}\r\n", self.nts)); |
| 511 | if let Some(server) = &self.server { |
| 512 | out.push_str(&fmt!("SERVER: {}\r\n", server)); |
| 513 | } |
| 514 | out.push_str(&fmt!("USN: {}\r\n", self.usn)); |
| 515 | if let Some(id) = self.boot_id { |
| 516 | out.push_str(&fmt!("BOOTID.UPNP.ORG: {}\r\n", id)); |
| 517 | } |
| 518 | if let Some(id) = self.config_id { |
| 519 | out.push_str(&fmt!("CONFIGID.UPNP.ORG: {}\r\n", id)); |
| 520 | } |
| 521 | out.push_str(&extra_lines(&self.extra)); |
| 522 | out.push_str("\r\n"); |
| 523 | out |
| 524 | } |
| 525 | } |
| 526 | |
| 527 | |
| 528 | // ┌───────────────────────────────────────────────────────────────────────────┐ |
| 529 | // │ THE SOCKET │ |
| 530 | // └───────────────────────────────────────────────────────────────────────────┘ |
| 531 | |
| 532 | /// A socket bound to the SSDP group on one interface, reading and writing the |
| 533 | /// messages above. |
| 534 | /// |
| 535 | /// # One interface at a time |
| 536 | /// |
| 537 | /// A machine with two networks on it is on two SSDP groups, and a single socket |
| 538 | /// joined to both cannot say which one a datagram came from or choose which one an |
| 539 | /// announcement goes out on. So a responder is per-interface, named by the local |
| 540 | /// address to bind and join on, and a caller with several interfaces holds several |
| 541 | /// responders. |
| 542 | /// |
| 543 | /// # What this cannot do yet |
| 544 | /// |
| 545 | /// Two processes cannot both listen on port 1900: that needs `SO_REUSEADDR`, which |
| 546 | /// neither the standard library nor tokio exposes on a bound socket, and setting |
| 547 | /// it means a socket option this crate has no dependency to set. A machine already |
| 548 | /// running a UPnP daemon will therefore refuse the bind, with the address in the |
| 549 | /// error. Likewise the outgoing interface for a multicast send is chosen by |
| 550 | /// binding to that interface's address rather than by `IP_MULTICAST_IF`, which is |
| 551 | /// what the kernel does with a bound source address on Linux and every BSD. |
| 552 | #[derive(Debug)] |
| 553 | pub struct Responder { |
| 554 | socket: UdpSocket, |
| 555 | iface: Ipv4Addr, // bound to, and announced on |
| 556 | } |
| 557 | |
| 558 | impl Responder { |
| 559 | |
| 560 | /// `iface` is the local IPv4 address of the interface to speak on; |
| 561 | /// `Ipv4Addr::UNSPECIFIED` lets the kernel choose, which is right on a machine |
| 562 | /// with one network and wrong on a machine with two. |
| 563 | pub async fn bind(iface: Ipv4Addr) -> Outcome<Self> { |
| 564 | Self::bind_to_port(iface, PORT).await |
| 565 | } |
| 566 | |
| 567 | /// A test wants an ephemeral port; a real responder wants 1900, because that |
| 568 | /// is where searches are sent. |
| 569 | pub async fn bind_to_port(iface: Ipv4Addr, port: u16) -> Outcome<Self> { |
| 570 | let bind_addr = SocketAddrV4::new(Ipv4Addr::UNSPECIFIED, port); |
| 571 | let socket = match UdpSocket::bind(bind_addr).await { |
| 572 | Ok(s) => s, |
| 573 | Err(e) => return Err(err!(e, |
| 574 | "Binding the SSDP port {}. A UPnP daemon already listening there \ |
| 575 | holds it exclusively, since this socket does not ask to share it.", |
| 576 | port; |
| 577 | IO, Network, Init)), |
| 578 | }; |
| 579 | let result = socket.join_multicast_v4(MULTICAST_ADDR, iface); |
| 580 | res!(result, IO, Network, Init); |
| 581 | // An announcement should be heard by the machine that sent it, since a |
| 582 | // control point may be running beside the device it is discovering. |
| 583 | let result = socket.set_multicast_loop_v4(true); |
| 584 | res!(result, IO, Network, Init); |
| 585 | Ok(Self { |
| 586 | socket, |
| 587 | iface, |
| 588 | }) |
| 589 | } |
| 590 | |
| 591 | pub fn local_addr(&self) -> Outcome<SocketAddr> { |
| 592 | let result = self.socket.local_addr(); |
| 593 | Ok(res!(result, IO, Network)) |
| 594 | } |
| 595 | |
| 596 | pub fn interface(&self) -> Ipv4Addr { |
| 597 | self.iface |
| 598 | } |
| 599 | |
| 600 | /// A datagram that does not parse is not an error the caller can do anything |
| 601 | /// about -- the network carries plenty of them -- so it is logged and the wait |
| 602 | /// resumes. |
| 603 | pub async fn recv(&self) -> Outcome<(SsdpMessage, SocketAddr)> { |
| 604 | let mut buf = vec![0u8; MAX_DATAGRAM]; |
| 605 | loop { |
| 606 | let (n, from) = match self.socket.recv_from(&mut buf).await { |
| 607 | Ok(pair) => pair, |
| 608 | Err(e) => return Err(err!(e, |
| 609 | "Reading an SSDP datagram."; IO, Network, Wire, Read)), |
| 610 | }; |
| 611 | match SsdpMessage::parse(&buf[..n]) { |
| 612 | Ok(msg) => return Ok((msg, from)), |
| 613 | Err(e) => { |
| 614 | debug!("An SSDP datagram from {} did not parse: {}", from, e); |
| 615 | continue; |
| 616 | } |
| 617 | } |
| 618 | } |
| 619 | } |
| 620 | |
| 621 | pub async fn multicast(&self, msg: &SsdpMessage) -> Outcome<()> { |
| 622 | let to = SocketAddrV4::new(MULTICAST_ADDR, PORT); |
| 623 | res!(self.send_to(msg, SocketAddr::V4(to)).await); |
| 624 | Ok(()) |
| 625 | } |
| 626 | |
| 627 | pub async fn send_to(&self, msg: &SsdpMessage, to: SocketAddr) -> Outcome<()> { |
| 628 | let bytes = msg.as_bytes(); |
| 629 | let result = self.socket.send_to(&bytes, to).await; |
| 630 | let sent = res!(result, IO, Network, Wire, Write); |
| 631 | if sent != bytes.len() { |
| 632 | return Err(err!( |
| 633 | "An SSDP datagram of {} bytes went out as {}.", bytes.len(), sent; |
| 634 | IO, Network, Wire, Write, Size)); |
| 635 | } |
| 636 | Ok(()) |
| 637 | } |
| 638 | |
| 639 | /// Goes back to the address the search came from. The `ST` of the answer is |
| 640 | /// the target actually being announced, not the `ssdp:all` that may have been |
| 641 | /// asked: a control point matches the two, and an answer that echoes |
| 642 | /// `ssdp:all` is discarded. |
| 643 | pub async fn answer( |
| 644 | &self, |
| 645 | to: SocketAddr, |
| 646 | target: Target, |
| 647 | usn: String, |
| 648 | location: String, |
| 649 | server: String, |
| 650 | ) |
| 651 | -> Outcome<()> |
| 652 | { |
| 653 | let response = SearchResponse { |
| 654 | max_age: DEFAULT_MAX_AGE, |
| 655 | date: Some(res!(http_date())), |
| 656 | location, |
| 657 | server, |
| 658 | target, |
| 659 | usn, |
| 660 | boot_id: None, |
| 661 | config_id: None, |
| 662 | extra: BTreeMap::new(), |
| 663 | }; |
| 664 | res!(self.send_to(&SsdpMessage::Response(response), to).await); |
| 665 | Ok(()) |
| 666 | } |
| 667 | } |
| 668 | |
| 669 | |
| 670 | // ┌───────────────────────────────────────────────────────────────────────────┐ |
| 671 | // │ THE SAME SOCKET, WITHOUT A RUNTIME │ |
| 672 | // └───────────────────────────────────────────────────────────────────────────┘ |
| 673 | |
| 674 | /// A blocking SSDP socket: the same protocol as [`Responder`] over `std::net`. |
| 675 | /// |
| 676 | /// A media server is a thread that blocks on a read and answers what arrives. |
| 677 | /// That is the whole of discovery, and it needs no executor. |
| 678 | /// |
| 679 | /// # Where the datagrams go |
| 680 | /// |
| 681 | /// One socket is bound to `0.0.0.0` on the SSDP port and joined to the group on |
| 682 | /// each interface named by [`SyncResponder::join`]. A search is answered by |
| 683 | /// unicast back to whoever sent it, which the routing table places correctly |
| 684 | /// however many interfaces there are. A multicast announcement, though, leaves by |
| 685 | /// the interface the kernel picks, because `IP_MULTICAST_IF` is not reachable |
| 686 | /// through the standard library and this crate carries no socket-options |
| 687 | /// dependency. On a machine with one network -- which is what a household has -- |
| 688 | /// the two are the same interface. |
| 689 | /// |
| 690 | /// # Sharing the port |
| 691 | /// |
| 692 | /// Two processes cannot both hold port 1900 here: `SO_REUSEADDR` and |
| 693 | /// `SO_REUSEPORT` are likewise out of the standard library's reach. A machine |
| 694 | /// already running a UPnP daemon refuses the bind, with the address in the error. |
| 695 | #[derive(Debug)] |
| 696 | pub struct SyncResponder { |
| 697 | socket: std::net::UdpSocket, |
| 698 | joined: Vec<Ipv4Addr>, // the interfaces the group was successfully joined on |
| 699 | } |
| 700 | |
| 701 | impl SyncResponder { |
| 702 | |
| 703 | /// On every address, and joining no group yet. |
| 704 | pub fn bind() -> Outcome<Self> { |
| 705 | Self::bind_to_port(PORT) |
| 706 | } |
| 707 | |
| 708 | /// The same, on a port of the caller's choosing. A test wants an ephemeral |
| 709 | /// one; a real responder wants 1900, because that is where searches are sent. |
| 710 | pub fn bind_to_port(port: u16) -> Outcome<Self> { |
| 711 | let bind_addr = SocketAddrV4::new(Ipv4Addr::UNSPECIFIED, port); |
| 712 | let socket = match std::net::UdpSocket::bind(bind_addr) { |
| 713 | Ok(s) => s, |
| 714 | Err(e) => return Err(err!(e, |
| 715 | "Binding the SSDP port {}. A UPnP daemon already listening there \ |
| 716 | holds it exclusively, since this socket does not ask to share it.", |
| 717 | port; |
| 718 | IO, Network, Init)), |
| 719 | }; |
| 720 | // An announcement should be heard by the machine that sent it, since a |
| 721 | // control point may be running beside the device it is discovering. |
| 722 | let result = socket.set_multicast_loop_v4(true); |
| 723 | res!(result, IO, Network, Init); |
| 724 | Ok(Self { |
| 725 | socket, |
| 726 | joined: Vec::new(), |
| 727 | }) |
| 728 | } |
| 729 | |
| 730 | /// `iface` is the interface's local IPv4 address. Joining the same interface |
| 731 | /// twice is an error from the kernel and is |
| 732 | /// reported as one; [`Self::join_every_interface`] is the call that tolerates |
| 733 | /// it, since it does not know what is already joined. |
| 734 | pub fn join(&mut self, iface: Ipv4Addr) -> Outcome<()> { |
| 735 | let result = self.socket.join_multicast_v4(&MULTICAST_ADDR, &iface); |
| 736 | res!(result, IO, Network, Init); |
| 737 | self.joined.push(iface); |
| 738 | Ok(()) |
| 739 | } |
| 740 | |
| 741 | /// Best effort by design: an interface that refuses the join is skipped rather |
| 742 | /// than failing the lot, because one unusable interface on a machine with three |
| 743 | /// must not stop a television on the other two from finding anything. An answer |
| 744 | /// of zero is the caller's cue to complain. |
| 745 | pub fn join_every_interface(&mut self) -> Outcome<usize> { |
| 746 | let mut joined = 0usize; |
| 747 | for iface in local_interfaces() { |
| 748 | if self.joined.contains(&iface) { |
| 749 | continue; |
| 750 | } |
| 751 | match self.join(iface) { |
| 752 | Ok(()) => joined += 1, |
| 753 | Err(e) => debug!("SSDP could not join the group on {}: {}", iface, e), |
| 754 | } |
| 755 | } |
| 756 | Ok(joined) |
| 757 | } |
| 758 | |
| 759 | pub fn interfaces(&self) -> &[Ipv4Addr] { |
| 760 | &self.joined |
| 761 | } |
| 762 | |
| 763 | pub fn local_addr(&self) -> Outcome<SocketAddr> { |
| 764 | let result = self.socket.local_addr(); |
| 765 | Ok(res!(result, IO, Network)) |
| 766 | } |
| 767 | |
| 768 | /// How long a read may block before it gives up. `None` blocks for ever, |
| 769 | /// which leaves no way to shut the thread down. |
| 770 | pub fn set_timeout(&self, how_long: Option<Duration>) -> Outcome<()> { |
| 771 | let result = self.socket.set_read_timeout(how_long); |
| 772 | res!(result, IO, Network); |
| 773 | Ok(()) |
| 774 | } |
| 775 | |
| 776 | /// Another handle on the same socket, so that one thread can read while |
| 777 | /// another announces. |
| 778 | pub fn try_clone(&self) -> Outcome<Self> { |
| 779 | let result = self.socket.try_clone(); |
| 780 | let socket = res!(result, IO, Network); |
| 781 | Ok(Self { |
| 782 | socket, |
| 783 | joined: self.joined.clone(), |
| 784 | }) |
| 785 | } |
| 786 | |
| 787 | /// `Ok(None)` means the read timed out, which is how a shutdown flag gets |
| 788 | /// looked at. A datagram that does not parse is not an error the caller can do |
| 789 | /// anything about -- the network carries plenty of them -- so it is logged and |
| 790 | /// the wait resumes. |
| 791 | pub fn recv(&self) -> Outcome<Option<(SsdpMessage, SocketAddr)>> { |
| 792 | let mut buf = vec![0u8; MAX_DATAGRAM]; |
| 793 | loop { |
| 794 | let (n, from) = match self.socket.recv_from(&mut buf) { |
| 795 | Ok(pair) => pair, |
| 796 | Err(e) => { |
| 797 | // A timeout is spelled two ways depending on the platform, |
| 798 | // and neither is a failure. |
| 799 | if matches!(e.kind(), |
| 800 | std::io::ErrorKind::WouldBlock | std::io::ErrorKind::TimedOut) |
| 801 | { |
| 802 | return Ok(None); |
| 803 | } |
| 804 | return Err(err!(e, |
| 805 | "Reading an SSDP datagram."; IO, Network, Wire, Read)); |
| 806 | }, |
| 807 | }; |
| 808 | match SsdpMessage::parse(&buf[..n]) { |
| 809 | Ok(msg) => return Ok(Some((msg, from))), |
| 810 | Err(e) => { |
| 811 | debug!("An SSDP datagram from {} did not parse: {}", from, e); |
| 812 | continue; |
| 813 | }, |
| 814 | } |
| 815 | } |
| 816 | } |
| 817 | |
| 818 | pub fn multicast(&self, msg: &SsdpMessage) -> Outcome<()> { |
| 819 | let to = SocketAddrV4::new(MULTICAST_ADDR, PORT); |
| 820 | res!(self.send_to(msg, SocketAddr::V4(to))); |
| 821 | Ok(()) |
| 822 | } |
| 823 | |
| 824 | pub fn send_to(&self, msg: &SsdpMessage, to: SocketAddr) -> Outcome<()> { |
| 825 | let bytes = msg.as_bytes(); |
| 826 | let result = self.socket.send_to(&bytes, to); |
| 827 | let sent = res!(result, IO, Network, Wire, Write); |
| 828 | if sent != bytes.len() { |
| 829 | return Err(err!( |
| 830 | "An SSDP datagram of {} bytes went out as {}.", bytes.len(), sent; |
| 831 | IO, Network, Wire, Write, Size)); |
| 832 | } |
| 833 | Ok(()) |
| 834 | } |
| 835 | |
| 836 | /// Goes back to the address the search came from. The `ST` of the answer is |
| 837 | /// the target actually being announced, not the `ssdp:all` that may have been |
| 838 | /// asked: a control point matches the two, and an answer that echoes |
| 839 | /// `ssdp:all` is discarded. |
| 840 | pub fn answer( |
| 841 | &self, |
| 842 | to: SocketAddr, |
| 843 | target: Target, |
| 844 | usn: String, |
| 845 | location: String, |
| 846 | server: String, |
| 847 | ) |
| 848 | -> Outcome<()> |
| 849 | { |
| 850 | let response = SearchResponse { |
| 851 | max_age: DEFAULT_MAX_AGE, |
| 852 | date: Some(res!(http_date())), |
| 853 | location, |
| 854 | server, |
| 855 | target, |
| 856 | usn, |
| 857 | boot_id: None, |
| 858 | config_id: None, |
| 859 | extra: BTreeMap::new(), |
| 860 | }; |
| 861 | res!(self.send_to(&SsdpMessage::Response(response), to)); |
| 862 | Ok(()) |
| 863 | } |
| 864 | } |
| 865 | |
| 866 | /// The local address of the interface that would carry a packet to the SSDP |
| 867 | /// group, which on a machine with one network is the only answer there is. |
| 868 | /// |
| 869 | /// No packet is sent: connecting a UDP socket only chooses a route. |
| 870 | pub fn route_interface() -> Option<Ipv4Addr> { |
| 871 | let socket = match std::net::UdpSocket::bind(SocketAddrV4::new(Ipv4Addr::UNSPECIFIED, 0)) { |
| 872 | Ok(s) => s, |
| 873 | Err(_) => return None, |
| 874 | }; |
| 875 | if socket.connect(SocketAddrV4::new(MULTICAST_ADDR, PORT)).is_err() { |
| 876 | return None; |
| 877 | } |
| 878 | match socket.local_addr() { |
| 879 | Ok(SocketAddr::V4(addr)) => Some(*addr.ip()), |
| 880 | _ => None, |
| 881 | } |
| 882 | } |
| 883 | |
| 884 | /// Every IPv4 address this machine holds, as interfaces to join a group on. |
| 885 | /// |
| 886 | /// The route this machine would use to reach the group is always first, because |
| 887 | /// it is the one that matters and the only one some platforms can be asked for. |
| 888 | /// On Linux the rest are read from the kernel's routing table; elsewhere the list |
| 889 | /// is the route interface and the loopback, which covers a machine with one |
| 890 | /// network and understates a machine with several. |
| 891 | pub fn local_interfaces() -> Vec<Ipv4Addr> { |
| 892 | let mut out: Vec<Ipv4Addr> = Vec::new(); |
| 893 | if let Some(addr) = route_interface() { |
| 894 | out.push(addr); |
| 895 | } |
| 896 | for addr in host_addresses() { |
| 897 | if !out.contains(&addr) { |
| 898 | out.push(addr); |
| 899 | } |
| 900 | } |
| 901 | if !out.contains(&Ipv4Addr::LOCALHOST) { |
| 902 | out.push(Ipv4Addr::LOCALHOST); |
| 903 | } |
| 904 | out |
| 905 | } |
| 906 | |
| 907 | /// Every address the kernel calls local, read from `/proc/net/fib_trie`. |
| 908 | /// |
| 909 | /// The file lists the routing trie, and a machine's own addresses appear in it as |
| 910 | /// a `/32 host LOCAL` route under the address itself. That is a Linux detail and a |
| 911 | /// stable one; on any other platform this answers nothing and |
| 912 | /// [`local_interfaces`] falls back to the route interface alone. |
| 913 | #[cfg(target_os = "linux")] |
| 914 | fn host_addresses() -> Vec<Ipv4Addr> { |
| 915 | let text = match std::fs::read_to_string("/proc/net/fib_trie") { |
| 916 | Ok(t) => t, |
| 917 | Err(_) => return Vec::new(), |
| 918 | }; |
| 919 | let mut out: Vec<Ipv4Addr> = Vec::new(); |
| 920 | let mut candidate: Option<Ipv4Addr> = None; |
| 921 | for line in text.lines() { |
| 922 | let trimmed = line.trim(); |
| 923 | if let Some(rest) = trimmed.strip_prefix("|-- ") { |
| 924 | candidate = Ipv4Addr::from_str(rest.trim()).ok(); |
| 925 | continue; |
| 926 | } |
| 927 | // The line after an address says what kind of route it is. Only a host |
| 928 | // route to the machine itself is an interface address. |
| 929 | if trimmed.starts_with("/32 host LOCAL") { |
| 930 | if let Some(addr) = candidate.take() { |
| 931 | if !addr.is_loopback() && !out.contains(&addr) { |
| 932 | out.push(addr); |
| 933 | } |
| 934 | } |
| 935 | } |
| 936 | } |
| 937 | out |
| 938 | } |
| 939 | |
| 940 | /// Nothing, on a platform whose routing table is not a file. |
| 941 | #[cfg(not(target_os = "linux"))] |
| 942 | fn host_addresses() -> Vec<Ipv4Addr> { |
| 943 | Vec::new() |
| 944 | } |
| 945 | |
| 946 | /// Now, as an HTTP date, for the `DATE` field of an answer. |
| 947 | fn http_date() -> Outcome<String> { |
| 948 | let since = match SystemTime::now().duration_since(UNIX_EPOCH) { |
| 949 | Ok(d) => d, |
| 950 | Err(_) => Duration::from_secs(0), // A clock before the epoch is odd, not fatal. |
| 951 | }; |
| 952 | Ok(res!(Time::fmt_http(&since))) |
| 953 | } |
| 954 | |
| 955 | |
| 956 | #[cfg(test)] |
| 957 | mod tests { |
| 958 | use super::*; |
| 959 | |
| 960 | // A search as a real control point sends one, line endings and all. |
| 961 | const A_SEARCH: &str = "M-SEARCH * HTTP/1.1\r\n\ |
| 962 | HOST: 239.255.255.250:1900\r\n\ |
| 963 | MAN: \"ssdp:discover\"\r\n\ |
| 964 | MX: 3\r\n\ |
| 965 | ST: urn:schemas-upnp-org:device:MediaServer:1\r\n\ |
| 966 | USER-AGENT: Linux/6.1 UPnP/2.0 Probe/1.0\r\n\ |
| 967 | \r\n"; |
| 968 | |
| 969 | const AN_ANSWER: &str = "HTTP/1.1 200 OK\r\n\ |
| 970 | CACHE-CONTROL: max-age=1800\r\n\ |
| 971 | DATE: Mon, 28 Jul 2026 10:00:00 GMT\r\n\ |
| 972 | EXT:\r\n\ |
| 973 | LOCATION: http://192.168.1.10:8200/rootDesc.xml\r\n\ |
| 974 | SERVER: Linux/6.1 UPnP/2.0 Server/1.0\r\n\ |
| 975 | ST: urn:schemas-upnp-org:device:MediaServer:1\r\n\ |
| 976 | USN: uuid:4d696e69-444c-164e-9d41-0011328c0e2f::urn:schemas-upnp-org:device:MediaServer:1\r\n\ |
| 977 | \r\n"; |
| 978 | |
| 979 | const AN_ALIVE: &str = "NOTIFY * HTTP/1.1\r\n\ |
| 980 | HOST: 239.255.255.250:1900\r\n\ |
| 981 | CACHE-CONTROL: max-age=1800\r\n\ |
| 982 | LOCATION: http://192.168.1.10:8200/rootDesc.xml\r\n\ |
| 983 | NT: upnp:rootdevice\r\n\ |
| 984 | NTS: ssdp:alive\r\n\ |
| 985 | SERVER: Linux/6.1 UPnP/2.0 Server/1.0\r\n\ |
| 986 | USN: uuid:4d696e69-444c-164e-9d41-0011328c0e2f::upnp:rootdevice\r\n\ |
| 987 | \r\n"; |
| 988 | |
| 989 | const A_BYEBYE: &str = "NOTIFY * HTTP/1.1\r\n\ |
| 990 | HOST: 239.255.255.250:1900\r\n\ |
| 991 | NT: upnp:rootdevice\r\n\ |
| 992 | NTS: ssdp:byebye\r\n\ |
| 993 | USN: uuid:4d696e69-444c-164e-9d41-0011328c0e2f::upnp:rootdevice\r\n\ |
| 994 | \r\n"; |
| 995 | |
| 996 | // ┌───────────────────────────────────────────────────────────────────────┐ |
| 997 | // │ TARGETS │ |
| 998 | // └───────────────────────────────────────────────────────────────────────┘ |
| 999 | |
| 1000 | #[test] |
| 1001 | fn test_a_target_survives_being_written_and_read() -> Outcome<()> { |
| 1002 | let targets = [ |
| 1003 | Target::All, |
| 1004 | Target::RootDevice, |
| 1005 | Target::Uuid("4d696e69-444c-164e-9d41-0011328c0e2f".to_string()), |
| 1006 | Target::Urn("schemas-upnp-org:device:MediaServer:1".to_string()), |
| 1007 | Target::Other("some-vendor-thing".to_string()), |
| 1008 | ]; |
| 1009 | for target in targets { |
| 1010 | let written = fmt!("{}", target); |
| 1011 | assert_eq!(res!(Target::from_str(&written)), target, |
| 1012 | "{:?} did not survive being written as {:?}", target, written); |
| 1013 | } |
| 1014 | Ok(()) |
| 1015 | } |
| 1016 | |
| 1017 | /// `ssdp:all` is answered by everything; everything else by itself alone. A |
| 1018 | /// responder that gets this wrong announces a service nobody asked about, and |
| 1019 | /// control points ignore the answer. |
| 1020 | #[test] |
| 1021 | fn test_who_answers_whom() { |
| 1022 | let root = Target::RootDevice; |
| 1023 | let server = Target::Urn("schemas-upnp-org:device:MediaServer:1".to_string()); |
| 1024 | assert!(root.answers(&Target::All)); |
| 1025 | assert!(server.answers(&Target::All)); |
| 1026 | assert!(root.answers(&Target::RootDevice)); |
| 1027 | assert!(!root.answers(&server)); |
| 1028 | assert!(!server.answers(&Target::RootDevice)); |
| 1029 | } |
| 1030 | |
| 1031 | // ┌───────────────────────────────────────────────────────────────────────┐ |
| 1032 | // │ PARSING │ |
| 1033 | // └───────────────────────────────────────────────────────────────────────┘ |
| 1034 | |
| 1035 | #[test] |
| 1036 | fn test_a_search_is_read() -> Outcome<()> { |
| 1037 | match res!(SsdpMessage::parse(A_SEARCH.as_bytes())) { |
| 1038 | SsdpMessage::Search(s) => { |
| 1039 | assert_eq!(s.target, |
| 1040 | Target::Urn("schemas-upnp-org:device:MediaServer:1".to_string())); |
| 1041 | assert_eq!(s.mx, 3); |
| 1042 | assert_eq!(s.user_agent.as_deref(), Some("Linux/6.1 UPnP/2.0 Probe/1.0")); |
| 1043 | assert!(s.extra.is_empty(), "unexpected extras: {:?}", s.extra); |
| 1044 | } |
| 1045 | other => return Err(err!("A search was read as {:?}.", other; Test, Mismatch)), |
| 1046 | } |
| 1047 | Ok(()) |
| 1048 | } |
| 1049 | |
| 1050 | #[test] |
| 1051 | fn test_an_answer_is_read() -> Outcome<()> { |
| 1052 | match res!(SsdpMessage::parse(AN_ANSWER.as_bytes())) { |
| 1053 | SsdpMessage::Response(r) => { |
| 1054 | assert_eq!(r.max_age, 1800); |
| 1055 | assert_eq!(r.location, "http://192.168.1.10:8200/rootDesc.xml"); |
| 1056 | assert_eq!(r.target, |
| 1057 | Target::Urn("schemas-upnp-org:device:MediaServer:1".to_string())); |
| 1058 | assert!(r.usn.starts_with("uuid:4d696e69")); |
| 1059 | assert_eq!(r.date.as_deref(), Some("Mon, 28 Jul 2026 10:00:00 GMT")); |
| 1060 | } |
| 1061 | other => return Err(err!("An answer was read as {:?}.", other; Test, Mismatch)), |
| 1062 | } |
| 1063 | Ok(()) |
| 1064 | } |
| 1065 | |
| 1066 | #[test] |
| 1067 | fn test_an_alive_and_a_byebye_are_read() -> Outcome<()> { |
| 1068 | match res!(SsdpMessage::parse(AN_ALIVE.as_bytes())) { |
| 1069 | SsdpMessage::Notify(n) => { |
| 1070 | assert_eq!(n.nts, Nts::Alive); |
| 1071 | assert_eq!(n.target, Target::RootDevice); |
| 1072 | assert_eq!(n.max_age, Some(1800)); |
| 1073 | assert_eq!(n.location.as_deref(), |
| 1074 | Some("http://192.168.1.10:8200/rootDesc.xml")); |
| 1075 | } |
| 1076 | other => return Err(err!("An alive was read as {:?}.", other; Test, Mismatch)), |
| 1077 | } |
| 1078 | // A byebye carries neither a lifetime nor a location: there is nothing |
| 1079 | // left to fetch. |
| 1080 | match res!(SsdpMessage::parse(A_BYEBYE.as_bytes())) { |
| 1081 | SsdpMessage::Notify(n) => { |
| 1082 | assert_eq!(n.nts, Nts::ByeBye); |
| 1083 | assert_eq!(n.max_age, None); |
| 1084 | assert_eq!(n.location, None); |
| 1085 | } |
| 1086 | other => return Err(err!("A byebye was read as {:?}.", other; Test, Mismatch)), |
| 1087 | } |
| 1088 | Ok(()) |
| 1089 | } |
| 1090 | |
| 1091 | /// Devices on a real network get the details wrong. Field names in any case, |
| 1092 | /// bare newlines instead of CRLF, and spacing around the colon are all things |
| 1093 | /// that must not stop a message being understood. |
| 1094 | #[test] |
| 1095 | fn test_a_sloppy_sender_is_still_understood() -> Outcome<()> { |
| 1096 | let sloppy = "m-search * HTTP/1.1\n\ |
| 1097 | host:239.255.255.250:1900\n\ |
| 1098 | Man: \"ssdp:discover\"\n\ |
| 1099 | mx: 1\n\ |
| 1100 | st:ssdp:all\n\ |
| 1101 | \n"; |
| 1102 | match res!(SsdpMessage::parse(sloppy.as_bytes())) { |
| 1103 | SsdpMessage::Search(s) => { |
| 1104 | assert_eq!(s.target, Target::All); |
| 1105 | assert_eq!(s.mx, 1); |
| 1106 | } |
| 1107 | other => return Err(err!("A sloppy search was read as {:?}.", other; Test, Mismatch)), |
| 1108 | } |
| 1109 | Ok(()) |
| 1110 | } |
| 1111 | |
| 1112 | /// A field this crate does not model is kept rather than dropped, so a |
| 1113 | /// responder can pass a vendor extension through and a reader can see it. |
| 1114 | #[test] |
| 1115 | fn test_an_unmodelled_field_is_kept() -> Outcome<()> { |
| 1116 | let with_extra = "NOTIFY * HTTP/1.1\r\n\ |
| 1117 | HOST: 239.255.255.250:1900\r\n\ |
| 1118 | NT: upnp:rootdevice\r\n\ |
| 1119 | NTS: ssdp:alive\r\n\ |
| 1120 | USN: uuid:x::upnp:rootdevice\r\n\ |
| 1121 | X-VENDOR-THING: 42\r\n\ |
| 1122 | \r\n"; |
| 1123 | match res!(SsdpMessage::parse(with_extra.as_bytes())) { |
| 1124 | SsdpMessage::Notify(n) => { |
| 1125 | assert_eq!(n.extra.get("X-VENDOR-THING").map(String::as_str), Some("42")); |
| 1126 | // And it goes back out again. |
| 1127 | assert!(n.as_text().contains("X-VENDOR-THING: 42\r\n")); |
| 1128 | } |
| 1129 | other => return Err(err!("A notify was read as {:?}.", other; Test, Mismatch)), |
| 1130 | } |
| 1131 | Ok(()) |
| 1132 | } |
| 1133 | |
| 1134 | #[test] |
| 1135 | fn test_what_is_not_an_ssdp_message_is_refused() { |
| 1136 | for bad in [ |
| 1137 | "", // Nothing at all. |
| 1138 | "GET / HTTP/1.1\r\nHost: x\r\n\r\n", // HTTP, but not SSDP. |
| 1139 | "M-SEARCH * HTTP/1.1\r\nMX: 1\r\n\r\n", // A search naming nothing. |
| 1140 | "NOTIFY * HTTP/1.1\r\nNT: upnp:rootdevice\r\n\r\n", // No NTS, no USN. |
| 1141 | "HTTP/1.1 404 Not Found\r\n\r\n", // A refusal is not a discovery. |
| 1142 | ] { |
| 1143 | assert!(SsdpMessage::parse(bad.as_bytes()).is_err(), |
| 1144 | "{:?} should not have parsed", bad); |
| 1145 | } |
| 1146 | } |
| 1147 | |
| 1148 | // ┌───────────────────────────────────────────────────────────────────────┐ |
| 1149 | // │ SERIALISATION │ |
| 1150 | // └───────────────────────────────────────────────────────────────────────┘ |
| 1151 | |
| 1152 | /// The bytes on the wire, compared with a message from a real device rather |
| 1153 | /// than with what this module happens to produce. |
| 1154 | #[test] |
| 1155 | fn test_a_search_goes_out_as_a_search() { |
| 1156 | let mut search = Search::new( |
| 1157 | Target::Urn("schemas-upnp-org:device:MediaServer:1".to_string())); |
| 1158 | search.mx = 3; |
| 1159 | search.user_agent = Some("Linux/6.1 UPnP/2.0 Probe/1.0".to_string()); |
| 1160 | assert_eq!(search.as_text(), A_SEARCH); |
| 1161 | } |
| 1162 | |
| 1163 | /// `EXT:` is empty, means nothing, and is required: a control point that does |
| 1164 | /// not see it discards the answer (UPnP DA 2.0 §1.3.3). |
| 1165 | #[test] |
| 1166 | fn test_an_answer_carries_the_empty_ext_field() { |
| 1167 | let answer = SearchResponse { |
| 1168 | max_age: 1800, |
| 1169 | date: Some("Mon, 28 Jul 2026 10:00:00 GMT".to_string()), |
| 1170 | location: "http://192.168.1.10:8200/rootDesc.xml".to_string(), |
| 1171 | server: "Linux/6.1 UPnP/2.0 Server/1.0".to_string(), |
| 1172 | target: Target::Urn("schemas-upnp-org:device:MediaServer:1".to_string()), |
| 1173 | usn: "uuid:4d696e69-444c-164e-9d41-0011328c0e2f::\ |
| 1174 | urn:schemas-upnp-org:device:MediaServer:1".to_string(), |
| 1175 | boot_id: None, |
| 1176 | config_id: None, |
| 1177 | extra: BTreeMap::new(), |
| 1178 | }; |
| 1179 | assert_eq!(answer.as_text(), AN_ANSWER); |
| 1180 | } |
| 1181 | |
| 1182 | #[test] |
| 1183 | fn test_an_alive_goes_out_as_an_alive() { |
| 1184 | let alive = Notify::alive( |
| 1185 | Target::RootDevice, |
| 1186 | "uuid:4d696e69-444c-164e-9d41-0011328c0e2f::upnp:rootdevice".to_string(), |
| 1187 | "http://192.168.1.10:8200/rootDesc.xml".to_string(), |
| 1188 | "Linux/6.1 UPnP/2.0 Server/1.0".to_string(), |
| 1189 | ); |
| 1190 | assert_eq!(alive.as_text(), AN_ALIVE); |
| 1191 | } |
| 1192 | |
| 1193 | #[test] |
| 1194 | fn test_a_byebye_promises_nothing() { |
| 1195 | let bye = Notify::byebye( |
| 1196 | Target::RootDevice, |
| 1197 | "uuid:4d696e69-444c-164e-9d41-0011328c0e2f::upnp:rootdevice".to_string(), |
| 1198 | ); |
| 1199 | assert_eq!(bye.as_text(), A_BYEBYE); |
| 1200 | } |
| 1201 | |
| 1202 | /// Everything this module writes, it can read back. |
| 1203 | #[test] |
| 1204 | fn test_every_message_survives_a_round_trip() -> Outcome<()> { |
| 1205 | for text in [A_SEARCH, AN_ANSWER, AN_ALIVE, A_BYEBYE] { |
| 1206 | let msg = res!(SsdpMessage::parse(text.as_bytes())); |
| 1207 | let again = res!(SsdpMessage::parse(&msg.as_bytes())); |
| 1208 | assert_eq!(msg, again, "a message changed on the way round: {:?}", text); |
| 1209 | } |
| 1210 | Ok(()) |
| 1211 | } |
| 1212 | |
| 1213 | #[test] |
| 1214 | fn test_a_cache_control_yields_its_seconds() { |
| 1215 | assert_eq!(max_age_of("max-age=1800"), Some(1800)); |
| 1216 | assert_eq!(max_age_of("max-age = 60"), Some(60)); |
| 1217 | assert_eq!(max_age_of("public, max-age=120"), Some(120)); |
| 1218 | assert_eq!(max_age_of("no-cache"), None); |
| 1219 | } |
| 1220 | |
| 1221 | // ┌───────────────────────────────────────────────────────────────────────┐ |
| 1222 | // │ THE SOCKET │ |
| 1223 | // └───────────────────────────────────────────────────────────────────────┘ |
| 1224 | |
| 1225 | /// The group is joined and a message goes round it, on the loopback interface |
| 1226 | /// and an ephemeral port so nothing on the machine is disturbed. Live |
| 1227 | /// behaviour on a real network is not tested here: it wants a second machine. |
| 1228 | #[test] |
| 1229 | fn test_a_responder_binds_joins_and_carries_a_message() -> Outcome<()> { |
| 1230 | let rt = res!(tokio::runtime::Runtime::new()); |
| 1231 | rt.block_on(async { |
| 1232 | let responder = match Responder::bind_to_port(Ipv4Addr::LOCALHOST, 0).await { |
| 1233 | Ok(r) => r, |
| 1234 | Err(e) => { |
| 1235 | // A machine with no loopback multicast cannot run this, and |
| 1236 | // that is the machine's business rather than a failure of the |
| 1237 | // code under test. |
| 1238 | warn!("SSDP loopback multicast is unavailable here: {}", e); |
| 1239 | return Ok(()); |
| 1240 | } |
| 1241 | }; |
| 1242 | let addr = res!(responder.local_addr()); |
| 1243 | let alive = SsdpMessage::Notify(Notify::alive( |
| 1244 | Target::RootDevice, |
| 1245 | "uuid:test::upnp:rootdevice".to_string(), |
| 1246 | "http://127.0.0.1:8200/rootDesc.xml".to_string(), |
| 1247 | "Test/1.0 UPnP/2.0 Test/1.0".to_string(), |
| 1248 | )); |
| 1249 | res!(responder.send_to(&alive, addr).await); |
| 1250 | let got = tokio::time::timeout( |
| 1251 | Duration::from_secs(2), |
| 1252 | responder.recv(), |
| 1253 | ).await; |
| 1254 | match got { |
| 1255 | Ok(Ok((msg, _from))) => assert_eq!(msg, alive), |
| 1256 | Ok(Err(e)) => return Err(e), |
| 1257 | Err(_) => return Err(err!( |
| 1258 | "The responder never heard its own announcement."; Test, Timeout)), |
| 1259 | } |
| 1260 | Ok(()) |
| 1261 | }) |
| 1262 | } |
| 1263 | |
| 1264 | /// The same, with no runtime under it: a search goes out, comes back, and is |
| 1265 | /// answered to the address it came from. |
| 1266 | #[test] |
| 1267 | fn test_a_blocking_responder_carries_a_search_and_its_answer() -> Outcome<()> { |
| 1268 | let mut responder = res!(SyncResponder::bind_to_port(0)); |
| 1269 | match responder.join(Ipv4Addr::LOCALHOST) { |
| 1270 | Ok(()) => {}, |
| 1271 | Err(e) => { |
| 1272 | // A machine with no loopback multicast cannot run this, and that |
| 1273 | // is the machine's business rather than a failure of the code. |
| 1274 | warn!("SSDP loopback multicast is unavailable here: {}", e); |
| 1275 | return Ok(()); |
| 1276 | }, |
| 1277 | } |
| 1278 | res!(responder.set_timeout(Some(Duration::from_secs(2)))); |
| 1279 | let addr = res!(responder.local_addr()); |
| 1280 | |
| 1281 | let search = SsdpMessage::Search(Search::new( |
| 1282 | Target::Urn("schemas-upnp-org:device:MediaServer:1".to_string()))); |
| 1283 | res!(responder.send_to(&search, addr)); |
| 1284 | let (msg, from) = match res!(responder.recv()) { |
| 1285 | Some(pair) => pair, |
| 1286 | None => return Err(err!( |
| 1287 | "The responder never heard the search."; Test, Timeout)), |
| 1288 | }; |
| 1289 | req!(msg, search); |
| 1290 | |
| 1291 | // And the answer goes back to the asker, echoing the target that was |
| 1292 | // asked for rather than the one that was searched under. |
| 1293 | res!(responder.answer( |
| 1294 | from, |
| 1295 | Target::Urn("schemas-upnp-org:device:MediaServer:1".to_string()), |
| 1296 | "uuid:test::urn:schemas-upnp-org:device:MediaServer:1".to_string(), |
| 1297 | "http://127.0.0.1:8336/dlna/desc.xml".to_string(), |
| 1298 | "Test/1.0 UPnP/1.0 Ochre/0.1".to_string(), |
| 1299 | )); |
| 1300 | match res!(responder.recv()) { |
| 1301 | Some((SsdpMessage::Response(r), _)) => { |
| 1302 | req!(r.location, "http://127.0.0.1:8336/dlna/desc.xml".to_string()); |
| 1303 | req!(r.target, |
| 1304 | Target::Urn("schemas-upnp-org:device:MediaServer:1".to_string())); |
| 1305 | assert!(r.usn.ends_with(&fmt!("::{}", r.target)), |
| 1306 | "the answer's USN and ST disagree: {} against {}", r.usn, r.target); |
| 1307 | }, |
| 1308 | other => return Err(err!( |
| 1309 | "The answer came back as {:?}.", other; Test, Mismatch)), |
| 1310 | } |
| 1311 | Ok(()) |
| 1312 | } |
| 1313 | |
| 1314 | /// A read that finds nothing is not a failure, and is what lets a serving |
| 1315 | /// thread look at its shutdown flag. |
| 1316 | #[test] |
| 1317 | fn test_a_blocking_read_that_finds_nothing_says_so() -> Outcome<()> { |
| 1318 | let responder = res!(SyncResponder::bind_to_port(0)); |
| 1319 | res!(responder.set_timeout(Some(Duration::from_millis(200)))); |
| 1320 | req!(res!(responder.recv()).is_none(), true); |
| 1321 | Ok(()) |
| 1322 | } |
| 1323 | |
| 1324 | /// The machine is asked what interfaces it has, and every answer is an |
| 1325 | /// address a group can actually be joined on. |
| 1326 | #[test] |
| 1327 | fn test_the_interfaces_offered_are_addresses_of_this_machine() -> Outcome<()> { |
| 1328 | let ifaces = local_interfaces(); |
| 1329 | assert!(ifaces.contains(&Ipv4Addr::LOCALHOST), |
| 1330 | "the loopback is always usable and was not offered: {:?}", ifaces); |
| 1331 | for addr in &ifaces { |
| 1332 | assert!(!addr.is_multicast(), "{} is a group, not an interface", addr); |
| 1333 | assert!(!addr.is_unspecified(), "0.0.0.0 is not an interface"); |
| 1334 | } |
| 1335 | Ok(()) |
| 1336 | } |
| 1337 | } |