1use std::net::TcpStream;
11use std::os::fd::RawFd;
12use std::path::Path;
13use std::sync::Arc;
14use std::time::{Duration, Instant};
15
16use parking_lot::{Condvar, Mutex};
17use serde::{Deserialize, Serialize};
18use tungstenite::stream::MaybeTlsStream;
19
20use crate::config::{self, Config};
21use crate::helpers::{subsonic_auth, subsonic_client};
22use crate::remote::client::SubsonicAuth;
23use crate::remote::profile;
24use crate::remote::wire::{self, Waker};
25
26#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
28#[serde(tag = "type", rename_all = "camelCase")]
29pub enum LinkCommand {
30 #[serde(rename_all = "camelCase")]
33 Play {
34 track_ids: Vec<String>,
35 #[serde(default)]
36 start_at: u32,
37 #[serde(default, skip_serializing_if = "is_zero")]
38 position_ms: u64,
39 },
40 #[serde(rename_all = "camelCase")]
42 Enqueue {
43 track_ids: Vec<String>,
44 },
45 #[serde(rename_all = "camelCase")]
47 PlayNext {
48 track_ids: Vec<String>,
49 },
50 #[serde(rename_all = "camelCase")]
52 Remove {
53 track_ids: Vec<String>,
54 },
55 Clear,
56 Radio {
57 enabled: bool,
58 },
59 Sync {
62 #[serde(default)]
63 full: bool,
64 },
65 #[serde(rename_all = "camelCase")]
68 Evict {
69 track_ids: Vec<String>,
70 },
71 #[serde(rename_all = "camelCase")]
74 JumpTo {
75 track_id: String,
76 },
77 #[serde(rename_all = "camelCase")]
78 Seek {
79 position_ms: u64,
80 },
81 Pause,
82 Resume,
83 Next,
84 Previous,
85 PlayItem {
89 id: String,
90 },
91 RemoveItems {
92 ids: Vec<String>,
93 },
94 MoveItems {
96 ids: Vec<String>,
97 target: String,
98 after: bool,
99 },
100 #[serde(rename_all = "camelCase")]
102 Insert {
103 track_ids: Vec<String>,
104 after: String,
105 },
106 Undo,
107 Redo,
108 HandOff {
112 to: String,
113 },
114 Devices {
117 devices: Vec<LinkDevice>,
118 },
119}
120
121fn is_zero(n: &u64) -> bool {
122 *n == 0
123}
124
125impl LinkCommand {
126 pub fn allowed_nearby(&self) -> bool {
130 !matches!(
131 self,
132 Self::Sync { .. } | Self::Evict { .. } | Self::Devices { .. }
133 )
134 }
135}
136
137#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
139#[serde(rename_all = "camelCase")]
140pub struct LinkDevice {
141 pub id: String,
143 pub name: String,
144 pub platform: String,
145 pub linked: bool,
148 pub state: Option<LinkState>,
151}
152
153impl LinkCommand {
154 pub fn track_ids_mut(&mut self) -> Vec<&mut String> {
157 match self {
158 Self::Play { track_ids, .. }
159 | Self::Enqueue { track_ids }
160 | Self::PlayNext { track_ids }
161 | Self::Remove { track_ids }
162 | Self::Evict { track_ids }
163 | Self::Insert { track_ids, .. } => track_ids.iter_mut().collect(),
164 Self::JumpTo { track_id } => vec![track_id],
165 Self::PlayItem { .. }
166 | Self::RemoveItems { .. }
167 | Self::MoveItems { .. }
168 | Self::Undo
169 | Self::Redo
170 | Self::HandOff { .. }
171 | Self::Devices { .. }
172 | Self::Clear
173 | Self::Radio { .. }
174 | Self::Sync { .. }
175 | Self::Seek { .. }
176 | Self::Pause
177 | Self::Resume
178 | Self::Next
179 | Self::Previous => Vec::new(),
180 }
181 }
182}
183
184#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
186#[serde(rename_all = "camelCase")]
187pub struct LinkState {
188 pub playing: bool,
189 pub title: Option<String>,
190 pub artist: Option<String>,
191 #[serde(default)]
192 pub album: Option<String>,
193 #[serde(default)]
195 pub position_ms: u64,
196 #[serde(default)]
197 pub duration_ms: u64,
198 #[serde(default)]
199 pub radio: bool,
200 #[serde(default)]
202 pub queue: Vec<LinkQueueEntry>,
203}
204
205#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
206#[serde(rename_all = "camelCase")]
207pub struct LinkQueueEntry {
208 #[serde(default)]
210 pub id: Option<String>,
211 pub track_id: Option<String>,
213 pub title: String,
214 pub artist: String,
215 #[serde(default)]
216 pub album: String,
217 #[serde(default)]
218 pub duration_ms: u64,
219 pub current: bool,
220}
221
222impl LinkState {
223 pub fn differs(&self, sent: &LinkState, elapsed: Duration) -> bool {
227 let strip = |s: &LinkState| LinkState {
228 position_ms: 0,
229 ..s.clone()
230 };
231 if strip(self) != strip(sent) {
232 return true;
233 }
234 let expected = if sent.playing {
235 sent.position_ms + elapsed.as_millis() as u64
236 } else {
237 sent.position_ms
238 };
239 self.position_ms.abs_diff(expected) > 3000
240 }
241}
242
243#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
245#[serde(tag = "type", rename_all = "camelCase")]
246pub enum LinkReport {
247 State(LinkState),
248 Push {
252 token: String,
253 sandbox: bool,
254 },
255 Command {
257 to: String,
258 command: LinkCommand,
259 },
260 Activity {
263 token: Option<String>,
264 device: Option<String>,
265 #[serde(default)]
266 sandbox: bool,
267 },
268 Hello(LinkHello),
270}
271
272#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
275#[serde(rename_all = "camelCase")]
276pub struct LinkHello {
277 pub id: String,
278 pub name: String,
279 pub platform: String,
280 pub library: Option<String>,
284}
285
286pub fn library_fingerprint(cfg: &Config) -> Option<String> {
289 let auth = subsonic_auth(cfg)?;
290 let url = auth.base_url.trim_end_matches('/').to_ascii_lowercase();
291 Some(format!("{:x}", md5::compute(url.as_bytes())))
292}
293
294static PUSH_TOKEN: Mutex<Option<(String, bool)>> = Mutex::new(None);
297
298type ActivityToken = (String, String, bool);
302static ACTIVITY: Mutex<Option<Option<ActivityToken>>> = Mutex::new(None);
303
304pub fn set_activity(activity: Option<ActivityToken>) {
307 *ACTIVITY.lock() = Some(activity);
308 if let Some(up) = LINK.lock().as_ref() {
309 up.waker.wake();
310 }
311}
312
313pub fn parse_command(json: &str) -> Result<LinkCommand, String> {
316 serde_json::from_str(json).map_err(|e| e.to_string())
317}
318
319pub fn set_push_token(token: String, sandbox: bool) {
321 *PUSH_TOKEN.lock() = Some((token, sandbox));
322 nudge();
323 if let Some(up) = LINK.lock().as_ref() {
324 up.waker.wake();
325 }
326}
327
328#[derive(Debug, Clone)]
330pub struct LinkIdentity {
331 pub name: String,
333 pub platform: String,
335 pub device_id: String,
338}
339
340impl LinkIdentity {
341 pub fn this_device(name: Option<String>) -> Self {
343 let (platform, label) = if cfg!(target_os = "ios") {
344 ("ios", "iPhone")
345 } else if cfg!(target_os = "macos") {
346 ("macos", "Mac")
347 } else {
348 ("linux", "Linux")
349 };
350 Self {
351 name: name
352 .filter(|n| !n.trim().is_empty())
353 .or_else(hostname)
354 .unwrap_or_else(|| label.to_string()),
355 platform: platform.to_string(),
356 device_id: device_id(&config::config_dir()),
357 }
358 }
359}
360
361#[derive(Clone)]
365pub struct Local {
366 pub identity: LinkIdentity,
367 pub state: Arc<dyn Fn() -> LinkState + Send + Sync>,
368 pub on_command: Arc<dyn Fn(LinkCommand) + Send + Sync>,
369}
370
371pub fn spawn(local: Local) {
380 std::thread::Builder::new()
381 .name("koan-link".into())
382 .spawn(move || run(local))
383 .expect("failed to spawn the link thread");
384}
385
386const RETRY_MIN: Duration = Duration::from_secs(2);
387const RETRY_MAX: Duration = Duration::from_secs(60);
388
389fn run(local: Local) {
390 let mut wait = RETRY_MIN;
391 loop {
392 let cfg = Config::load().unwrap_or_default();
393 let Some(auth) = subsonic_auth(&cfg) else {
394 rest(RETRY_MAX);
395 continue;
396 };
397 match profile::for_auth(&auth) {
398 Some(p) if p.links() => {}
399 Some(_) => {
401 rest(RETRY_MAX);
402 continue;
403 }
404 None => {
405 rest(wait);
406 wait = (wait * 2).min(RETRY_MAX);
407 continue;
408 }
409 }
410
411 match connect(&auth, &local.identity) {
412 Ok((mut socket, fd)) => {
413 log::info!("link: connected to {}", auth.base_url);
414 wait = RETRY_MIN;
415 if let Err(e) = serve(&mut socket, fd, &local) {
416 log::info!("link: closed: {e}");
417 }
418 *LINK.lock() = None;
419 crate::remote::devices::set_linked(false);
420 }
421 Err(e) => {
422 log::warn!("link: {e}");
423 profile::forget();
429 }
430 }
431 if rest(wait) {
432 wait = RETRY_MIN;
433 continue;
434 }
435 wait = (wait * 2).min(RETRY_MAX);
436 }
437}
438
439static NUDGE: (Mutex<bool>, Condvar) = (Mutex::new(false), Condvar::new());
440
441pub fn nudge() {
447 *NUDGE.0.lock() = true;
448 NUDGE.1.notify_all();
449}
450
451fn rest(d: Duration) -> bool {
453 let mut nudged = NUDGE.0.lock();
454 if !*nudged {
455 NUDGE.1.wait_for(&mut nudged, d);
456 }
457 std::mem::replace(&mut *nudged, false)
458}
459
460struct Up {
462 waker: Arc<Waker>,
463 outbox: Vec<LinkReport>,
464}
465
466static LINK: Mutex<Option<Up>> = Mutex::new(None);
467
468pub fn report(report: LinkReport) -> bool {
470 let mut link = LINK.lock();
471 let Some(up) = link.as_mut() else {
472 return false;
473 };
474 up.outbox.push(report);
475 up.waker.wake();
476 true
477}
478
479pub fn is_up() -> bool {
481 LINK.lock().is_some()
482}
483
484type Socket = tungstenite::WebSocket<MaybeTlsStream<TcpStream>>;
485
486fn connect(auth: &SubsonicAuth, identity: &LinkIdentity) -> Result<(Socket, RawFd), String> {
487 let url = link_url(auth, identity)?;
488 let (socket, _) = tungstenite::connect(url).map_err(|e| e.to_string())?;
489 let fd = wire::prepare(socket.get_ref())?;
490 Ok((socket, fd))
491}
492
493fn link_url(auth: &SubsonicAuth, identity: &LinkIdentity) -> Result<String, String> {
495 let base = if let Some(rest) = auth.base_url.strip_prefix("https://") {
496 format!("wss://{rest}")
497 } else if let Some(rest) = auth.base_url.strip_prefix("http://") {
498 format!("ws://{rest}")
499 } else {
500 return Err(format!("not an http(s) server: {}", auth.base_url));
501 };
502 let mut query = auth.query().map_err(|e| e.to_string())?;
503 for (k, v) in [
504 ("client", identity.name.as_str()),
505 ("platform", identity.platform.as_str()),
506 ("device", identity.device_id.as_str()),
507 ("devices", "1"),
510 ] {
511 query.push('&');
512 query.push_str(k);
513 query.push('=');
514 query.push_str(&percent_encode(v));
515 }
516 Ok(format!("{base}/rest/koanLink?{query}"))
517}
518
519fn serve(socket: &mut Socket, fd: RawFd, local: &Local) -> Result<(), String> {
520 let waker = Waker::new().map_err(|e| e.to_string())?;
521 wire::wake_on_engine_change(&waker);
522 *LINK.lock() = Some(Up {
523 waker: waker.clone(),
524 outbox: Vec::new(),
525 });
526 crate::remote::devices::set_linked(true);
527 let mut session = LinkSession {
528 local,
529 sent: None,
530 sent_push: None,
531 sent_activity: None,
532 };
533 wire::drive(socket, fd, &waker, &mut session)
534}
535
536struct LinkSession<'a> {
537 local: &'a Local,
538 sent: Option<(LinkState, Instant)>,
539 sent_push: Option<(String, bool)>,
540 sent_activity: Option<Option<ActivityToken>>,
541}
542
543impl wire::Session for LinkSession<'_> {
544 fn outgoing(&mut self) -> Vec<String> {
545 let mut out = Vec::new();
546 let push = PUSH_TOKEN.lock().clone();
547 if let Some((token, sandbox)) = push.clone()
548 && push != self.sent_push
549 {
550 out.push(LinkReport::Push { token, sandbox });
551 self.sent_push = push;
552 }
553 let activity = ACTIVITY.lock().clone();
554 if activity.is_some() && activity != self.sent_activity {
555 let (token, device, sandbox) = match activity.clone().flatten() {
556 Some((t, d, s)) => (Some(t), Some(d), s),
557 None => (None, None, false),
558 };
559 out.push(LinkReport::Activity {
560 token,
561 device,
562 sandbox,
563 });
564 self.sent_activity = activity;
565 }
566 if let Some(up) = LINK.lock().as_mut() {
567 out.append(&mut up.outbox);
568 }
569 let now = (self.local.state)();
570 if self
571 .sent
572 .as_ref()
573 .is_none_or(|(s, at)| now.differs(s, at.elapsed()))
574 {
575 out.push(LinkReport::State(now.clone()));
576 self.sent = Some((now, Instant::now()));
577 }
578 out.iter()
579 .filter_map(|r| serde_json::to_string(r).ok())
580 .collect()
581 }
582
583 fn incoming(&mut self, text: &str) {
584 match serde_json::from_str::<LinkCommand>(text) {
585 Ok(LinkCommand::Devices { devices }) => {
586 crate::remote::devices::set_account(devices);
587 }
588 Ok(cmd) => (self.local.on_command)(cmd),
589 Err(e) => log::warn!("link: not a command ({e}): {text}"),
590 }
591 }
592}
593
594pub fn resolve_tracks(
602 db: &crate::db::connection::Database,
603 remote_ids: &[String],
604) -> (Vec<i64>, bool) {
605 let lookup = |db: &crate::db::connection::Database| {
606 let mut stmt = db
607 .conn
608 .prepare_cached(
609 "SELECT id FROM tracks WHERE uid = ?1
610 UNION ALL SELECT id FROM tracks WHERE remote_id = ?1 LIMIT 1",
611 )
612 .ok();
613 remote_ids
614 .iter()
615 .map(|rid| {
616 stmt.as_mut()
617 .and_then(|s| s.query_row([rid], |r| r.get::<_, i64>(0)).ok())
618 })
619 .collect::<Vec<_>>()
620 };
621 let found = lookup(db);
622 if found.iter().all(Option::is_some) {
623 return (found.into_iter().flatten().collect(), false);
624 }
625 sync(db, false);
626 (lookup(db).into_iter().flatten().collect(), true)
627}
628
629pub fn sync(db: &crate::db::connection::Database, full: bool) {
632 let cfg = Config::load().unwrap_or_default();
633 if let Some(client) = subsonic_client(&cfg)
634 && let Err(e) = crate::helpers::sync_remote(
635 db,
636 &client,
637 full,
638 &cfg.remote.url,
639 &cfg.remote.username,
640 &|_| {},
641 )
642 {
643 log::warn!("link: sync failed: {e}");
644 }
645}
646
647fn device_id(dir: &Path) -> String {
649 let path = dir.join("device-id");
650 if let Ok(id) = std::fs::read_to_string(&path) {
651 let id = id.trim();
652 if !id.is_empty() {
653 return id.to_string();
654 }
655 }
656 let id = uuid::Uuid::now_v7().to_string();
657 let _ = std::fs::create_dir_all(dir);
658 let _ = std::fs::write(&path, &id);
659 id
660}
661
662fn hostname() -> Option<String> {
663 let mut buf = [0u8; 256];
664 let ok = unsafe { libc::gethostname(buf.as_mut_ptr().cast(), buf.len()) } == 0;
667 if !ok {
668 return None;
669 }
670 let end = buf.iter().position(|&b| b == 0).unwrap_or(buf.len());
671 let name = String::from_utf8_lossy(&buf[..end]);
672 let name = name.trim_end_matches(".local").trim();
673 (!name.is_empty() && name != "localhost").then(|| name.to_string())
674}
675
676fn percent_encode(s: &str) -> String {
677 let mut out = String::with_capacity(s.len());
678 for b in s.bytes() {
679 if b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_' | b'.' | b'~') {
680 out.push(b as char);
681 } else {
682 out.push_str(&format!("%{b:02X}"));
683 }
684 }
685 out
686}
687
688#[cfg(test)]
689mod tests {
690 use super::*;
691
692 #[test]
693 fn a_push_token_is_a_tagged_report() {
694 let report = LinkReport::Push {
695 token: "ab12".into(),
696 sandbox: true,
697 };
698 let text = serde_json::to_string(&report).unwrap();
699 assert_eq!(text, r#"{"type":"push","token":"ab12","sandbox":true}"#);
700 assert_eq!(serde_json::from_str::<LinkReport>(&text).unwrap(), report);
701 }
702
703 #[test]
704 fn commands_are_tagged_json() {
705 let play = LinkCommand::Play {
706 track_ids: vec!["12".into(), "34".into()],
707 start_at: 1,
708 position_ms: 0,
709 };
710 let json = serde_json::to_string(&play).unwrap();
711 assert_eq!(
712 json,
713 r#"{"type":"play","trackIds":["12","34"],"startAt":1}"#
714 );
715 assert_eq!(serde_json::from_str::<LinkCommand>(&json).unwrap(), play);
716 assert_eq!(
717 serde_json::from_str::<LinkCommand>(r#"{"type":"pause"}"#).unwrap(),
718 LinkCommand::Pause
719 );
720 let report = LinkReport::State(LinkState {
721 playing: true,
722 title: Some("Portions for Foxes".into()),
723 ..Default::default()
724 });
725 let json = serde_json::to_string(&report).unwrap();
726 assert!(json.starts_with(r#"{"type":"state","playing":true,"title":"Portions for Foxes""#));
727 assert_eq!(serde_json::from_str::<LinkReport>(&json).unwrap(), report);
728 }
729
730 #[test]
731 fn a_playhead_moving_on_time_is_not_news() {
732 let sent = LinkState {
733 playing: true,
734 position_ms: 10_000,
735 ..Default::default()
736 };
737 let later = |pos| LinkState {
738 position_ms: pos,
739 ..sent.clone()
740 };
741 let five = Duration::from_secs(5);
742 assert!(!later(15_000).differs(&sent, five));
743 assert!(later(60_000).differs(&sent, five), "a seek");
744 let paused = LinkState {
745 playing: false,
746 ..later(15_000)
747 };
748 assert!(paused.differs(&sent, five));
749 }
750
751 #[test]
752 fn the_url_follows_the_scheme_and_names_the_device() {
753 let identity = LinkIdentity {
754 name: "J's iPhone".into(),
755 platform: "ios".into(),
756 device_id: "abc".into(),
757 };
758 let url = link_url(
759 &SubsonicAuth::new("https://music.example.com", "j", "pw"),
760 &identity,
761 )
762 .unwrap();
763 assert!(url.starts_with("wss://music.example.com/rest/koanLink?"));
764 assert!(url.contains("client=J%27s%20iPhone"));
765 assert!(url.contains("device=abc"));
766 assert!(url.contains("devices=1"));
767 assert!(
768 link_url(&SubsonicAuth::new("http://h:4000", "j", "pw"), &identity)
769 .unwrap()
770 .starts_with("ws://h:4000/")
771 );
772 }
773
774 #[test]
775 fn a_relayed_command_nests_the_command() {
776 let report = LinkReport::Command {
777 to: "phone".into(),
778 command: LinkCommand::HandOff { to: "mac".into() },
779 };
780 let json = serde_json::to_string(&report).unwrap();
781 assert_eq!(
782 json,
783 r#"{"type":"command","to":"phone","command":{"type":"handOff","to":"mac"}}"#
784 );
785 assert_eq!(serde_json::from_str::<LinkReport>(&json).unwrap(), report);
786 }
787
788 #[test]
789 fn an_older_queue_entry_still_reads() {
790 let e: LinkQueueEntry =
791 serde_json::from_str(r#"{"trackId":"7","title":"t","artist":"a","current":true}"#)
792 .unwrap();
793 assert_eq!(e.id, None);
794 assert_eq!(e.duration_ms, 0);
795 }
796
797 #[test]
798 fn strangers_cannot_touch_the_library() {
799 assert!(LinkCommand::Pause.allowed_nearby());
800 assert!(LinkCommand::HandOff { to: "x".into() }.allowed_nearby());
801 assert!(!LinkCommand::Sync { full: true }.allowed_nearby());
802 assert!(!LinkCommand::Evict { track_ids: vec![] }.allowed_nearby());
803 }
804
805 #[test]
806 fn the_device_id_is_kept() {
807 let dir = tempfile::tempdir().unwrap();
808 let first = device_id(dir.path());
809 assert_eq!(device_id(dir.path()), first);
810 }
811}