1use std::collections::HashMap;
11use std::net::TcpStream;
12use std::os::fd::RawFd;
13use std::path::Path;
14use std::sync::atomic::{AtomicU64, Ordering};
15use std::sync::{Arc, LazyLock};
16use std::time::{Duration, Instant};
17
18use parking_lot::{Condvar, Mutex};
19use serde::{Deserialize, Serialize};
20use tungstenite::stream::MaybeTlsStream;
21
22use crate::config::{self, Config};
23use crate::helpers::{subsonic_auth, subsonic_client};
24use crate::player::state::QueueMode;
25use crate::remote::client::SubsonicAuth;
26pub use crate::remote::outputs::{LinkOutput, LinkOutputs, OutputChoice};
27use crate::remote::profile;
28use crate::remote::wire::{self, Waker};
29
30#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
32#[serde(tag = "type", rename_all = "camelCase")]
33pub enum LinkCommand {
34 #[serde(rename_all = "camelCase")]
37 Play {
38 track_ids: Vec<String>,
39 #[serde(default)]
40 start_at: u32,
41 #[serde(default, skip_serializing_if = "is_zero")]
42 position_ms: u64,
43 #[serde(default, skip_serializing_if = "std::ops::Not::not")]
44 paused: bool,
45 #[serde(default, skip_serializing_if = "std::ops::Not::not")]
51 handoff: bool,
52 #[serde(default, skip_serializing_if = "QueueMode::is_keep")]
55 mode: QueueMode,
56 },
57 #[serde(rename_all = "camelCase")]
59 Enqueue {
60 track_ids: Vec<String>,
61 },
62 #[serde(rename_all = "camelCase")]
64 PlayNext {
65 track_ids: Vec<String>,
66 },
67 #[serde(rename_all = "camelCase")]
69 Remove {
70 track_ids: Vec<String>,
71 },
72 Clear,
73 Sync {
77 #[serde(default)]
78 full: bool,
79 },
80 #[serde(rename_all = "camelCase")]
83 Evict {
84 track_ids: Vec<String>,
85 },
86 #[serde(rename_all = "camelCase")]
89 JumpTo {
90 track_id: String,
91 },
92 #[serde(rename_all = "camelCase")]
93 Seek {
94 position_ms: u64,
95 },
96 Pause,
97 Resume,
98 Next,
99 Previous,
100 PlayItem {
104 id: String,
105 },
106 RemoveItems {
107 ids: Vec<String>,
108 },
109 MoveItems {
111 ids: Vec<String>,
112 target: String,
113 after: bool,
114 },
115 #[serde(rename_all = "camelCase")]
117 Insert {
118 track_ids: Vec<String>,
119 after: String,
120 },
121 Undo,
122 Redo,
123 Shuffle {
126 on: bool,
127 },
128 Repeat {
130 mode: crate::player::state::Repeat,
131 },
132 SleepTimer {
134 timer: Option<crate::player::state::SleepTimer>,
135 },
136 HandOff {
140 to: String,
141 },
142 Devices {
145 devices: Vec<LinkDevice>,
146 },
147 DeviceKeys {
153 keys: Vec<LinkDeviceKey>,
154 },
155 Shares {
158 grantees: Vec<String>,
159 #[serde(default, skip_serializing_if = "Option::is_none")]
161 error: Option<String>,
162 #[serde(default, skip_serializing_if = "Vec::is_empty")]
164 accounts: Vec<String>,
165 },
166 Shared {
169 command: Box<LinkCommand>,
170 },
171 Forgotten {
174 device: String,
175 },
176 HistoryChanged,
180 DspProfilesChanged,
184 SetOutput {
187 output: OutputChoice,
188 },
189 RefreshOutputs,
192 SetRendererVolume {
194 volume: u8,
195 },
196 SetPreset {
199 device: String,
200 profile: Option<String>,
201 },
202 WatchLevels {
206 on: bool,
207 },
208 Levels {
211 from: String,
212 f: crate::remote::levels::Frame,
213 },
214 Acked {
217 from: String,
218 ack: u64,
219 outcome: crate::remote::acks::AckOutcome,
220 },
221}
222
223fn is_zero(n: &u64) -> bool {
224 *n == 0
225}
226
227impl LinkCommand {
228 pub fn allowed_playback(&self) -> bool {
239 match self {
240 Self::Play { .. }
241 | Self::Enqueue { .. }
242 | Self::PlayNext { .. }
243 | Self::Remove { .. }
244 | Self::Clear
245 | Self::JumpTo { .. }
246 | Self::Seek { .. }
247 | Self::Pause
248 | Self::Resume
249 | Self::Next
250 | Self::Previous
251 | Self::PlayItem { .. }
252 | Self::RemoveItems { .. }
253 | Self::MoveItems { .. }
254 | Self::Insert { .. }
255 | Self::Undo
256 | Self::Redo
257 | Self::Shuffle { .. }
258 | Self::Repeat { .. }
259 | Self::SleepTimer { .. }
260 | Self::HandOff { .. }
261 | Self::SetOutput { .. }
262 | Self::RefreshOutputs
263 | Self::SetRendererVolume { .. }
264 | Self::SetPreset { .. }
265 | Self::WatchLevels { .. } => true,
266 Self::Sync { .. }
268 | Self::Evict { .. }
269 | Self::Devices { .. }
270 | Self::DeviceKeys { .. }
271 | Self::Shares { .. }
272 | Self::Shared { .. }
273 | Self::Forgotten { .. }
274 | Self::HistoryChanged
275 | Self::DspProfilesChanged
276 | Self::Levels { .. }
277 | Self::Acked { .. } => false,
278 }
279 }
280
281 pub fn from_the_network(&self, full: bool) -> Option<CommandSource> {
285 if full {
286 self.allowed_playback().then_some(CommandSource::Nearby)
287 } else {
288 self.allowed_nearby().then_some(CommandSource::Stranger)
289 }
290 }
291
292 pub fn relayable(&self) -> bool {
299 match self {
300 Self::Play { .. }
301 | Self::Enqueue { .. }
302 | Self::PlayNext { .. }
303 | Self::Remove { .. }
304 | Self::Clear
305 | Self::Sync { .. }
306 | Self::Evict { .. }
307 | Self::JumpTo { .. }
308 | Self::Seek { .. }
309 | Self::Pause
310 | Self::Resume
311 | Self::Next
312 | Self::Previous
313 | Self::PlayItem { .. }
314 | Self::RemoveItems { .. }
315 | Self::MoveItems { .. }
316 | Self::Insert { .. }
317 | Self::Undo
318 | Self::Redo
319 | Self::Shuffle { .. }
320 | Self::Repeat { .. }
321 | Self::SleepTimer { .. }
322 | Self::HandOff { .. }
323 | Self::SetOutput { .. }
324 | Self::RefreshOutputs
325 | Self::SetRendererVolume { .. }
326 | Self::SetPreset { .. } => true,
327 Self::Devices { .. }
328 | Self::DeviceKeys { .. }
329 | Self::Acked { .. }
330 | Self::Shares { .. }
331 | Self::Shared { .. }
332 | Self::Forgotten { .. }
333 | Self::HistoryChanged
334 | Self::DspProfilesChanged
335 | Self::WatchLevels { .. }
336 | Self::Levels { .. } => false,
337 }
338 }
339
340 pub fn live_only(&self) -> bool {
344 match self {
345 Self::Shared { command } => command.live_only(),
346 cmd => matches!(cmd, Self::WatchLevels { .. } | Self::RefreshOutputs),
347 }
348 }
349
350 pub fn allowed_nearby(&self) -> bool {
357 !matches!(
358 self,
359 Self::Sync { .. }
360 | Self::Evict { .. }
361 | Self::Devices { .. }
362 | Self::DeviceKeys { .. }
363 | Self::Forgotten { .. }
364 | Self::HistoryChanged
365 | Self::DspProfilesChanged
366 | Self::Levels { .. }
367 | Self::Acked { .. }
368 | Self::SetOutput { .. }
369 | Self::RefreshOutputs
370 | Self::SetRendererVolume { .. }
371 | Self::SetPreset { .. }
372 | Self::Shares { .. }
373 | Self::Shared { .. }
374 )
375 }
376
377 pub fn track_ids(&self) -> &[String] {
379 match self {
380 Self::Play { track_ids, .. }
381 | Self::Enqueue { track_ids }
382 | Self::PlayNext { track_ids }
383 | Self::Remove { track_ids }
384 | Self::Evict { track_ids }
385 | Self::Insert { track_ids, .. } => track_ids,
386 Self::JumpTo { track_id } => std::slice::from_ref(track_id),
387 Self::Shared { command } => command.track_ids(),
388 _ => &[],
389 }
390 }
391}
392
393#[derive(Debug, Clone, Copy, PartialEq, Eq)]
395pub enum CommandSource {
396 Account,
399 Shared,
402 Nearby,
405 Stranger,
409}
410
411#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
413pub struct LinkDeviceKey {
414 pub id: String,
416 pub key: String,
418 #[serde(default, skip_serializing_if = "Option::is_none")]
421 pub owner: Option<String>,
422}
423
424pub fn valid_device_key(key: &str) -> bool {
427 use base64::Engine as _;
428 base64::engine::general_purpose::STANDARD
429 .decode(key)
430 .is_ok_and(|bytes| bytes.len() == 32)
431}
432
433#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
435#[serde(rename_all = "camelCase")]
436pub struct LinkDevice {
437 pub id: String,
439 pub name: String,
440 pub platform: String,
441 pub linked: bool,
444 pub state: Option<LinkState>,
447 #[serde(default, skip_serializing_if = "Option::is_none")]
449 pub last_seen: Option<i64>,
450 #[serde(default, skip_serializing_if = "Option::is_none")]
454 pub wakeable: Option<bool>,
455 #[serde(default, skip_serializing_if = "Option::is_none")]
458 pub owner: Option<String>,
459 #[serde(default, skip_serializing_if = "std::ops::Not::not")]
462 pub acks: bool,
463}
464
465impl LinkCommand {
466 pub fn track_ids_mut(&mut self) -> Vec<&mut String> {
469 match self {
470 Self::Play { track_ids, .. }
471 | Self::Enqueue { track_ids }
472 | Self::PlayNext { track_ids }
473 | Self::Remove { track_ids }
474 | Self::Evict { track_ids }
475 | Self::Insert { track_ids, .. } => track_ids.iter_mut().collect(),
476 Self::JumpTo { track_id } => vec![track_id],
477 Self::Shared { command } => command.track_ids_mut(),
478 Self::PlayItem { .. }
479 | Self::RemoveItems { .. }
480 | Self::MoveItems { .. }
481 | Self::Undo
482 | Self::Redo
483 | Self::Shuffle { .. }
484 | Self::Repeat { .. }
485 | Self::SleepTimer { .. }
486 | Self::HandOff { .. }
487 | Self::Devices { .. }
488 | Self::DeviceKeys { .. }
489 | Self::Shares { .. }
490 | Self::Forgotten { .. }
491 | Self::HistoryChanged
492 | Self::DspProfilesChanged
493 | Self::WatchLevels { .. }
494 | Self::Levels { .. }
495 | Self::Acked { .. }
496 | Self::SetOutput { .. }
497 | Self::RefreshOutputs
498 | Self::SetRendererVolume { .. }
499 | Self::SetPreset { .. }
500 | Self::Clear
501 | Self::Sync { .. }
502 | Self::Seek { .. }
503 | Self::Pause
504 | Self::Resume
505 | Self::Next
506 | Self::Previous => Vec::new(),
507 }
508 }
509}
510
511#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
513#[serde(rename_all = "camelCase")]
514pub struct LinkState {
515 pub playing: bool,
516 pub title: Option<String>,
517 pub artist: Option<String>,
518 #[serde(default)]
519 pub album: Option<String>,
520 #[serde(default)]
522 pub position_ms: u64,
523 #[serde(default)]
524 pub duration_ms: u64,
525 #[serde(default)]
527 pub queue: Vec<LinkQueueEntry>,
528 #[serde(default, skip_serializing_if = "Option::is_none")]
530 pub outputs: Option<LinkOutputs>,
531 #[serde(default, skip_serializing_if = "std::ops::Not::not")]
532 pub shuffle: bool,
533 #[serde(default, skip_serializing_if = "crate::player::state::Repeat::is_off")]
534 pub repeat: crate::player::state::Repeat,
535 #[serde(default, skip_serializing_if = "Option::is_none")]
536 pub sleep: Option<crate::player::state::Sleep>,
537 #[serde(default, skip_serializing_if = "std::ops::Not::not")]
539 pub sleep_fading: bool,
540}
541
542#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
543#[serde(rename_all = "camelCase")]
544pub struct LinkQueueEntry {
545 #[serde(default)]
547 pub id: Option<String>,
548 pub track_id: Option<String>,
550 pub title: String,
551 pub artist: String,
552 #[serde(default)]
553 pub album: String,
554 #[serde(default)]
555 pub duration_ms: u64,
556 pub current: bool,
557}
558
559impl LinkState {
560 pub fn differs(&self, sent: &LinkState, elapsed: Duration) -> bool {
564 let strip = |s: &LinkState| LinkState {
565 position_ms: 0,
566 ..s.clone()
567 };
568 if strip(self) != strip(sent) {
569 return true;
570 }
571 let expected = if sent.playing {
572 sent.position_ms + elapsed.as_millis() as u64
573 } else {
574 sent.position_ms
575 };
576 self.position_ms.abs_diff(expected) > 3000
577 }
578}
579
580#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
582#[serde(tag = "type", rename_all = "camelCase")]
583pub enum LinkReport {
584 State(LinkState),
585 Push {
589 token: String,
590 sandbox: bool,
591 },
592 Command {
596 to: String,
597 command: LinkCommand,
598 #[serde(default, skip_serializing_if = "Option::is_none")]
599 ack: Option<u64>,
600 },
601 Received {
605 ack: u64,
606 },
607 Ack {
609 ack: u64,
610 outcome: crate::remote::acks::AckOutcome,
611 },
612 Activity {
615 token: Option<String>,
616 device: Option<String>,
617 #[serde(default)]
618 sandbox: bool,
619 },
620 Hello(LinkHello),
622 Wake {
625 to: String,
626 #[serde(default)]
627 notify: bool,
628 },
629 Share {
632 grantee: String,
633 allow: bool,
634 },
635 Forget {
638 device: String,
639 },
640 Levels {
643 f: crate::remote::levels::Frame,
644 },
645}
646
647#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
650#[serde(rename_all = "camelCase")]
651pub struct LinkHello {
652 pub id: String,
653 pub name: String,
654 pub platform: String,
655 pub library: Option<String>,
659 #[serde(default)]
662 pub acks: bool,
663 #[serde(default, skip_serializing_if = "Option::is_none")]
666 pub nonce: Option<String>,
667}
668
669pub fn library_fingerprint(cfg: &Config) -> Option<String> {
672 let auth = subsonic_auth(cfg)?;
673 let url = auth.base_url.trim_end_matches('/').to_ascii_lowercase();
674 Some(format!("{:x}", md5::compute(url.as_bytes())))
675}
676
677static PUSH_TOKEN: Mutex<Option<(String, bool)>> = Mutex::new(None);
680
681type ActivityToken = (String, String, bool);
685static ACTIVITY: Mutex<Option<Option<ActivityToken>>> = Mutex::new(None);
686
687pub fn set_activity(activity: Option<ActivityToken>) {
690 *ACTIVITY.lock() = Some(activity);
691 if let Some(up) = LINK.lock().as_ref() {
692 up.waker.wake();
693 }
694}
695
696pub fn parse_command(json: &str) -> Result<LinkCommand, String> {
699 serde_json::from_str(json).map_err(|e| e.to_string())
700}
701
702pub fn set_push_token(token: String, sandbox: bool) {
704 *PUSH_TOKEN.lock() = Some((token, sandbox));
705 nudge();
706 if let Some(up) = LINK.lock().as_ref() {
707 up.waker.wake();
708 }
709}
710
711#[derive(Debug, Clone)]
713pub struct LinkIdentity {
714 pub name: String,
716 pub platform: String,
718 pub device_id: String,
721}
722
723impl LinkIdentity {
724 pub fn this_device(name: Option<String>) -> Self {
726 let (platform, label) = if cfg!(target_os = "ios") {
727 ("ios", "iPhone")
728 } else if cfg!(target_os = "tvos") {
729 ("tvos", "Apple TV")
730 } else if cfg!(target_os = "macos") {
731 ("macos", "Mac")
732 } else {
733 ("linux", "Linux")
734 };
735 Self {
736 name: name
737 .filter(|n| !n.trim().is_empty())
738 .or_else(hostname)
739 .unwrap_or_else(|| label.to_string()),
740 platform: platform.to_string(),
741 device_id: device_id(&config::config_dir()),
742 }
743 }
744}
745
746#[derive(Clone)]
750pub struct Local {
751 pub identity: LinkIdentity,
752 pub state: Arc<dyn Fn() -> LinkState + Send + Sync>,
753 pub on_command:
756 Arc<dyn Fn(LinkCommand, CommandSource, Option<crate::remote::acks::Pending>) + Send + Sync>,
757}
758
759pub fn spawn(local: Local) {
768 std::thread::Builder::new()
769 .name("koan-link".into())
770 .spawn(move || run(local))
771 .expect("failed to spawn the link thread");
772}
773
774const RETRY_MIN: Duration = Duration::from_secs(2);
775const RETRY_MAX: Duration = Duration::from_secs(60);
776
777static SIGN_IN: AtomicU64 = AtomicU64::new(0);
780
781pub fn relink() {
785 SIGN_IN.fetch_add(1, Ordering::SeqCst);
786 if let Some(up) = LINK.lock().as_ref() {
787 up.waker.wake();
788 }
789 nudge();
790}
791
792fn run(local: Local) {
793 let mut wait = RETRY_MIN;
794 loop {
795 crate::quiet::wait_until_awake();
796 let signed_in = SIGN_IN.load(Ordering::SeqCst);
799 let cfg = Config::load().unwrap_or_default();
800 let account = crate::remote::proof::account_of(&cfg);
801 let Some(auth) = subsonic_auth(&cfg) else {
802 rest(RETRY_MAX);
803 continue;
804 };
805 match profile::for_auth(&auth) {
806 Some(p) if p.links() => {}
807 Some(_) => {
809 rest(RETRY_MAX);
810 continue;
811 }
812 None => {
813 rest(wait);
814 wait = (wait * 2).min(RETRY_MAX);
815 continue;
816 }
817 }
818
819 match connect(&auth, &local.identity) {
820 Ok((mut socket, fd)) => {
821 log::info!("link: connected to {}", auth.base_url);
822 if let Some(client) = crate::helpers::subsonic_client(&cfg) {
825 client.outage().retry_now();
826 }
827 wait = RETRY_MIN;
828 if let Err(e) = serve(&mut socket, fd, &local, signed_in, account) {
829 log::info!("link: closed: {e}");
830 }
831 *LINK.lock() = None;
832 crate::remote::devices::set_linked(false);
833 }
834 Err(e) => {
835 log::warn!("link: {e}");
836 profile::forget();
842 }
843 }
844 if rest(wait) {
845 wait = RETRY_MIN;
846 continue;
847 }
848 wait = (wait * 2).min(RETRY_MAX);
849 }
850}
851
852static NUDGE: (Mutex<bool>, Condvar) = (Mutex::new(false), Condvar::new());
853
854pub fn nudge() {
860 *NUDGE.0.lock() = true;
861 NUDGE.1.notify_all();
862}
863
864fn rest(d: Duration) -> bool {
866 let mut nudged = NUDGE.0.lock();
867 if !*nudged {
868 NUDGE.1.wait_for(&mut nudged, d);
869 }
870 std::mem::replace(&mut *nudged, false)
871}
872
873pub fn hang_up() {
875 if let Some(up) = LINK.lock().as_ref() {
876 up.waker.wake();
877 }
878}
879
880struct Up {
882 waker: Arc<Waker>,
883 outbox: Vec<LinkReport>,
884}
885
886static LINK: Mutex<Option<Up>> = Mutex::new(None);
887
888pub fn report(report: LinkReport) -> bool {
890 let mut link = LINK.lock();
891 let Some(up) = link.as_mut() else {
892 return false;
893 };
894 up.outbox.push(report);
895 up.waker.wake();
896 true
897}
898
899type Socket = tungstenite::WebSocket<MaybeTlsStream<TcpStream>>;
900
901fn connect(auth: &SubsonicAuth, identity: &LinkIdentity) -> Result<(Socket, RawFd), String> {
902 let url = link_url(auth, identity)?;
903 let (socket, _) = tungstenite::connect(url).map_err(|e| e.to_string())?;
904 let fd = wire::prepare(socket.get_ref())?;
905 Ok((socket, fd))
906}
907
908fn link_url(auth: &SubsonicAuth, identity: &LinkIdentity) -> Result<String, String> {
910 let base = if let Some(rest) = auth.base_url.strip_prefix("https://") {
911 format!("wss://{rest}")
912 } else if let Some(rest) = auth.base_url.strip_prefix("http://") {
913 format!("ws://{rest}")
914 } else {
915 return Err(format!("not an http(s) server: {}", auth.base_url));
916 };
917 let mut query = auth.query().map_err(|e| e.to_string())?;
918 let device_key = crate::remote::proof::public_key().unwrap_or_default();
922 for (k, v) in [
923 ("client", identity.name.as_str()),
924 ("platform", identity.platform.as_str()),
925 ("device", identity.device_id.as_str()),
926 ("devices", "1"),
929 ("acks", "1"),
932 ("deviceKey", device_key.as_str()),
933 ] {
934 if v.is_empty() {
935 continue;
936 }
937 query.push('&');
938 query.push_str(k);
939 query.push('=');
940 query.push_str(&percent_encode(v));
941 }
942 Ok(format!("{base}/rest/koanLink?{query}"))
943}
944
945fn serve(
946 socket: &mut Socket,
947 fd: RawFd,
948 local: &Local,
949 signed_in: u64,
950 account: Option<String>,
951) -> Result<(), String> {
952 let waker = Waker::new().map_err(|e| e.to_string())?;
953 let watcher = waker.clone();
954 wire::wake_on_engine_change(&waker);
955 *LINK.lock() = Some(Up {
956 waker: waker.clone(),
957 outbox: Vec::new(),
958 });
959 crate::remote::devices::set_linked(true);
960 let mut session = LinkSession {
961 local,
962 sent: None,
963 sent_push: None,
964 sent_activity: None,
965 waker: watcher,
966 levels: None,
967 signed_in,
968 account,
969 };
970 wire::drive(socket, fd, &waker, &mut session)
971}
972
973struct LinkSession<'a> {
974 local: &'a Local,
975 sent: Option<(LinkState, Instant)>,
976 sent_push: Option<(String, bool)>,
977 sent_activity: Option<Option<ActivityToken>>,
978 waker: Arc<Waker>,
979 levels: Option<crate::remote::levels::Watch>,
982 signed_in: u64,
984 account: Option<String>,
987}
988
989impl wire::Session for LinkSession<'_> {
990 fn outgoing(&mut self) -> Vec<String> {
991 let mut out = Vec::new();
992 let push = PUSH_TOKEN.lock().clone();
993 if let Some((token, sandbox)) = push.clone()
994 && push != self.sent_push
995 {
996 out.push(LinkReport::Push { token, sandbox });
997 self.sent_push = push;
998 }
999 let activity = ACTIVITY.lock().clone();
1000 if activity.is_some() && activity != self.sent_activity {
1001 let (token, device, sandbox) = match activity.clone().flatten() {
1002 Some((t, d, s)) => (Some(t), Some(d), s),
1003 None => (None, None, false),
1004 };
1005 out.push(LinkReport::Activity {
1006 token,
1007 device,
1008 sandbox,
1009 });
1010 self.sent_activity = activity;
1011 }
1012 if let Some(up) = LINK.lock().as_mut() {
1013 out.append(&mut up.outbox);
1014 }
1015 let now = (self.local.state)();
1016 if self
1017 .sent
1018 .as_ref()
1019 .is_none_or(|(s, at)| now.differs(s, at.elapsed()))
1020 {
1021 out.push(LinkReport::State(now.clone()));
1022 self.sent = Some((now, Instant::now()));
1023 }
1024 if let Some(f) = self.levels.as_mut().and_then(|w| w.take()) {
1025 out.push(LinkReport::Levels { f });
1026 }
1027 out.iter()
1028 .filter_map(|r| serde_json::to_string(r).ok())
1029 .collect()
1030 }
1031
1032 fn incoming(&mut self, text: &str) {
1033 let envelope = match serde_json::from_str::<crate::remote::acks::Envelope>(text) {
1034 Ok(envelope) => envelope,
1035 Err(e) => {
1036 log::warn!("link: not a command ({e}): {text}");
1037 return;
1038 }
1039 };
1040 if let Some(ack) = envelope.ack {
1041 report(LinkReport::Received { ack });
1042 }
1043 let Some((command, pending)) = crate::remote::acks::take(envelope, |ack, outcome| {
1045 report(LinkReport::Ack { ack, outcome });
1046 }) else {
1047 return;
1048 };
1049 match command {
1050 LinkCommand::Acked { from, ack, outcome } => {
1051 crate::remote::acks::resolve(ack, &from, outcome);
1052 }
1053 LinkCommand::Devices { devices } => {
1054 crate::remote::devices::set_account(devices);
1055 }
1056 LinkCommand::DeviceKeys { keys } => {
1057 crate::remote::proof::keep(keys, self.account.clone());
1058 }
1059 LinkCommand::Shares {
1060 grantees,
1061 error,
1062 accounts,
1063 } => {
1064 crate::remote::devices::set_shares(grantees, error, accounts);
1065 }
1066 LinkCommand::Shared { command } => {
1067 if command.allowed_playback() {
1070 (self.local.on_command)(*command, CommandSource::Shared, pending);
1071 } else {
1072 log::warn!("link: refused from a shared account: {command:?}");
1073 if let Some(pending) = pending {
1074 pending.finish(crate::remote::acks::AckOutcome::Refused {
1075 reason: "not allowed from another account".into(),
1076 });
1077 }
1078 }
1079 }
1080 LinkCommand::Forgotten { device } => {
1081 crate::remote::devices::forgotten(&device);
1082 }
1083 LinkCommand::WatchLevels { on } => {
1084 self.levels = on.then(|| crate::remote::levels::feed().watch(&self.waker));
1085 }
1086 LinkCommand::Levels { from, f } => {
1087 crate::remote::levels::remote().received(&from, f);
1088 }
1089 cmd => (self.local.on_command)(cmd, CommandSource::Account, pending),
1090 }
1091 }
1092
1093 fn done(&self) -> bool {
1094 !crate::quiet::awake() || SIGN_IN.load(Ordering::SeqCst) != self.signed_in
1095 }
1096}
1097
1098pub fn resolve_tracks(
1109 db: &crate::db::connection::Database,
1110 remote_ids: &[String],
1111 may_sync: bool,
1112) -> (Vec<i64>, bool) {
1113 let lookup = |db: &crate::db::connection::Database| {
1114 let mut stmt = db
1115 .conn
1116 .prepare_cached(
1117 "SELECT id FROM tracks WHERE uid = ?1
1118 UNION ALL SELECT id FROM tracks WHERE remote_id = ?1 LIMIT 1",
1119 )
1120 .ok();
1121 remote_ids
1122 .iter()
1123 .map(|rid| {
1124 stmt.as_mut()
1125 .and_then(|s| s.query_row([rid], |r| r.get::<_, i64>(0)).ok())
1126 })
1127 .collect::<Vec<_>>()
1128 };
1129 let found = lookup(db);
1130 let missing: Vec<&String> = remote_ids
1131 .iter()
1132 .zip(&found)
1133 .filter(|(_, f)| f.is_none())
1134 .map(|(id, _)| id)
1135 .collect();
1136 if missing.is_empty() || !may_sync || missing.iter().all(|id| recently_missed(id)) {
1137 return (found.into_iter().flatten().collect(), false);
1138 }
1139 sync(db, crate::helpers::Walk::IfChanged);
1140 let found = lookup(db);
1141 let mut missed = MISSED.lock();
1142 let now = Instant::now();
1143 missed.retain(|_, at| now.duration_since(*at) < MISS_TTL);
1144 for (id, _) in remote_ids.iter().zip(&found).filter(|(_, f)| f.is_none()) {
1145 missed.insert(id.clone(), now);
1146 }
1147 (found.into_iter().flatten().collect(), true)
1148}
1149
1150static MISSED: LazyLock<Mutex<HashMap<String, Instant>>> = LazyLock::new(Default::default);
1152const MISS_TTL: Duration = Duration::from_secs(300);
1153
1154fn recently_missed(id: &str) -> bool {
1155 MISSED
1156 .lock()
1157 .get(id)
1158 .is_some_and(|at| at.elapsed() < MISS_TTL)
1159}
1160
1161pub fn sync(db: &crate::db::connection::Database, walk: crate::helpers::Walk) {
1164 let cfg = Config::load().unwrap_or_default();
1165 if let Some(client) = subsonic_client(&cfg)
1166 && let Err(e) = crate::helpers::sync_remote(
1167 db,
1168 &client,
1169 walk,
1170 &cfg.remote.url,
1171 &cfg.remote.username,
1172 &|_| {},
1173 )
1174 {
1175 log::warn!("link: sync failed: {e}");
1176 }
1177}
1178
1179fn device_id(dir: &Path) -> String {
1183 #[cfg(any(target_os = "ios", target_os = "tvos"))]
1184 {
1185 use security_framework::passwords::{get_generic_password, set_generic_password};
1186 const SERVICE: &str = "cc.blit.koan.link";
1187 if let Some(id) = get_generic_password(SERVICE, "device-id")
1188 .ok()
1189 .and_then(|b| String::from_utf8(b).ok())
1190 .filter(|id| !id.trim().is_empty())
1191 {
1192 return id;
1193 }
1194 let id = file_device_id(dir);
1195 let _ = set_generic_password(SERVICE, "device-id", id.as_bytes());
1196 id
1197 }
1198 #[cfg(not(any(target_os = "ios", target_os = "tvos")))]
1199 file_device_id(dir)
1200}
1201
1202fn file_device_id(dir: &Path) -> String {
1203 let path = dir.join("device-id");
1204 if let Ok(id) = std::fs::read_to_string(&path) {
1205 let id = id.trim();
1206 if !id.is_empty() {
1207 return id.to_string();
1208 }
1209 }
1210 let id = uuid::Uuid::now_v7().to_string();
1211 let _ = std::fs::create_dir_all(dir);
1212 let _ = std::fs::write(&path, &id);
1213 id
1214}
1215
1216fn hostname() -> Option<String> {
1217 let mut buf = [0u8; 256];
1218 let ok = unsafe { libc::gethostname(buf.as_mut_ptr().cast(), buf.len()) } == 0;
1221 if !ok {
1222 return None;
1223 }
1224 let end = buf.iter().position(|&b| b == 0).unwrap_or(buf.len());
1225 let name = String::from_utf8_lossy(&buf[..end]);
1226 let name = name.trim_end_matches(".local").trim();
1227 (!name.is_empty() && name != "localhost").then(|| name.to_string())
1228}
1229
1230fn percent_encode(s: &str) -> String {
1231 let mut out = String::with_capacity(s.len());
1232 for b in s.bytes() {
1233 if b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_' | b'.' | b'~') {
1234 out.push(b as char);
1235 } else {
1236 out.push_str(&format!("%{b:02X}"));
1237 }
1238 }
1239 out
1240}
1241
1242#[cfg(test)]
1243mod tests {
1244
1245 #[test]
1246 fn levels_cross_the_link_compactly() {
1247 use crate::remote::levels::Frame;
1248 let f = Frame(61_250, 512, 300, 40);
1249
1250 let report = serde_json::to_string(&LinkReport::Levels { f }).unwrap();
1251 assert_eq!(report, r#"{"type":"levels","f":[61250,512,300,40]}"#);
1252 assert_eq!(
1253 serde_json::from_str::<LinkReport>(&report).unwrap(),
1254 LinkReport::Levels { f }
1255 );
1256
1257 let relayed = LinkCommand::Levels {
1258 from: "dev-phone".into(),
1259 f,
1260 };
1261 let text = serde_json::to_string(&relayed).unwrap();
1262 assert!(text.len() < 80, "{} bytes: {text}", text.len());
1263 assert_eq!(serde_json::from_str::<LinkCommand>(&text).unwrap(), relayed);
1264
1265 let watch = LinkCommand::WatchLevels { on: true };
1266 let text = serde_json::to_string(&watch).unwrap();
1267 assert_eq!(serde_json::from_str::<LinkCommand>(&text).unwrap(), watch);
1268 assert!(watch.allowed_nearby(), "a stranger may watch the bars");
1269 assert!(!relayed.allowed_nearby(), "only the server relays frames");
1270 }
1271 use super::*;
1272
1273 #[test]
1274 fn a_sleep_timer_travels_as_playback_and_comes_back_in_the_state() {
1275 use crate::player::state::{Sleep, SleepTimer};
1276 let cmd: LinkCommand =
1277 serde_json::from_str(r#"{"type":"sleepTimer","timer":{"kind":"after","minutes":30}}"#)
1278 .unwrap();
1279 assert_eq!(
1280 cmd,
1281 LinkCommand::SleepTimer {
1282 timer: Some(SleepTimer::After { minutes: 30 })
1283 }
1284 );
1285 let cancel: LinkCommand =
1286 serde_json::from_str(r#"{"type":"sleepTimer","timer":null}"#).unwrap();
1287 assert_eq!(cancel, LinkCommand::SleepTimer { timer: None });
1288 for cmd in [cmd, cancel] {
1289 assert!(cmd.allowed_playback() && cmd.allowed_nearby(), "{cmd:?}");
1290 }
1291
1292 let state = LinkState {
1293 sleep: Some(Sleep::EndOfRecord),
1294 ..Default::default()
1295 };
1296 let json = serde_json::to_string(&state).unwrap();
1297 assert!(json.contains(r#""sleep":{"kind":"endOfRecord"}"#), "{json}");
1298 assert_eq!(serde_json::from_str::<LinkState>(&json).unwrap(), state);
1299 assert!(
1300 !serde_json::to_string(&LinkState::default())
1301 .unwrap()
1302 .contains("sleep"),
1303 "nothing said with none set"
1304 );
1305 }
1306
1307 #[test]
1311 fn the_playback_set_is_playback_and_nothing_of_the_library() {
1312 let ids = vec!["t".to_string()];
1313 for cmd in [
1314 LinkCommand::Pause,
1315 LinkCommand::JumpTo {
1316 track_id: "t".into(),
1317 },
1318 LinkCommand::Enqueue {
1319 track_ids: ids.clone(),
1320 },
1321 LinkCommand::HandOff { to: "x".into() },
1322 LinkCommand::SetRendererVolume { volume: 1 },
1323 ] {
1324 assert!(cmd.allowed_playback(), "{cmd:?}");
1325 assert_eq!(cmd.from_the_network(true), Some(CommandSource::Nearby));
1326 }
1327 for cmd in [
1328 LinkCommand::Sync { full: false },
1329 LinkCommand::Evict {
1330 track_ids: ids.clone(),
1331 },
1332 LinkCommand::Devices { devices: vec![] },
1333 LinkCommand::Shares {
1334 grantees: vec![],
1335 error: None,
1336 accounts: vec![],
1337 },
1338 LinkCommand::Shared {
1339 command: Box::new(LinkCommand::Pause),
1340 },
1341 ] {
1342 assert!(!cmd.allowed_playback(), "{cmd:?}");
1343 assert_eq!(cmd.from_the_network(true), None, "{cmd:?}");
1344 assert_eq!(cmd.from_the_network(false), None, "{cmd:?}");
1345 }
1346 assert_eq!(
1348 LinkCommand::Pause.from_the_network(false),
1349 Some(CommandSource::Stranger)
1350 );
1351 assert_eq!(
1352 LinkCommand::SetRendererVolume { volume: 1 }.from_the_network(false),
1353 None
1354 );
1355 }
1356
1357 #[test]
1360 fn unknown_ids_sync_only_when_allowed_and_not_recently_missed() {
1361 let dir = tempfile::tempdir().unwrap();
1362 let db = crate::db::connection::Database::open(&dir.path().join("koan.db")).unwrap();
1363
1364 let (found, synced) = resolve_tracks(&db, &["from-a-stranger".into()], false);
1365 assert!(found.is_empty());
1366 assert!(!synced, "a nearby peer's unknown id must not start a sync");
1367
1368 MISSED
1369 .lock()
1370 .insert("deleted-on-server".into(), Instant::now());
1371 let (_, synced) = resolve_tracks(&db, &["deleted-on-server".into()], true);
1372 assert!(
1373 !synced,
1374 "an id a sync just failed to find does not start another"
1375 );
1376 }
1377
1378 #[test]
1379 fn commands_name_their_tracks() {
1380 let play: LinkCommand =
1381 serde_json::from_str(r#"{"type":"play","trackIds":["a","b"]}"#).unwrap();
1382 assert_eq!(play.track_ids(), ["a", "b"]);
1383 let jump: LinkCommand = serde_json::from_str(r#"{"type":"jumpTo","trackId":"c"}"#).unwrap();
1384 assert_eq!(jump.track_ids(), ["c"]);
1385 assert!(LinkCommand::Pause.track_ids().is_empty());
1386 }
1387
1388 #[test]
1389 fn a_push_token_is_a_tagged_report() {
1390 let report = LinkReport::Push {
1391 token: "ab12".into(),
1392 sandbox: true,
1393 };
1394 let text = serde_json::to_string(&report).unwrap();
1395 assert_eq!(text, r#"{"type":"push","token":"ab12","sandbox":true}"#);
1396 assert_eq!(serde_json::from_str::<LinkReport>(&text).unwrap(), report);
1397 }
1398
1399 #[test]
1400 fn commands_are_tagged_json() {
1401 let play = LinkCommand::Play {
1402 track_ids: vec!["12".into(), "34".into()],
1403 start_at: 1,
1404 position_ms: 0,
1405 paused: false,
1406 handoff: false,
1407 mode: QueueMode::Keep,
1408 };
1409 let json = serde_json::to_string(&play).unwrap();
1410 assert_eq!(
1411 json,
1412 r#"{"type":"play","trackIds":["12","34"],"startAt":1}"#
1413 );
1414 assert_eq!(serde_json::from_str::<LinkCommand>(&json).unwrap(), play);
1415 let held = LinkCommand::Play {
1416 track_ids: vec!["12".into()],
1417 start_at: 0,
1418 position_ms: 61_250,
1419 paused: true,
1420 handoff: false,
1421 mode: QueueMode::Keep,
1422 };
1423 let json = serde_json::to_string(&held).unwrap();
1424 assert_eq!(
1425 json,
1426 r#"{"type":"play","trackIds":["12"],"startAt":0,"positionMs":61250,"paused":true}"#
1427 );
1428 assert_eq!(serde_json::from_str::<LinkCommand>(&json).unwrap(), held);
1429 assert_eq!(
1430 serde_json::from_str::<LinkCommand>(r#"{"type":"pause"}"#).unwrap(),
1431 LinkCommand::Pause
1432 );
1433 let report = LinkReport::State(LinkState {
1434 playing: true,
1435 title: Some("Portions for Foxes".into()),
1436 ..Default::default()
1437 });
1438 let json = serde_json::to_string(&report).unwrap();
1439 assert!(json.starts_with(r#"{"type":"state","playing":true,"title":"Portions for Foxes""#));
1440 assert_eq!(serde_json::from_str::<LinkReport>(&json).unwrap(), report);
1441 }
1442
1443 #[test]
1444 fn a_playhead_moving_on_time_is_not_news() {
1445 let sent = LinkState {
1446 playing: true,
1447 position_ms: 10_000,
1448 ..Default::default()
1449 };
1450 let later = |pos| LinkState {
1451 position_ms: pos,
1452 ..sent.clone()
1453 };
1454 let five = Duration::from_secs(5);
1455 assert!(!later(15_000).differs(&sent, five));
1456 assert!(later(60_000).differs(&sent, five), "a seek");
1457 let paused = LinkState {
1458 playing: false,
1459 ..later(15_000)
1460 };
1461 assert!(paused.differs(&sent, five));
1462 }
1463
1464 #[test]
1465 fn the_url_follows_the_scheme_and_names_the_device() {
1466 let identity = LinkIdentity {
1467 name: "J's iPhone".into(),
1468 platform: "ios".into(),
1469 device_id: "abc".into(),
1470 };
1471 let url = link_url(
1472 &SubsonicAuth::new("https://music.example.com", "j", "pw"),
1473 &identity,
1474 )
1475 .unwrap();
1476 assert!(url.starts_with("wss://music.example.com/rest/koanLink?"));
1477 assert!(url.contains("client=J%27s%20iPhone"));
1478 assert!(url.contains("device=abc"));
1479 assert!(url.contains("devices=1"));
1480 assert!(
1481 link_url(&SubsonicAuth::new("http://h:4000", "j", "pw"), &identity)
1482 .unwrap()
1483 .starts_with("ws://h:4000/")
1484 );
1485 }
1486
1487 #[test]
1488 fn a_relayed_command_nests_the_command() {
1489 let report = LinkReport::Command {
1490 to: "phone".into(),
1491 command: LinkCommand::HandOff { to: "mac".into() },
1492 ack: None,
1493 };
1494 let json = serde_json::to_string(&report).unwrap();
1495 assert_eq!(
1496 json,
1497 r#"{"type":"command","to":"phone","command":{"type":"handOff","to":"mac"}}"#
1498 );
1499 assert_eq!(serde_json::from_str::<LinkReport>(&json).unwrap(), report);
1500 }
1501
1502 #[test]
1503 fn an_older_queue_entry_still_reads() {
1504 let e: LinkQueueEntry =
1505 serde_json::from_str(r#"{"trackId":"7","title":"t","artist":"a","current":true}"#)
1506 .unwrap();
1507 assert_eq!(e.id, None);
1508 assert_eq!(e.duration_ms, 0);
1509 }
1510
1511 #[test]
1512 fn strangers_cannot_touch_the_library() {
1513 assert!(LinkCommand::Pause.allowed_nearby());
1514 assert!(LinkCommand::HandOff { to: "x".into() }.allowed_nearby());
1515 assert!(!LinkCommand::Sync { full: true }.allowed_nearby());
1516 assert!(!LinkCommand::Evict { track_ids: vec![] }.allowed_nearby());
1517 }
1518
1519 #[test]
1520 fn the_device_id_is_kept() {
1521 let dir = tempfile::tempdir().unwrap();
1522 let first = device_id(dir.path());
1523 assert_eq!(device_id(dir.path()), first);
1524 }
1525}
1526
1527#[cfg(test)]
1528mod device_key_tests {
1529 use super::*;
1530
1531 #[test]
1535 fn device_keys_come_from_the_server_alone() {
1536 let cmd = LinkCommand::DeviceKeys { keys: vec![] };
1537 assert!(!cmd.allowed_nearby());
1538 assert!(!cmd.allowed_playback());
1539 assert_eq!(cmd.from_the_network(true), None);
1540 assert_eq!(cmd.from_the_network(false), None);
1541 }
1542
1543 #[test]
1544 fn the_server_never_relays_its_own_news() {
1545 for news in [
1546 LinkCommand::DeviceKeys { keys: vec![] },
1547 LinkCommand::Devices { devices: vec![] },
1548 LinkCommand::HistoryChanged,
1549 LinkCommand::DspProfilesChanged,
1550 LinkCommand::Shared {
1551 command: Box::new(LinkCommand::Pause),
1552 },
1553 ] {
1554 assert!(!news.relayable(), "{news:?}");
1555 }
1556 assert!(LinkCommand::Pause.relayable());
1557 }
1558
1559 #[test]
1560 fn a_device_key_is_32_bytes_of_base64() {
1561 use base64::Engine as _;
1562 let b64 = base64::engine::general_purpose::STANDARD;
1563 assert!(valid_device_key(&b64.encode([7u8; 32])));
1564 assert!(!valid_device_key(&b64.encode([7u8; 31])));
1565 assert!(!valid_device_key(&b64.encode([7u8; 33])));
1566 assert!(!valid_device_key("not base64!"));
1567 assert!(!valid_device_key(""));
1568 }
1569}