use std::net::TcpStream;
use std::os::fd::RawFd;
use std::path::Path;
use std::sync::Arc;
use std::time::{Duration, Instant};
use parking_lot::{Condvar, Mutex};
use serde::{Deserialize, Serialize};
use tungstenite::stream::MaybeTlsStream;
use crate::config::{self, Config};
use crate::helpers::{subsonic_auth, subsonic_client};
use crate::remote::client::SubsonicAuth;
use crate::remote::profile;
use crate::remote::wire::{self, Waker};
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "camelCase")]
pub enum LinkCommand {
#[serde(rename_all = "camelCase")]
Play {
track_ids: Vec<String>,
#[serde(default)]
start_at: u32,
#[serde(default, skip_serializing_if = "is_zero")]
position_ms: u64,
},
#[serde(rename_all = "camelCase")]
Enqueue {
track_ids: Vec<String>,
},
#[serde(rename_all = "camelCase")]
PlayNext {
track_ids: Vec<String>,
},
#[serde(rename_all = "camelCase")]
Remove {
track_ids: Vec<String>,
},
Clear,
Radio {
enabled: bool,
},
Sync {
#[serde(default)]
full: bool,
},
#[serde(rename_all = "camelCase")]
Evict {
track_ids: Vec<String>,
},
#[serde(rename_all = "camelCase")]
JumpTo {
track_id: String,
},
#[serde(rename_all = "camelCase")]
Seek {
position_ms: u64,
},
Pause,
Resume,
Next,
Previous,
PlayItem {
id: String,
},
RemoveItems {
ids: Vec<String>,
},
MoveItems {
ids: Vec<String>,
target: String,
after: bool,
},
#[serde(rename_all = "camelCase")]
Insert {
track_ids: Vec<String>,
after: String,
},
Undo,
Redo,
HandOff {
to: String,
},
Devices {
devices: Vec<LinkDevice>,
},
}
fn is_zero(n: &u64) -> bool {
*n == 0
}
impl LinkCommand {
pub fn allowed_nearby(&self) -> bool {
!matches!(
self,
Self::Sync { .. } | Self::Evict { .. } | Self::Devices { .. }
)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct LinkDevice {
pub id: String,
pub name: String,
pub platform: String,
pub linked: bool,
pub state: Option<LinkState>,
}
impl LinkCommand {
pub fn track_ids_mut(&mut self) -> Vec<&mut String> {
match self {
Self::Play { track_ids, .. }
| Self::Enqueue { track_ids }
| Self::PlayNext { track_ids }
| Self::Remove { track_ids }
| Self::Evict { track_ids }
| Self::Insert { track_ids, .. } => track_ids.iter_mut().collect(),
Self::JumpTo { track_id } => vec![track_id],
Self::PlayItem { .. }
| Self::RemoveItems { .. }
| Self::MoveItems { .. }
| Self::Undo
| Self::Redo
| Self::HandOff { .. }
| Self::Devices { .. }
| Self::Clear
| Self::Radio { .. }
| Self::Sync { .. }
| Self::Seek { .. }
| Self::Pause
| Self::Resume
| Self::Next
| Self::Previous => Vec::new(),
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct LinkState {
pub playing: bool,
pub title: Option<String>,
pub artist: Option<String>,
#[serde(default)]
pub album: Option<String>,
#[serde(default)]
pub position_ms: u64,
#[serde(default)]
pub duration_ms: u64,
#[serde(default)]
pub radio: bool,
#[serde(default)]
pub queue: Vec<LinkQueueEntry>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct LinkQueueEntry {
#[serde(default)]
pub id: Option<String>,
pub track_id: Option<String>,
pub title: String,
pub artist: String,
#[serde(default)]
pub album: String,
#[serde(default)]
pub duration_ms: u64,
pub current: bool,
}
impl LinkState {
pub fn differs(&self, sent: &LinkState, elapsed: Duration) -> bool {
let strip = |s: &LinkState| LinkState {
position_ms: 0,
..s.clone()
};
if strip(self) != strip(sent) {
return true;
}
let expected = if sent.playing {
sent.position_ms + elapsed.as_millis() as u64
} else {
sent.position_ms
};
self.position_ms.abs_diff(expected) > 3000
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "camelCase")]
pub enum LinkReport {
State(LinkState),
Push {
token: String,
sandbox: bool,
},
Command {
to: String,
command: LinkCommand,
},
Activity {
token: Option<String>,
device: Option<String>,
#[serde(default)]
sandbox: bool,
},
Hello(LinkHello),
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct LinkHello {
pub id: String,
pub name: String,
pub platform: String,
pub library: Option<String>,
}
pub fn library_fingerprint(cfg: &Config) -> Option<String> {
let auth = subsonic_auth(cfg)?;
let url = auth.base_url.trim_end_matches('/').to_ascii_lowercase();
Some(format!("{:x}", md5::compute(url.as_bytes())))
}
static PUSH_TOKEN: Mutex<Option<(String, bool)>> = Mutex::new(None);
type ActivityToken = (String, String, bool);
static ACTIVITY: Mutex<Option<Option<ActivityToken>>> = Mutex::new(None);
pub fn set_activity(activity: Option<ActivityToken>) {
*ACTIVITY.lock() = Some(activity);
if let Some(up) = LINK.lock().as_ref() {
up.waker.wake();
}
}
pub fn parse_command(json: &str) -> Result<LinkCommand, String> {
serde_json::from_str(json).map_err(|e| e.to_string())
}
pub fn set_push_token(token: String, sandbox: bool) {
*PUSH_TOKEN.lock() = Some((token, sandbox));
nudge();
if let Some(up) = LINK.lock().as_ref() {
up.waker.wake();
}
}
#[derive(Debug, Clone)]
pub struct LinkIdentity {
pub name: String,
pub platform: String,
pub device_id: String,
}
impl LinkIdentity {
pub fn this_device(name: Option<String>) -> Self {
let (platform, label) = if cfg!(target_os = "ios") {
("ios", "iPhone")
} else if cfg!(target_os = "macos") {
("macos", "Mac")
} else {
("linux", "Linux")
};
Self {
name: name
.filter(|n| !n.trim().is_empty())
.or_else(hostname)
.unwrap_or_else(|| label.to_string()),
platform: platform.to_string(),
device_id: device_id(&config::config_dir()),
}
}
}
#[derive(Clone)]
pub struct Local {
pub identity: LinkIdentity,
pub state: Arc<dyn Fn() -> LinkState + Send + Sync>,
pub on_command: Arc<dyn Fn(LinkCommand) + Send + Sync>,
}
pub fn spawn(local: Local) {
std::thread::Builder::new()
.name("koan-link".into())
.spawn(move || run(local))
.expect("failed to spawn the link thread");
}
const RETRY_MIN: Duration = Duration::from_secs(2);
const RETRY_MAX: Duration = Duration::from_secs(60);
fn run(local: Local) {
let mut wait = RETRY_MIN;
loop {
let cfg = Config::load().unwrap_or_default();
let Some(auth) = subsonic_auth(&cfg) else {
rest(RETRY_MAX);
continue;
};
match profile::for_auth(&auth) {
Some(p) if p.links() => {}
Some(_) => {
rest(RETRY_MAX);
continue;
}
None => {
rest(wait);
wait = (wait * 2).min(RETRY_MAX);
continue;
}
}
match connect(&auth, &local.identity) {
Ok((mut socket, fd)) => {
log::info!("link: connected to {}", auth.base_url);
wait = RETRY_MIN;
if let Err(e) = serve(&mut socket, fd, &local) {
log::info!("link: closed: {e}");
}
*LINK.lock() = None;
crate::remote::devices::set_linked(false);
}
Err(e) => {
log::warn!("link: {e}");
profile::forget();
}
}
if rest(wait) {
wait = RETRY_MIN;
continue;
}
wait = (wait * 2).min(RETRY_MAX);
}
}
static NUDGE: (Mutex<bool>, Condvar) = (Mutex::new(false), Condvar::new());
pub fn nudge() {
*NUDGE.0.lock() = true;
NUDGE.1.notify_all();
}
fn rest(d: Duration) -> bool {
let mut nudged = NUDGE.0.lock();
if !*nudged {
NUDGE.1.wait_for(&mut nudged, d);
}
std::mem::replace(&mut *nudged, false)
}
struct Up {
waker: Arc<Waker>,
outbox: Vec<LinkReport>,
}
static LINK: Mutex<Option<Up>> = Mutex::new(None);
pub fn report(report: LinkReport) -> bool {
let mut link = LINK.lock();
let Some(up) = link.as_mut() else {
return false;
};
up.outbox.push(report);
up.waker.wake();
true
}
pub fn is_up() -> bool {
LINK.lock().is_some()
}
type Socket = tungstenite::WebSocket<MaybeTlsStream<TcpStream>>;
fn connect(auth: &SubsonicAuth, identity: &LinkIdentity) -> Result<(Socket, RawFd), String> {
let url = link_url(auth, identity)?;
let (socket, _) = tungstenite::connect(url).map_err(|e| e.to_string())?;
let fd = wire::prepare(socket.get_ref())?;
Ok((socket, fd))
}
fn link_url(auth: &SubsonicAuth, identity: &LinkIdentity) -> Result<String, String> {
let base = if let Some(rest) = auth.base_url.strip_prefix("https://") {
format!("wss://{rest}")
} else if let Some(rest) = auth.base_url.strip_prefix("http://") {
format!("ws://{rest}")
} else {
return Err(format!("not an http(s) server: {}", auth.base_url));
};
let mut query = auth.query().map_err(|e| e.to_string())?;
for (k, v) in [
("client", identity.name.as_str()),
("platform", identity.platform.as_str()),
("device", identity.device_id.as_str()),
("devices", "1"),
] {
query.push('&');
query.push_str(k);
query.push('=');
query.push_str(&percent_encode(v));
}
Ok(format!("{base}/rest/koanLink?{query}"))
}
fn serve(socket: &mut Socket, fd: RawFd, local: &Local) -> Result<(), String> {
let waker = Waker::new().map_err(|e| e.to_string())?;
wire::wake_on_engine_change(&waker);
*LINK.lock() = Some(Up {
waker: waker.clone(),
outbox: Vec::new(),
});
crate::remote::devices::set_linked(true);
let mut session = LinkSession {
local,
sent: None,
sent_push: None,
sent_activity: None,
};
wire::drive(socket, fd, &waker, &mut session)
}
struct LinkSession<'a> {
local: &'a Local,
sent: Option<(LinkState, Instant)>,
sent_push: Option<(String, bool)>,
sent_activity: Option<Option<ActivityToken>>,
}
impl wire::Session for LinkSession<'_> {
fn outgoing(&mut self) -> Vec<String> {
let mut out = Vec::new();
let push = PUSH_TOKEN.lock().clone();
if let Some((token, sandbox)) = push.clone()
&& push != self.sent_push
{
out.push(LinkReport::Push { token, sandbox });
self.sent_push = push;
}
let activity = ACTIVITY.lock().clone();
if activity.is_some() && activity != self.sent_activity {
let (token, device, sandbox) = match activity.clone().flatten() {
Some((t, d, s)) => (Some(t), Some(d), s),
None => (None, None, false),
};
out.push(LinkReport::Activity {
token,
device,
sandbox,
});
self.sent_activity = activity;
}
if let Some(up) = LINK.lock().as_mut() {
out.append(&mut up.outbox);
}
let now = (self.local.state)();
if self
.sent
.as_ref()
.is_none_or(|(s, at)| now.differs(s, at.elapsed()))
{
out.push(LinkReport::State(now.clone()));
self.sent = Some((now, Instant::now()));
}
out.iter()
.filter_map(|r| serde_json::to_string(r).ok())
.collect()
}
fn incoming(&mut self, text: &str) {
match serde_json::from_str::<LinkCommand>(text) {
Ok(LinkCommand::Devices { devices }) => {
crate::remote::devices::set_account(devices);
}
Ok(cmd) => (self.local.on_command)(cmd),
Err(e) => log::warn!("link: not a command ({e}): {text}"),
}
}
}
pub fn resolve_tracks(
db: &crate::db::connection::Database,
remote_ids: &[String],
) -> (Vec<i64>, bool) {
let lookup = |db: &crate::db::connection::Database| {
let mut stmt = db
.conn
.prepare_cached(
"SELECT id FROM tracks WHERE uid = ?1
UNION ALL SELECT id FROM tracks WHERE remote_id = ?1 LIMIT 1",
)
.ok();
remote_ids
.iter()
.map(|rid| {
stmt.as_mut()
.and_then(|s| s.query_row([rid], |r| r.get::<_, i64>(0)).ok())
})
.collect::<Vec<_>>()
};
let found = lookup(db);
if found.iter().all(Option::is_some) {
return (found.into_iter().flatten().collect(), false);
}
sync(db, false);
(lookup(db).into_iter().flatten().collect(), true)
}
pub fn sync(db: &crate::db::connection::Database, full: bool) {
let cfg = Config::load().unwrap_or_default();
if let Some(client) = subsonic_client(&cfg)
&& let Err(e) = crate::helpers::sync_remote(
db,
&client,
full,
&cfg.remote.url,
&cfg.remote.username,
&|_| {},
)
{
log::warn!("link: sync failed: {e}");
}
}
fn device_id(dir: &Path) -> String {
let path = dir.join("device-id");
if let Ok(id) = std::fs::read_to_string(&path) {
let id = id.trim();
if !id.is_empty() {
return id.to_string();
}
}
let id = uuid::Uuid::now_v7().to_string();
let _ = std::fs::create_dir_all(dir);
let _ = std::fs::write(&path, &id);
id
}
fn hostname() -> Option<String> {
let mut buf = [0u8; 256];
let ok = unsafe { libc::gethostname(buf.as_mut_ptr().cast(), buf.len()) } == 0;
if !ok {
return None;
}
let end = buf.iter().position(|&b| b == 0).unwrap_or(buf.len());
let name = String::from_utf8_lossy(&buf[..end]);
let name = name.trim_end_matches(".local").trim();
(!name.is_empty() && name != "localhost").then(|| name.to_string())
}
fn percent_encode(s: &str) -> String {
let mut out = String::with_capacity(s.len());
for b in s.bytes() {
if b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_' | b'.' | b'~') {
out.push(b as char);
} else {
out.push_str(&format!("%{b:02X}"));
}
}
out
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_push_token_is_a_tagged_report() {
let report = LinkReport::Push {
token: "ab12".into(),
sandbox: true,
};
let text = serde_json::to_string(&report).unwrap();
assert_eq!(text, r#"{"type":"push","token":"ab12","sandbox":true}"#);
assert_eq!(serde_json::from_str::<LinkReport>(&text).unwrap(), report);
}
#[test]
fn commands_are_tagged_json() {
let play = LinkCommand::Play {
track_ids: vec!["12".into(), "34".into()],
start_at: 1,
position_ms: 0,
};
let json = serde_json::to_string(&play).unwrap();
assert_eq!(
json,
r#"{"type":"play","trackIds":["12","34"],"startAt":1}"#
);
assert_eq!(serde_json::from_str::<LinkCommand>(&json).unwrap(), play);
assert_eq!(
serde_json::from_str::<LinkCommand>(r#"{"type":"pause"}"#).unwrap(),
LinkCommand::Pause
);
let report = LinkReport::State(LinkState {
playing: true,
title: Some("Portions for Foxes".into()),
..Default::default()
});
let json = serde_json::to_string(&report).unwrap();
assert!(json.starts_with(r#"{"type":"state","playing":true,"title":"Portions for Foxes""#));
assert_eq!(serde_json::from_str::<LinkReport>(&json).unwrap(), report);
}
#[test]
fn a_playhead_moving_on_time_is_not_news() {
let sent = LinkState {
playing: true,
position_ms: 10_000,
..Default::default()
};
let later = |pos| LinkState {
position_ms: pos,
..sent.clone()
};
let five = Duration::from_secs(5);
assert!(!later(15_000).differs(&sent, five));
assert!(later(60_000).differs(&sent, five), "a seek");
let paused = LinkState {
playing: false,
..later(15_000)
};
assert!(paused.differs(&sent, five));
}
#[test]
fn the_url_follows_the_scheme_and_names_the_device() {
let identity = LinkIdentity {
name: "J's iPhone".into(),
platform: "ios".into(),
device_id: "abc".into(),
};
let url = link_url(
&SubsonicAuth::new("https://music.example.com", "j", "pw"),
&identity,
)
.unwrap();
assert!(url.starts_with("wss://music.example.com/rest/koanLink?"));
assert!(url.contains("client=J%27s%20iPhone"));
assert!(url.contains("device=abc"));
assert!(url.contains("devices=1"));
assert!(
link_url(&SubsonicAuth::new("http://h:4000", "j", "pw"), &identity)
.unwrap()
.starts_with("ws://h:4000/")
);
}
#[test]
fn a_relayed_command_nests_the_command() {
let report = LinkReport::Command {
to: "phone".into(),
command: LinkCommand::HandOff { to: "mac".into() },
};
let json = serde_json::to_string(&report).unwrap();
assert_eq!(
json,
r#"{"type":"command","to":"phone","command":{"type":"handOff","to":"mac"}}"#
);
assert_eq!(serde_json::from_str::<LinkReport>(&json).unwrap(), report);
}
#[test]
fn an_older_queue_entry_still_reads() {
let e: LinkQueueEntry =
serde_json::from_str(r#"{"trackId":"7","title":"t","artist":"a","current":true}"#)
.unwrap();
assert_eq!(e.id, None);
assert_eq!(e.duration_ms, 0);
}
#[test]
fn strangers_cannot_touch_the_library() {
assert!(LinkCommand::Pause.allowed_nearby());
assert!(LinkCommand::HandOff { to: "x".into() }.allowed_nearby());
assert!(!LinkCommand::Sync { full: true }.allowed_nearby());
assert!(!LinkCommand::Evict { track_ids: vec![] }.allowed_nearby());
}
#[test]
fn the_device_id_is_kept() {
let dir = tempfile::tempdir().unwrap();
let first = device_id(dir.path());
assert_eq!(device_id(dir.path()), first);
}
}