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