Skip to main content

koan_core/upnp/
mod.rs

1//! Playing to UPnP AV MediaRenderers: network amplifiers and streamers that
2//! fetch a URL they are given and decode it themselves.
3//!
4//! A renderer is this koan's output, the way a USB DAC is. The queue, cursor,
5//! history and everything else stay in the `Player`; only the audio goes
6//! elsewhere, as the original file, so playback stays bit-perfect up to the
7//! renderer's own DAC. Nothing here involves the server.
8//!
9//! Hand-rolled on std threads and the blocking reqwest koan already uses:
10//! eight AVTransport actions, two RenderingControl ones, SSDP and GENA.
11
12pub mod description;
13pub mod didl;
14pub mod discovery;
15pub mod serve;
16pub mod session;
17pub mod soap;
18pub mod stream;
19mod xml;
20
21#[cfg(test)]
22pub(crate) mod fake;
23
24pub use description::Renderer;
25
26/// The renderer this koan is playing to, as the front ends show it.
27#[derive(Debug, Clone, PartialEq)]
28pub struct Output {
29    pub udn: String,
30    pub name: String,
31    /// 0–100, when the renderer has a volume control.
32    pub volume: Option<u8>,
33    /// Why the current track is not playing there, when koan skipped it.
34    pub problem: Option<String>,
35}
36
37/// An open session with a renderer, for the player to play to.
38///
39/// What the renderer is heard to do reaches the player as
40/// `PlayerCommand::Renderer`, tagged with the number of the player session it
41/// was heard in: `tag`, which the player keeps up to date. An event from a
42/// session already over is recognised and dropped.
43#[derive(Debug)]
44pub struct Connection {
45    pub session: session::Session,
46    pub tag: std::sync::Arc<std::sync::atomic::AtomicU64>,
47}
48
49/// Open a session with `renderer` whose events go to `player`.
50pub fn open(
51    renderer: Renderer,
52    player: &crossbeam_channel::Sender<crate::player::commands::PlayerCommand>,
53) -> Result<Connection, String> {
54    let tag = std::sync::Arc::new(std::sync::atomic::AtomicU64::new(0));
55    let tx = player.clone();
56    let tagged = tag.clone();
57    let session = session::Session::open(renderer, move |event| {
58        let cmd = crate::player::commands::PlayerCommand::Renderer {
59            session: tagged.load(std::sync::atomic::Ordering::Acquire),
60            event,
61        };
62        match &cmd {
63            // Only a wake-up: the player reads the loss from the session. It
64            // can be raised on the player's own thread, or under the
65            // discovery lock, where waiting for room in the channel could
66            // never end.
67            crate::player::commands::PlayerCommand::Renderer {
68                event: session::Event::Gone,
69                ..
70            } => {
71                let _ = tx.try_send(cmd);
72            }
73            _ => {
74                let _ = tx.send(cmd);
75            }
76        }
77    })?;
78    Ok(Connection { session, tag })
79}
80
81/// Each choice of output, in the order they were made. Opening a session
82/// takes a few round trips; one that finishes after a later choice is
83/// dropped, so the last one picked is where the music goes.
84static CHOICE: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
85
86/// Record a choice of output, made now, and return its place. Called when the
87/// choice is made, before any work towards it: every output change (a
88/// renderer, a local device, the system default) takes one.
89pub fn choose() -> u64 {
90    CHOICE.fetch_add(1, std::sync::atomic::Ordering::AcqRel) + 1
91}
92
93/// Make the renderer `udn` this koan's output, carrying on from where the
94/// music is, unless a later choice was made while the session opened.
95/// `choice` is what `choose` returned when this one was picked. Blocks for
96/// the few round trips a session takes to open, so call it from a thread
97/// nothing is waiting on.
98pub fn connect(
99    udn: &str,
100    choice: u64,
101    player: &crossbeam_channel::Sender<crate::player::commands::PlayerCommand>,
102) -> Result<(), String> {
103    let renderer = discovery::find(udn)
104        .ok_or_else(|| "That renderer is no longer on the network.".to_string())?;
105    let connection = open(renderer, player)?;
106    if CHOICE.load(std::sync::atomic::Ordering::Acquire) != choice {
107        log::info!(
108            "upnp: {} opened after a later choice of output; not used",
109            connection.session.renderer().name
110        );
111        return Ok(());
112    }
113    player
114        .send(crate::player::commands::PlayerCommand::UseRenderer(Some(
115            Box::new(connection),
116        )))
117        .map_err(|_| "The player has stopped.".to_string())
118}
119
120/// How long after launch the renderer used last time is looked for.
121const RESUME_WINDOW: std::time::Duration = std::time::Duration::from_secs(6);
122
123/// Go back to the renderer `udn`, used last time, on a thread of its own:
124/// look for it for `RESUME_WINDOW`, waking on each change discovery
125/// announces, and hand the player an open session as
126/// `PlayerCommand::ResumeRenderer` if it turns up idle. The player takes it
127/// only if nobody has played anything or picked an output meanwhile.
128pub fn resume(
129    udn: String,
130    player: &crossbeam_channel::Sender<crate::player::commands::PlayerCommand>,
131) {
132    let choice = CHOICE.load(std::sync::atomic::Ordering::Acquire);
133    let player = player.clone();
134    let _ = std::thread::Builder::new()
135        .name("koan-upnp-resume".into())
136        .spawn(move || {
137            if let Some(connection) = resume_onto(&udn, RESUME_WINDOW, choice, &player) {
138                let _ = player.send(crate::player::commands::PlayerCommand::ResumeRenderer(
139                    Box::new(connection),
140                ));
141            }
142        });
143}
144
145/// The session `resume` hands over, if the renderer is found within `window`,
146/// nothing else has been picked since `choice`, and it is not playing or
147/// paused for something else: whatever is driving it, a phone on the same
148/// amplifier say, is not taken over by a launch nobody asked for.
149fn resume_onto(
150    udn: &str,
151    window: std::time::Duration,
152    choice: u64,
153    player: &crossbeam_channel::Sender<crate::player::commands::PlayerCommand>,
154) -> Option<Connection> {
155    let Some(renderer) = await_renderer(udn, window) else {
156        log::info!("upnp: {udn}, used last time, is not on the network; playing here");
157        return None;
158    };
159    if CHOICE.load(std::sync::atomic::Ordering::Acquire) != choice {
160        return None;
161    }
162    if discovery::in_use(&renderer) != Some(false) {
163        log::info!(
164            "upnp: {}, used last time, is busy or not answering; playing here",
165            renderer.name
166        );
167        return None;
168    }
169    open(renderer, player)
170        .inspect_err(|e| log::info!("upnp: could not go back to {udn}: {e}"))
171        .ok()
172}
173
174/// The renderer `udn` once discovery has found it, waiting at most `window`.
175fn await_renderer(udn: &str, window: std::time::Duration) -> Option<Renderer> {
176    let signal = crate::signal::engine_changed();
177    let mut seen = signal.generation();
178    let deadline = std::time::Instant::now() + window;
179    discovery::search();
180    loop {
181        if let Some(renderer) = discovery::find(udn) {
182            return Some(renderer);
183        }
184        let left = deadline.saturating_duration_since(std::time::Instant::now());
185        if left.is_zero() {
186            return None;
187        }
188        seen = signal.wait_until(seen, left);
189    }
190}
191
192/// Bring the music back to this device's own output.
193pub fn disconnect(player: &crossbeam_channel::Sender<crate::player::commands::PlayerCommand>) {
194    choose();
195    let _ = player.send(crate::player::commands::PlayerCommand::UseRenderer(None));
196}
197
198#[cfg(test)]
199mod tests {
200    use super::*;
201
202    /// A renderer discovery knows is found at once; one it never hears of
203    /// is given up on when the window closes, and the music stays here.
204    #[test]
205    fn a_remembered_renderer_is_found_or_given_up_on() {
206        let fake = fake::FakeRenderer::start("http-get:*:audio/wav:*", true, false);
207        let renderer = fake.renderer();
208        discovery::remember(renderer.clone(), std::time::Duration::from_secs(60));
209        assert_eq!(
210            await_renderer(&renderer.udn, std::time::Duration::from_secs(1)).map(|r| r.udn),
211            Some(renderer.udn)
212        );
213
214        let start = std::time::Instant::now();
215        let window = std::time::Duration::from_millis(300);
216        assert!(await_renderer("uuid:nowhere", window).is_none());
217        assert!(start.elapsed() >= window);
218    }
219
220    /// A renderer playing for something else when the app launches is left
221    /// to it; an idle one is gone back to.
222    #[test]
223    fn a_busy_renderer_is_not_taken_at_launch() {
224        let fake = fake::FakeRenderer::start("http-get:*:audio/wav:*", true, false);
225        let renderer = fake.renderer();
226        discovery::remember(renderer.clone(), std::time::Duration::from_secs(60));
227        let (tx, _rx) = crossbeam_channel::unbounded();
228        let window = std::time::Duration::from_secs(1);
229        let now = CHOICE.load(std::sync::atomic::Ordering::Acquire);
230
231        fake.play_foreign("http://phone/track.flac");
232        assert!(resume_onto(&renderer.udn, window, now, &tx).is_none());
233
234        fake.press_stop(0);
235        assert!(resume_onto(&renderer.udn, window, now, &tx).is_some());
236    }
237}