1use std::net::IpAddr;
18use std::sync::Arc;
19use std::sync::atomic::{AtomicBool, AtomicU16, AtomicU64, Ordering};
20use std::thread;
21use std::time::{Duration, Instant};
22
23use crossbeam_channel::{Receiver, RecvTimeoutError, Sender};
24use parking_lot::Mutex;
25
26use super::description::{Renderer, Service};
27use super::serve::{Listener, Served};
28use super::soap::{self, SoapError, arg};
29use super::{didl, xml};
30
31const SUBSCRIPTION_SECONDS: u64 = 1800;
32const INITIAL_EVENT: Duration = Duration::from_secs(3);
35const POLL: Duration = Duration::from_secs(1);
36const UNKNOWN_VOLUME: u16 = u16::MAX;
37
38#[derive(Debug, Clone, Copy, PartialEq, Eq)]
39pub enum Transport {
40 Playing,
41 Paused,
42 Stopped,
43 Transitioning,
45 NoMedia,
46}
47
48impl Transport {
49 fn parse(s: &str) -> Self {
50 match s.trim() {
51 "PLAYING" => Self::Playing,
52 "PAUSED_PLAYBACK" | "PAUSED_RECORDING" => Self::Paused,
53 "TRANSITIONING" => Self::Transitioning,
54 "NO_MEDIA_PRESENT" => Self::NoMedia,
55 _ => Self::Stopped,
56 }
57 }
58}
59
60#[derive(Debug, Clone, PartialEq)]
62pub struct Snapshot {
63 pub epoch: u64,
64 pub transport: Transport,
65 pub track_uri: String,
66 pub position_ms: Option<u64>,
67 pub duration_ms: Option<u64>,
68 pub at: Instant,
70}
71
72#[derive(Debug, Clone)]
73pub enum Event {
74 Snapshot(Snapshot),
75 Volume(u8),
76 Gone,
79}
80
81struct Subscription {
82 service: Service,
83 sid: String,
84 renew_at: Instant,
85}
86
87struct Shared {
88 renderer: Renderer,
89 http: reqwest::blocking::Client,
90 epoch: AtomicU64,
91 look: Sender<()>,
92 polling: AtomicBool,
93 playing: AtomicBool,
94 volume: Arc<AtomicU16>,
96 subscriptions: Mutex<Vec<Subscription>>,
97 trust: Arc<Trust>,
98 callback: String,
99 closed: AtomicBool,
100 unreachable: AtomicBool,
101 on_event: Arc<dyn Fn(Event) + Send + Sync>,
102}
103
104#[derive(Default)]
109struct Trust {
110 sids: Mutex<Vec<String>>,
111 subscribing: AtomicBool,
112}
113
114pub struct Session {
115 shared: Arc<Shared>,
116 listener: Listener,
117 base: String,
119 sink: Vec<Option<String>>,
120 stop: Sender<()>,
121 gone: u64,
122}
123
124impl Session {
125 pub fn open(
128 renderer: Renderer,
129 on_event: impl Fn(Event) + Send + Sync + 'static,
130 ) -> Result<Self, String> {
131 let host: IpAddr = renderer
132 .location
133 .host_str()
134 .and_then(|h| h.trim_matches(['[', ']']).parse().ok())
135 .ok_or_else(|| format!("{} has no address", renderer.name))?;
136 let ip = super::serve::local_ip_towards(host)
137 .ok_or_else(|| format!("no route to {}", renderer.name))?;
138
139 let on_event: Arc<dyn Fn(Event) + Send + Sync> = Arc::new(on_event);
140 let (look_tx, look_rx) = crossbeam_channel::bounded::<()>(1);
141 let seen = Arc::new(AtomicBool::new(false));
142 let trust = Arc::new(Trust::default());
143 let volume = Arc::new(AtomicU16::new(UNKNOWN_VOLUME));
144 let listener = {
145 let look = look_tx.clone();
146 let on_event = on_event.clone();
147 let seen = seen.clone();
148 let trust = trust.clone();
149 let heard = volume.clone();
150 Listener::start(Box::new(move |sid, body| {
151 let known = trust.sids.lock().iter().any(|s| s == sid);
152 if !known && !trust.subscribing.load(Ordering::Acquire) {
153 return;
154 }
155 seen.store(true, Ordering::Release);
156 let change = parse_notify(&body);
157 if let Some(volume) = change.volume.filter(|_| known) {
158 heard.store(u16::from(volume), Ordering::Release);
159 on_event(Event::Volume(volume));
160 }
161 if change.transport {
162 let _ = look.try_send(());
163 }
164 }))
165 .map_err(|e| format!("cannot listen for {}: {e}", renderer.name))?
166 };
167 let base = match ip {
168 IpAddr::V6(v6) => format!("http://[{v6}]:{}", listener.port()),
169 IpAddr::V4(v4) => format!("http://{v4}:{}", listener.port()),
170 };
171
172 let http = soap::client();
173 let sink = renderer
174 .connection_manager
175 .as_ref()
176 .and_then(|cm| soap::call(&http, cm, "GetProtocolInfo", &[]).ok())
177 .and_then(|args| arg(&args, "Sink").map(didl::sink_mimes))
178 .unwrap_or_default();
179 log::info!(
180 "upnp: {} plays {}",
181 renderer.name,
182 if sink.is_empty() {
183 "anything it is given (no protocol info)".to_string()
184 } else {
185 sink.iter()
186 .map(|m| m.as_deref().unwrap_or("*"))
187 .collect::<Vec<_>>()
188 .join(", ")
189 }
190 );
191
192 let shared = Arc::new(Shared {
193 renderer,
194 http,
195 epoch: AtomicU64::new(0),
196 look: look_tx,
197 polling: AtomicBool::new(false),
198 playing: AtomicBool::new(false),
199 volume: volume.clone(),
200 subscriptions: Mutex::new(Vec::new()),
201 trust,
202 callback: format!("<{base}/events/>"),
203 closed: AtomicBool::new(false),
204 unreachable: AtomicBool::new(false),
205 on_event: on_event.clone(),
206 });
207
208 let services: Vec<Service> = std::iter::once(shared.renderer.av_transport.clone())
209 .chain(shared.renderer.rendering_control.clone())
210 .collect();
211 shared.trust.subscribing.store(true, Ordering::Release);
212 for (i, service) in services.into_iter().enumerate() {
213 match shared.subscribe(&service) {
214 Ok(sub) => shared.subscriptions.lock().push(sub),
215 Err(e) => {
216 log::info!("upnp: {} refused a subscription: {e}", shared.renderer.name);
217 if i == 0 {
218 shared.polling.store(true, Ordering::Release);
219 }
220 }
221 }
222 }
223 shared.trust.subscribing.store(false, Ordering::Release);
224
225 let gone = {
228 let weak = Arc::downgrade(&shared);
229 super::discovery::on_gone(&shared.renderer.udn, move || {
230 if let Some(shared) = weak.upgrade() {
231 shared.lost();
232 }
233 })
234 };
235
236 if let Some(v) = shared.volume() {
239 volume.store(u16::from(v), Ordering::Release);
240 }
241
242 let (stop_tx, stop_rx) = crossbeam_channel::bounded::<()>(1);
243 spawn_watcher(shared.clone(), look_rx, seen, on_event, stop_rx.clone());
244 spawn_renewer(shared.clone(), stop_rx);
245
246 Ok(Self {
247 shared,
248 listener,
249 base,
250 sink,
251 stop: stop_tx,
252 gone,
253 })
254 }
255
256 pub fn renderer(&self) -> &Renderer {
257 &self.shared.renderer
258 }
259
260 pub fn mime_for(&self, extension: &str) -> Option<&'static str> {
263 didl::choose_mime(extension, &self.sink)
264 }
265
266 pub fn serve(
268 &self,
269 path: &std::path::Path,
270 mime: &str,
271 extension: &str,
272 ) -> (String, String, String) {
273 self.offer(
274 Served::File {
275 path: path.to_path_buf(),
276 mime: mime.to_string(),
277 },
278 extension,
279 )
280 }
281
282 pub fn serve_stream(
285 &self,
286 pipe: Arc<super::stream::Pipe>,
287 art: &std::path::Path,
288 extension: &str,
289 ) -> (String, String, String) {
290 self.offer(
291 Served::Stream {
292 pipe,
293 art: art.to_path_buf(),
294 },
295 extension,
296 )
297 }
298
299 fn offer(&self, served: Served, extension: &str) -> (String, String, String) {
300 let token = self.listener.add(served);
301 let url = format!("{}/t/{token}.{extension}", self.base);
302 let art = format!("{}/art/{token}", self.base);
303 (token, url, art)
304 }
305
306 pub fn retain(&self, tokens: &[&str]) {
308 self.listener.retain(tokens);
309 }
310
311 pub fn token_of<'a>(&self, uri: &'a str) -> Option<&'a str> {
313 let rest = uri.strip_prefix(&self.base)?.strip_prefix("/t/")?;
314 rest.split('.').next()
315 }
316
317 pub fn is_lost(&self) -> bool {
320 self.shared.unreachable.load(Ordering::Acquire)
321 }
322
323 pub fn epoch(&self) -> u64 {
324 self.shared.epoch.load(Ordering::Acquire)
325 }
326
327 pub fn look(&self) {
329 let _ = self.shared.look.try_send(());
330 }
331
332 pub fn set_uri(&self, uri: &str, metadata: &str) -> Result<(), SoapError> {
333 self.transport(
334 "SetAVTransportURI",
335 &[("CurrentURI", uri), ("CurrentURIMetaData", metadata)],
336 )
337 }
338
339 pub fn set_next(&self, uri: &str, metadata: &str) -> Result<(), SoapError> {
340 self.transport(
341 "SetNextAVTransportURI",
342 &[("NextURI", uri), ("NextURIMetaData", metadata)],
343 )
344 }
345
346 pub fn play(&self) -> Result<(), SoapError> {
347 self.shared.playing.store(true, Ordering::Release);
348 self.transport("Play", &[("Speed", "1")])
349 }
350
351 pub fn pause(&self) -> Result<(), SoapError> {
352 self.shared.playing.store(false, Ordering::Release);
353 self.transport("Pause", &[])
354 }
355
356 pub fn stop(&self) -> Result<(), SoapError> {
357 self.shared.playing.store(false, Ordering::Release);
358 self.transport("Stop", &[])
359 }
360
361 pub fn seek(&self, position_ms: u64) -> Result<(), SoapError> {
362 self.transport(
363 "Seek",
364 &[
365 ("Unit", "REL_TIME"),
366 ("Target", &soap::format_time(position_ms)),
367 ],
368 )
369 }
370
371 pub fn volume(&self) -> Option<u8> {
374 match self.shared.volume.load(Ordering::Acquire) {
375 UNKNOWN_VOLUME => None,
376 v => Some(v as u8),
377 }
378 }
379
380 pub fn set_volume(&self, volume: u8) -> Result<(), SoapError> {
381 let Some(rc) = &self.shared.renderer.rendering_control else {
382 return Ok(());
383 };
384 let volume = volume.min(100);
385 self.shared
386 .volume
387 .store(u16::from(volume), Ordering::Release);
388 soap::call(
389 &self.shared.http,
390 rc,
391 "SetVolume",
392 &[
393 ("InstanceID", "0"),
394 ("Channel", "Master"),
395 ("DesiredVolume", &volume.min(100).to_string()),
396 ],
397 )
398 .map(|_| ())
399 }
400
401 pub fn has_volume(&self) -> bool {
402 self.shared.renderer.rendering_control.is_some()
403 }
404
405 fn transport(&self, action: &str, args: &[(&str, &str)]) -> Result<(), SoapError> {
412 if self.shared.unreachable.load(Ordering::Acquire) {
413 return Err(SoapError::Unreachable(format!(
414 "{} stopped answering",
415 self.shared.renderer.name
416 )));
417 }
418 let mut full = vec![("InstanceID", "0")];
419 full.extend_from_slice(args);
420 let result = soap::call(
421 &self.shared.http,
422 &self.shared.renderer.av_transport,
423 action,
424 &full,
425 );
426 self.shared.epoch.fetch_add(1, Ordering::AcqRel);
427 if let Err(e) = &result {
428 log::warn!("upnp: {}: {e}", self.shared.renderer.name);
429 if matches!(e, SoapError::Unreachable(_)) {
430 self.shared.lost();
431 }
432 }
433 result.map(|_| ())
434 }
435}
436
437impl std::fmt::Debug for Session {
438 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
439 f.debug_struct("Session")
440 .field("renderer", &self.shared.renderer.name)
441 .field("base", &self.base)
442 .finish()
443 }
444}
445
446impl Drop for Session {
447 fn drop(&mut self) {
448 self.shared.closed.store(true, Ordering::Release);
449 let _ = self.stop.try_send(());
450 super::discovery::forget_gone(self.gone);
451 let shared = self.shared.clone();
454 let _ = thread::Builder::new()
455 .name("koan-upnp-close".into())
456 .spawn(move || {
457 for sub in shared.subscriptions.lock().drain(..) {
458 let _ = shared
459 .http
460 .request(method("UNSUBSCRIBE"), sub.service.events.clone())
461 .header("SID", &sub.sid)
462 .send();
463 }
464 });
465 }
466}
467
468fn method(name: &str) -> reqwest::Method {
469 reqwest::Method::from_bytes(name.as_bytes()).expect("valid method")
470}
471
472impl Shared {
473 fn lost(&self) {
475 if !self.unreachable.swap(true, Ordering::AcqRel) && !self.closed.load(Ordering::Acquire) {
476 log::info!("upnp: {} is gone", self.renderer.name);
477 (self.on_event)(Event::Gone);
478 }
479 }
480
481 fn subscribe(&self, service: &Service) -> Result<Subscription, String> {
482 let response = self
483 .http
484 .request(method("SUBSCRIBE"), service.events.clone())
485 .header("CALLBACK", &self.callback)
486 .header("NT", "upnp:event")
487 .header("TIMEOUT", format!("Second-{SUBSCRIPTION_SECONDS}"))
488 .send()
489 .and_then(|r| r.error_for_status())
490 .map_err(|e| e.to_string())?;
491 let sub = subscription(service, &response)?;
492 self.trust.sids.lock().push(sub.sid.clone());
493 Ok(sub)
494 }
495
496 fn renew(&self, sub: &Subscription) -> Result<Subscription, String> {
497 let response = self
498 .http
499 .request(method("SUBSCRIBE"), sub.service.events.clone())
500 .header("SID", &sub.sid)
501 .header("TIMEOUT", format!("Second-{SUBSCRIPTION_SECONDS}"))
502 .send()
503 .and_then(|r| r.error_for_status())
504 .map_err(|e| e.to_string())?;
505 let renewed = subscription(&sub.service, &response)?;
506 self.trust.sids.lock().push(renewed.sid.clone());
507 Ok(renewed)
508 }
509
510 fn snapshot(&self) -> Result<Snapshot, SoapError> {
511 let epoch = self.epoch.load(Ordering::Acquire);
512 let avt = &self.renderer.av_transport;
513 let id = [("InstanceID", "0")];
514 let info = soap::call(&self.http, avt, "GetTransportInfo", &id)?;
515 let position = soap::call(&self.http, avt, "GetPositionInfo", &id)?;
516 Ok(Snapshot {
517 epoch,
518 transport: Transport::parse(arg(&info, "CurrentTransportState").unwrap_or_default()),
519 track_uri: arg(&position, "TrackURI").unwrap_or_default().to_string(),
520 position_ms: arg(&position, "RelTime").and_then(soap::parse_time),
521 duration_ms: arg(&position, "TrackDuration")
522 .and_then(soap::parse_time)
523 .filter(|d| *d > 0),
524 at: Instant::now(),
525 })
526 }
527
528 fn volume(&self) -> Option<u8> {
529 let rc = self.renderer.rendering_control.as_ref()?;
530 let args = soap::call(
531 &self.http,
532 rc,
533 "GetVolume",
534 &[("InstanceID", "0"), ("Channel", "Master")],
535 )
536 .ok()?;
537 arg(&args, "CurrentVolume")?
538 .parse::<u32>()
539 .ok()
540 .map(|v| v.min(100) as u8)
541 }
542}
543
544fn subscription(
545 service: &Service,
546 response: &reqwest::blocking::Response,
547) -> Result<Subscription, String> {
548 let header = |name: &str| {
549 response
550 .headers()
551 .get(name)
552 .and_then(|v| v.to_str().ok())
553 .map(str::to_string)
554 };
555 let sid = header("SID").ok_or("no SID in the reply")?;
556 let seconds = header("TIMEOUT")
557 .and_then(|t| {
558 t.trim()
559 .strip_prefix("Second-")
560 .and_then(|s| s.parse().ok())
561 })
562 .unwrap_or(SUBSCRIPTION_SECONDS)
563 .max(30);
564 Ok(Subscription {
565 service: service.clone(),
566 sid,
567 renew_at: Instant::now() + Duration::from_secs(seconds / 2),
570 })
571}
572
573fn spawn_watcher(
574 shared: Arc<Shared>,
575 look: Receiver<()>,
576 seen: Arc<AtomicBool>,
577 on_event: Arc<dyn Fn(Event) + Send + Sync>,
578 stop: Receiver<()>,
579) {
580 let spawned = thread::Builder::new()
581 .name("koan-upnp-watch".into())
582 .spawn(move || {
583 let started = Instant::now();
584 loop {
585 let polling = shared.polling.load(Ordering::Acquire);
586 let wait = if polling {
587 shared.playing.load(Ordering::Acquire).then_some(POLL)
588 } else if !seen.load(Ordering::Acquire) {
589 Some(INITIAL_EVENT.saturating_sub(started.elapsed()))
590 } else {
591 None
592 };
593 let woke = crossbeam_channel::select! {
594 recv(look) -> r => r.is_ok(),
595 recv(stop) -> _ => return,
596 default(wait.unwrap_or(Duration::from_secs(3600))) => wait.is_some(),
597 };
598 if shared.closed.load(Ordering::Acquire) {
599 return;
600 }
601 if !woke {
602 continue;
603 }
604 if !polling && !seen.load(Ordering::Acquire) && started.elapsed() >= INITIAL_EVENT {
605 log::info!(
606 "upnp: {} sent no events; asking it once a second while playing",
607 shared.renderer.name
608 );
609 shared.polling.store(true, Ordering::Release);
610 }
611 match shared.snapshot() {
612 Ok(snapshot) => on_event(Event::Snapshot(snapshot)),
613 Err(e) => log::debug!("upnp: {}: {e}", shared.renderer.name),
614 }
615 }
616 });
617 if let Err(e) = spawned {
618 log::warn!("upnp: could not start the watcher: {e}");
619 }
620}
621
622fn spawn_renewer(shared: Arc<Shared>, stop: Receiver<()>) {
623 let spawned = thread::Builder::new()
624 .name("koan-upnp-renew".into())
625 .spawn(move || {
626 loop {
627 let next = shared.subscriptions.lock().iter().map(|s| s.renew_at).min();
628 let wait = next
629 .map(|at| at.saturating_duration_since(Instant::now()))
630 .unwrap_or(Duration::from_secs(3600));
631 match stop.recv_timeout(wait) {
632 Err(RecvTimeoutError::Timeout) => {}
633 _ => return,
634 }
635 let now = Instant::now();
636 let due: Vec<Subscription> = {
637 let mut subs = shared.subscriptions.lock();
638 let (due, keep) = subs.drain(..).partition(|s| s.renew_at <= now);
639 *subs = keep;
640 due
641 };
642 for sub in due {
643 let renewed = shared
646 .renew(&sub)
647 .or_else(|_| shared.subscribe(&sub.service));
648 match renewed {
649 Ok(s) => shared.subscriptions.lock().push(s),
650 Err(e) => {
651 log::info!("upnp: lost events from {}: {e}", shared.renderer.name);
652 if sub.service == shared.renderer.av_transport {
653 shared.polling.store(true, Ordering::Release);
654 let _ = shared.look.try_send(());
655 }
656 }
657 }
658 }
659 }
660 });
661 if let Err(e) = spawned {
662 log::warn!("upnp: could not start subscription renewal: {e}");
663 }
664}
665
666#[derive(Debug, Default, PartialEq)]
667pub(crate) struct Change {
668 pub transport: bool,
670 pub volume: Option<u8>,
671}
672
673pub(crate) fn parse_notify(body: &str) -> Change {
676 let mut change = Change::default();
677 let Ok(set) = xml::parse(body) else {
678 return change;
679 };
680 for property in set.children_named("property") {
681 let Some(last) = property.child("LastChange") else {
682 continue;
683 };
684 let Ok(event) = xml::parse(last.text.trim()) else {
685 continue;
686 };
687 for instance in event.children_named("InstanceID") {
688 for var in &instance.children {
689 match var.name.as_str() {
690 "TransportState" | "CurrentTrackURI" | "AVTransportURI" => {
691 change.transport = true;
692 }
693 "Volume" if var.attr("channel").is_none_or(|c| c == "Master") => {
694 change.volume = var
695 .attr("val")
696 .and_then(|v| v.parse::<u32>().ok())
697 .map(|v| v.min(100) as u8);
698 }
699 _ => {}
700 }
701 }
702 }
703 }
704 change
705}
706
707#[cfg(test)]
708mod tests {
709 use super::*;
710
711 fn wrap(event: &str) -> String {
712 format!(
713 "<?xml version=\"1.0\"?><e:propertyset xmlns:e=\"urn:schemas-upnp-org:event-1-0\"><e:property><LastChange>{}</LastChange></e:property></e:propertyset>",
714 xml::escape(event)
715 )
716 }
717
718 #[test]
719 fn a_transport_change_asks_for_a_look() {
720 let change = parse_notify(&wrap(include_str!("fixtures/lastchange-avt.xml")));
721 assert_eq!(
722 change,
723 Change {
724 transport: true,
725 volume: None
726 }
727 );
728 }
729
730 #[test]
731 fn a_volume_change_carries_the_master_volume() {
732 let change = parse_notify(&wrap(include_str!("fixtures/lastchange-rc.xml")));
733 assert_eq!(
734 change,
735 Change {
736 transport: false,
737 volume: Some(37)
738 }
739 );
740 }
741
742 #[test]
743 fn garbage_is_no_change() {
744 assert_eq!(parse_notify("not xml <"), Change::default());
745 }
746
747 #[test]
748 fn transport_states_parse() {
749 assert_eq!(Transport::parse("PLAYING"), Transport::Playing);
750 assert_eq!(Transport::parse("PAUSED_PLAYBACK"), Transport::Paused);
751 assert_eq!(Transport::parse("STOPPED"), Transport::Stopped);
752 assert_eq!(Transport::parse("NO_MEDIA_PRESENT"), Transport::NoMedia);
753 }
754}