pub mod description;
pub mod didl;
pub mod discovery;
pub mod serve;
pub mod session;
pub mod soap;
pub mod stream;
mod xml;
#[cfg(test)]
pub(crate) mod fake;
pub use description::Renderer;
#[derive(Debug, Clone, PartialEq)]
pub struct Output {
pub udn: String,
pub name: String,
pub volume: Option<u8>,
pub problem: Option<String>,
}
#[derive(Debug)]
pub struct Connection {
pub session: session::Session,
pub tag: std::sync::Arc<std::sync::atomic::AtomicU64>,
}
pub fn open(
renderer: Renderer,
player: &crossbeam_channel::Sender<crate::player::commands::PlayerCommand>,
) -> Result<Connection, String> {
let tag = std::sync::Arc::new(std::sync::atomic::AtomicU64::new(0));
let tx = player.clone();
let tagged = tag.clone();
let session = session::Session::open(renderer, move |event| {
let cmd = crate::player::commands::PlayerCommand::Renderer {
session: tagged.load(std::sync::atomic::Ordering::Acquire),
event,
};
match &cmd {
crate::player::commands::PlayerCommand::Renderer {
event: session::Event::Gone,
..
} => {
let _ = tx.try_send(cmd);
}
_ => {
let _ = tx.send(cmd);
}
}
})?;
Ok(Connection { session, tag })
}
static CHOICE: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
pub fn choose() -> u64 {
CHOICE.fetch_add(1, std::sync::atomic::Ordering::AcqRel) + 1
}
pub fn connect(
udn: &str,
choice: u64,
player: &crossbeam_channel::Sender<crate::player::commands::PlayerCommand>,
) -> Result<(), String> {
let renderer = discovery::find(udn)
.ok_or_else(|| "That renderer is no longer on the network.".to_string())?;
let connection = open(renderer, player)?;
if CHOICE.load(std::sync::atomic::Ordering::Acquire) != choice {
log::info!(
"upnp: {} opened after a later choice of output; not used",
connection.session.renderer().name
);
return Ok(());
}
player
.send(crate::player::commands::PlayerCommand::UseRenderer(Some(
Box::new(connection),
)))
.map_err(|_| "The player has stopped.".to_string())
}
const RESUME_WINDOW: std::time::Duration = std::time::Duration::from_secs(6);
pub fn resume(
udn: String,
player: &crossbeam_channel::Sender<crate::player::commands::PlayerCommand>,
) {
let choice = CHOICE.load(std::sync::atomic::Ordering::Acquire);
let player = player.clone();
let _ = std::thread::Builder::new()
.name("koan-upnp-resume".into())
.spawn(move || {
let _ = player.send(match resume_onto(&udn, RESUME_WINDOW, choice, &player) {
Some(connection) => {
crate::player::commands::PlayerCommand::ResumeRenderer(Box::new(connection))
}
None => crate::player::commands::PlayerCommand::ResumeRendererMissed,
});
});
}
fn resume_onto(
udn: &str,
window: std::time::Duration,
choice: u64,
player: &crossbeam_channel::Sender<crate::player::commands::PlayerCommand>,
) -> Option<Connection> {
let Some(renderer) = await_renderer(udn, window) else {
log::info!("upnp: {udn}, used last time, is not on the network; playing here");
return None;
};
if CHOICE.load(std::sync::atomic::Ordering::Acquire) != choice {
return None;
}
match discovery::in_use(&renderer) {
Some(false) => {}
Some(true)
if discovery::playing_uri(&renderer)
.is_some_and(|uri| served_from_here(&renderer, &uri) && !still_served(&uri)) =>
{
log::info!(
"upnp: {} is still on what a kōan here sent it; taking it back",
renderer.name
);
}
_ => {
log::info!(
"upnp: {}, used last time, is busy or not answering; playing here",
renderer.name
);
return None;
}
}
open(renderer, player)
.inspect_err(|e| log::info!("upnp: could not go back to {udn}: {e}"))
.ok()
}
fn served_from_here(renderer: &Renderer, uri: &str) -> bool {
let Ok(url) = url::Url::parse(uri) else {
return false;
};
let host = |u: &url::Url| {
u.host_str()
.and_then(|h| h.parse::<std::net::IpAddr>().ok())
};
let Some(here) = host(&renderer.location).and_then(serve::local_ip_towards) else {
return false;
};
let token = url
.path()
.strip_prefix("/t/")
.and_then(|t| t.split('.').next())
.unwrap_or_default();
host(&url) == Some(here) && token.len() == 32 && token.bytes().all(|b| b.is_ascii_hexdigit())
}
fn still_served(uri: &str) -> bool {
use std::io::{Read, Write};
let Ok(url) = url::Url::parse(uri) else {
return false;
};
let Some(addr) = url
.socket_addrs(|| Some(80))
.ok()
.and_then(|a| a.into_iter().next())
else {
return false;
};
let timeout = std::time::Duration::from_millis(500);
let Ok(mut stream) = std::net::TcpStream::connect_timeout(&addr, timeout) else {
return false;
};
let _ = stream.set_read_timeout(Some(timeout));
if write!(stream, "HEAD {} HTTP/1.1\r\nHost: x\r\n\r\n", url.path()).is_err() {
return false;
}
let mut head = [0u8; 12];
stream.read_exact(&mut head).is_ok() && head.starts_with(b"HTTP/1.1 2")
}
fn await_renderer(udn: &str, window: std::time::Duration) -> Option<Renderer> {
let signal = crate::signal::engine_changed();
let mut seen = signal.generation();
let deadline = std::time::Instant::now() + window;
discovery::search();
loop {
if let Some(renderer) = discovery::find(udn) {
return Some(renderer);
}
let left = deadline.saturating_duration_since(std::time::Instant::now());
if left.is_zero() {
return None;
}
seen = signal.wait_until(seen, left);
}
}
pub fn disconnect(player: &crossbeam_channel::Sender<crate::player::commands::PlayerCommand>) {
choose();
let _ = player.send(crate::player::commands::PlayerCommand::UseRenderer(None));
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_remembered_renderer_is_found_or_given_up_on() {
let fake = fake::FakeRenderer::start("http-get:*:audio/wav:*", true, false);
let renderer = fake.renderer();
discovery::remember(renderer.clone(), std::time::Duration::from_secs(60));
assert_eq!(
await_renderer(&renderer.udn, std::time::Duration::from_secs(1)).map(|r| r.udn),
Some(renderer.udn)
);
let start = std::time::Instant::now();
let window = std::time::Duration::from_millis(300);
assert!(await_renderer("uuid:nowhere", window).is_none());
assert!(start.elapsed() >= window);
}
#[test]
fn a_busy_renderer_is_not_taken_at_launch() {
let fake = fake::FakeRenderer::start("http-get:*:audio/wav:*", true, false);
let renderer = fake.renderer();
discovery::remember(renderer.clone(), std::time::Duration::from_secs(60));
let (tx, _rx) = crossbeam_channel::unbounded();
let window = std::time::Duration::from_secs(1);
let now = CHOICE.load(std::sync::atomic::Ordering::Acquire);
fake.play_foreign("http://phone/track.flac");
assert!(resume_onto(&renderer.udn, window, now, &tx).is_none());
let gone = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let gone_port = gone.local_addr().unwrap().port();
drop(gone);
fake.play_foreign(&format!(
"http://127.0.0.1:{gone_port}/t/{}.flac",
"0a".repeat(16)
));
assert!(resume_onto(&renderer.udn, window, now, &tx).is_some());
let dir = tempfile::tempdir().unwrap();
let file = dir.path().join("t.flac");
std::fs::write(&file, b"fLaC").unwrap();
let live = serve::Listener::start(Box::new(|_, _| {})).unwrap();
let token = live.add(serve::Served::File {
path: file,
mime: "audio/flac".into(),
});
fake.play_foreign(&format!("http://127.0.0.1:{}/t/{token}.flac", live.port()));
assert!(resume_onto(&renderer.udn, window, now, &tx).is_none());
fake.press_stop(0);
assert!(resume_onto(&renderer.udn, window, now, &tx).is_some());
}
}