Skip to main content

koan_core/remote/
link.rs

1//! A koan client's standing connection to the koan server it syncs from.
2//!
3//! The server can then act on the client: play a list of tracks on a phone, or
4//! pause it. Subsonic has no way for a server to reach a client, so this is a
5//! koan extension: a WebSocket at `/rest/koanLink`, authenticated like every
6//! other `/rest` call, carrying [`LinkCommand`]s as JSON text frames from the
7//! server. Track ids are the server's, which a synced client holds as each
8//! track's `remote_id`.
9
10use 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/// What a server asks a linked client to do.
30#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
31#[serde(tag = "type", rename_all = "camelCase")]
32pub enum LinkCommand {
33    /// Replace the queue with these tracks and play from `start_at`,
34    /// `position_ms` into it — or load it there paused, when `paused`.
35    #[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        /// Sent by a device handing its music over. One that receives it
45        /// while controlling another device takes control back: the music
46        /// is here now, so this device is what the person is listening to.
47        /// A plain play leaves control alone, so two devices can still
48        /// control each other on purpose.
49        #[serde(default, skip_serializing_if = "std::ops::Not::not")]
50        handoff: bool,
51    },
52    /// Append these tracks to the queue.
53    #[serde(rename_all = "camelCase")]
54    Enqueue {
55        track_ids: Vec<String>,
56    },
57    /// Insert these tracks after the current one.
58    #[serde(rename_all = "camelCase")]
59    PlayNext {
60        track_ids: Vec<String>,
61    },
62    /// Take every queue entry for these tracks out of the queue.
63    #[serde(rename_all = "camelCase")]
64    Remove {
65        track_ids: Vec<String>,
66    },
67    Clear,
68    /// Pull what the server has changed: library, favourites and playlists.
69    /// The library is walked only if the server says it moved, unless `full`,
70    /// which walks it regardless.
71    Sync {
72        #[serde(default)]
73        full: bool,
74    },
75    /// Delete the downloaded copies of these tracks, so the next play fetches
76    /// them again: for a copy that was cached while the server's was bad.
77    #[serde(rename_all = "camelCase")]
78    Evict {
79        track_ids: Vec<String>,
80    },
81    /// Play this track: from where it sits in the queue, or slotted in after
82    /// the current one when the queue does not hold it.
83    #[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    /// Play this queue entry, by the id the device reported it under. Unlike
96    /// `JumpTo` it names one entry of a track queued twice, and reaches a file
97    /// only that device has.
98    PlayItem {
99        id: String,
100    },
101    RemoveItems {
102        ids: Vec<String>,
103    },
104    /// Move these queue entries before or after `target`, in the order given.
105    MoveItems {
106        ids: Vec<String>,
107        target: String,
108        after: bool,
109    },
110    /// Insert these tracks after the queue entry `after`.
111    #[serde(rename_all = "camelCase")]
112    Insert {
113        track_ids: Vec<String>,
114        after: String,
115    },
116    Undo,
117    Redo,
118    /// Turn shuffle on or off: the rest of the queue reordered at random, or
119    /// put back as it was.
120    Shuffle {
121        on: bool,
122    },
123    /// What follows a track at its end: `off`, `queue` or `one`.
124    Repeat {
125        mode: crate::player::state::Repeat,
126    },
127    /// Set the sleep timer, or with none cancel it.
128    SleepTimer {
129        timer: Option<crate::player::state::SleepTimer>,
130    },
131    /// Send this device's queue and playhead to the device `to`, as a `play`,
132    /// and pause here. The device holding the queue does it, so taking music
133    /// from another device and sending it there are one command.
134    HandOff {
135        to: String,
136    },
137    /// The account's other devices, as they are now. News rather than a
138    /// command: sent whenever one of them changes, to links that asked for it.
139    Devices {
140        devices: Vec<LinkDevice>,
141    },
142    /// The public keys the account's devices, and devices shared with it,
143    /// prove themselves with on the local network (`koanDeviceKeys`). News,
144    /// as `Devices` is: sent when a device that registered a key links, and
145    /// whenever the account's keys change. From the server alone, never from
146    /// the network: it is what the network's claims are checked against.
147    DeviceKeys {
148        keys: Vec<LinkDeviceKey>,
149    },
150    /// The accounts this device lets control it. News, as `Devices` is: sent
151    /// when it links and whenever the list changes.
152    Shares {
153        grantees: Vec<String>,
154        /// Why the last request to share or stop sharing was refused.
155        #[serde(default, skip_serializing_if = "Option::is_none")]
156        error: Option<String>,
157        /// Every account on the server, to choose from.
158        #[serde(default, skip_serializing_if = "Vec::is_empty")]
159        accounts: Vec<String>,
160    },
161    /// `command`, sent by another account this device is shared with: run as
162    /// that account's request, from the playback set (`allowed_playback`).
163    Shared {
164        command: Box<LinkCommand>,
165    },
166    /// `device` was forgotten: drop it, however it was last heard of. News,
167    /// as `Devices` is.
168    Forgotten {
169        device: String,
170    },
171    /// The account's play history on the server moved: a play recorded or
172    /// forgotten on one of its devices. The device reads what changed
173    /// (`remote::history::sync`).
174    HistoryChanged,
175    /// The account's EQ profiles kept everywhere moved on the server: one
176    /// changed or deleted on another of its devices. The device reads what
177    /// changed (`remote::dsp_sync::sync`).
178    DspProfilesChanged,
179    /// Play through this output from now on, carrying on from where the music
180    /// is, as the device's own output menu would.
181    SetOutput {
182        output: OutputChoice,
183    },
184    /// List the outputs again, and publish them if they moved: a controller
185    /// has opened its output menu. Cheap, and answered by the link state.
186    RefreshOutputs,
187    /// The volume of the renderer the device plays to, 0–100.
188    SetRendererVolume {
189        volume: u8,
190    },
191    /// Play the output `device` through the DSP profile `profile`, or
192    /// untouched with `None`. A renderer is named by its UDN.
193    SetPreset {
194        device: String,
195        profile: Option<String>,
196    },
197    /// Send this device's audio levels (`LinkReport::Levels`) while `on`: a
198    /// controller has playing bars on screen. A device on the same network may
199    /// ask; levels say no more than the now-playing it already sees.
200    WatchLevels {
201        on: bool,
202    },
203    /// A frame of the levels of `from`, a device this one watches, relayed by
204    /// the server. News, like `Devices`.
205    Levels {
206        from: String,
207        f: crate::remote::levels::Frame,
208    },
209    /// The device `from` answered the command this device sent it under
210    /// `ack`: see `remote::acks`. Relayed by the server; news, like `Levels`.
211    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    /// What a device that is not this account's may have it do: play, pause,
224    /// skip and seek, change the queue, jump, set the volume and the sleep
225    /// timer, choose the output and the preset, and move the music here or
226    /// away. What it asks for runs as the asker's request, never with this
227    /// account's powers: nothing here changes the library, the config beyond
228    /// the output in use, or the account's favourites, playlists or history,
229    /// and a track it names that the library lacks is not synced for (see
230    /// `CommandSource`).
231    /// For a device shared with another account, and one on the local network
232    /// under Full control.
233    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            // The library, and the server's own news.
262            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    /// How a command from a device on the local network runs here, if at
277    /// all: with `full` control (this device's setting) the playback set, as
278    /// `Nearby`; without, a stranger's narrower set, as `Stranger`.
279    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    /// Whether a client may have the server relay this to another device:
288    /// the commands one device gives another. The server's own news (the
289    /// device list, the device keys, shares, forgettings, history) it alone
290    /// originates; relayed, a forged copy would read as the server's. Levels
291    /// go only between live links, by their own route. Exhaustive, so a new
292    /// variant is relayable only once someone decides it is.
293    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    /// Only worth saying to a device that is there to hear it: sent over a
336    /// live route or not at all, never queued for an absent device and never
337    /// a push to wake one.
338    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    /// Whether a device on the same network, which may belong to anyone, may
346    /// send this under Playback only. Playback and the queue; nothing that touches the library or
347    /// the files on disk.
348    ///
349    /// Where the sound goes is not playback: an output switch reaches into
350    /// the room, and a preset into the config. Those are the account's own.
351    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    /// The server's ids for the tracks this command names, if any.
373    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/// Where a command came from, which decides what it may cost.
389#[derive(Debug, Clone, Copy, PartialEq, Eq)]
390pub enum CommandSource {
391    /// The signed-in server, or this person's own devices through it. The
392    /// only source that may cost a sync.
393    Account,
394    /// Another account this device is shared with, through the server: the
395    /// playback set (`allowed_playback`), as that account.
396    Shared,
397    /// A device on the local network, with this device set to Full control:
398    /// the playback set, whoever is signed in there.
399    Nearby,
400    /// A device on the local network, with this device set to Playback only:
401    /// a stranger's set (`allowed_nearby`), and a hand-off that stays on the
402    /// network.
403    Stranger,
404}
405
406/// A device's public key, as `LinkCommand::DeviceKeys` carries it.
407#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
408pub struct LinkDeviceKey {
409    /// The device's id.
410    pub id: String,
411    /// Base64 of its 32-byte Ed25519 public key.
412    pub key: String,
413    /// The account it belongs to, for a device shared with this one's: `None`
414    /// for the account's own.
415    #[serde(default, skip_serializing_if = "Option::is_none")]
416    pub owner: Option<String>,
417}
418
419/// Whether `key` is a public key a device may register: base64 of 32 bytes,
420/// the size of an Ed25519 public key.
421pub 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/// Another device on the same account, as the server sends it.
429#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
430#[serde(rename_all = "camelCase")]
431pub struct LinkDevice {
432    /// The device's own id: stable across its reconnects.
433    pub id: String,
434    pub name: String,
435    pub platform: String,
436    /// Linked now. A device that is not is one iOS has suspended: a command
437    /// wakes it, and music reaches it as a notification to tap.
438    pub linked: bool,
439    /// What it last reported, with the playhead placed as of sending. `None`
440    /// until it has reported at all.
441    pub state: Option<LinkState>,
442    /// Unix seconds when it last held a link, for one that does not now.
443    #[serde(default, skip_serializing_if = "Option::is_none")]
444    pub last_seen: Option<i64>,
445    /// Whether the server can wake it once it is not linked: it gave a push
446    /// token, and the server has a push key. `None` from a server older than
447    /// this, which listed an absent device only when it had a token.
448    #[serde(default, skip_serializing_if = "Option::is_none")]
449    pub wakeable: Option<bool>,
450    /// The account it belongs to, for a device another account shares: `None`
451    /// for the account's own.
452    #[serde(default, skip_serializing_if = "Option::is_none")]
453    pub owner: Option<String>,
454    /// It answers commands sent with an id, which the server relays: see
455    /// `remote::acks`. False from a server or a device that predates that.
456    #[serde(default, skip_serializing_if = "std::ops::Not::not")]
457    pub acks: bool,
458}
459
460impl LinkCommand {
461    /// Every track id the command carries, for translating between the
462    /// server's row ids and the uids it publishes.
463    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/// What a linked client tells the server about itself, as it changes.
507#[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    /// Into the current track, as of when this was sent.
516    #[serde(default)]
517    pub position_ms: u64,
518    #[serde(default)]
519    pub duration_ms: u64,
520    /// The queue, or the part of it around the current track when it is long.
521    #[serde(default)]
522    pub queue: Vec<LinkQueueEntry>,
523    /// What the device can play through, for the device controlling it.
524    #[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    /// The sleep timer is fading playback out.
533    #[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    /// The queue entry's own id on that device, for `playItem` and the rest.
541    #[serde(default)]
542    pub id: Option<String>,
543    /// The server's id for the track; `None` for a file only this device has.
544    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    /// Whether `self` says something `sent`, reported `elapsed` ago, did not.
556    /// A playhead moving at one second per second is not news; a seek, a
557    /// pause, another track or an edited queue is.
558    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/// A message from a client, up the same socket.
576#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
577#[serde(tag = "type", rename_all = "camelCase")]
578pub enum LinkReport {
579    State(LinkState),
580    /// Where Apple's push service reaches this device, so the server can wake
581    /// it once iOS has suspended it and the socket is gone. `sandbox` for a
582    /// development build, whose tokens only the sandbox gateway accepts.
583    Push {
584        token: String,
585        sandbox: bool,
586    },
587    /// Send `command` to the device `to` on the same account. `ack` asks the
588    /// server to relay `to`'s answer: see `remote::acks`. A server that
589    /// predates it ignores it.
590    Command {
591        to: String,
592        command: LinkCommand,
593        #[serde(default, skip_serializing_if = "Option::is_none")]
594        ack: Option<u64>,
595    },
596    /// A command sent under `ack` has come off this device's link, before
597    /// it is acted on: what tells the server the link is alive, however long
598    /// the command then takes.
599    Received {
600        ack: u64,
601    },
602    /// This device's answer to a command sent to it under `ack`.
603    Ack {
604        ack: u64,
605        outcome: crate::remote::acks::AckOutcome,
606    },
607    /// Where to push updates to a Live Activity showing the device `device`;
608    /// `None` for both when the activity has ended.
609    Activity {
610        token: Option<String>,
611        device: Option<String>,
612        #[serde(default)]
613        sandbox: bool,
614    },
615    /// Who is at the other end of a connection made on the local network.
616    Hello(LinkHello),
617    /// Wake the device `to`, which is not linked: with a background push, or
618    /// with `notify` a notification to tap, for when the push has not done it.
619    Wake {
620        to: String,
621        #[serde(default)]
622        notify: bool,
623    },
624    /// Let the account `grantee` control this device, or with `allow` false
625    /// stop letting it. Only ever about the device sending it.
626    Share {
627        grantee: String,
628        allow: bool,
629    },
630    /// Forget `device`, one of this account's that is not linked: its record
631    /// and its push token. It is listed again if it links again.
632    Forget {
633        device: String,
634    },
635    /// A frame of this device's audio levels, while it is watched: see
636    /// `LinkCommand::WatchLevels`. Sent at the analyser's rate, so short.
637    Levels {
638        f: crate::remote::levels::Frame,
639    },
640}
641
642/// How a device introduces itself to one that connected to it over the local
643/// network, before anything else.
644#[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    /// Which library it plays from, as `library_fingerprint` gives it; `None`
651    /// when it is signed in to none. Two devices with the same one share track
652    /// ids, so music can be handed between them.
653    pub library: Option<String>,
654    /// It answers commands sent with an id: see `remote::acks`. False from a
655    /// device that predates that.
656    #[serde(default)]
657    pub acks: bool,
658    /// From a listener that proves itself and asks the same of whoever
659    /// dialled it: the nonce to sign over. See `remote::proof`.
660    #[serde(default, skip_serializing_if = "Option::is_none")]
661    pub nonce: Option<String>,
662}
663
664/// The server this client plays from, as something two devices can compare
665/// without either saying its address to the network.
666pub 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
672/// This device's push token, once the OS has issued one. Set by the app; sent
673/// up each link as it opens, and again if it changes.
674static PUSH_TOKEN: Mutex<Option<(String, bool)>> = Mutex::new(None);
675
676/// A Live Activity on this device showing another: its push token, the device
677/// it shows, and whether the token is the sandbox's. Sent up each link as it
678/// opens, and again when it changes; `None` once the activity has ended.
679type ActivityToken = (String, String, bool);
680static ACTIVITY: Mutex<Option<Option<ActivityToken>>> = Mutex::new(None);
681
682/// Record where the server should push updates to this device's Live
683/// Activity, or that there is none now.
684pub 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
691/// A command as a push notification carries it: the same JSON as over the
692/// link.
693pub fn parse_command(json: &str) -> Result<LinkCommand, String> {
694    serde_json::from_str(json).map_err(|e| e.to_string())
695}
696
697/// Record the push token the OS issued this app, and link now to send it.
698pub 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/// How a client describes itself when it links.
707#[derive(Debug, Clone)]
708pub struct LinkIdentity {
709    /// Shown to whoever picks a client to play on: "James's iPhone".
710    pub name: String,
711    /// `ios`, `macos` or `linux`.
712    pub platform: String,
713    /// Stable across restarts, so a reconnect replaces its own entry on the
714    /// server rather than listing the device twice.
715    pub device_id: String,
716}
717
718impl LinkIdentity {
719    /// This machine, named `name` or else by its hostname.
720    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/// What this device offers whatever controls it: who it is, what it is doing,
742/// and what to do with a command. The link to the server and the connections
743/// made on the local network all serve the same one.
744#[derive(Clone)]
745pub struct Local {
746    pub identity: LinkIdentity,
747    pub state: Arc<dyn Fn() -> LinkState + Send + Sync>,
748    /// Act on a command. With a `Pending`, the sender waits for the answer:
749    /// finish it with the outcome once the command has been acted on.
750    pub on_command:
751        Arc<dyn Fn(LinkCommand, CommandSource, Option<crate::remote::acks::Pending>) + Send + Sync>,
752}
753
754/// Keep a link open to the configured server for as long as the process runs,
755/// handing each command to `on_command` on the link's own thread, and telling
756/// the server what `state` says whenever it changes: which of a person's
757/// devices is the one playing is how the server picks where to send music.
758///
759/// Reads the config before every attempt, so signing in later links without a
760/// restart. Links only to a server whose profile says it can: Navidrome has no
761/// such endpoint.
762pub 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
772/// Moved at every sign-in and sign-out. A link opened under an earlier one
773/// is the last account's, and closes, to open again as whoever is signed in.
774static SIGN_IN: AtomicU64 = AtomicU64::new(0);
775
776/// The account changed: close the link that is up, which was opened for the
777/// last one, and link again now as the new one. A link carries an account's
778/// news, its device keys among them, and none of it is the next account's.
779pub 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        // Before the config is read: a sign-in after this closes the link
792        // about to open with what it read.
793        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            // Not a server that links; asked again when the sign-in changes.
803            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                // The server is answering: downloads waiting out an outage
818                // against it need not wait for their backoff to find out.
819                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                // Asked again before the next attempt: the server may have
832                // been replaced by one that does not link. Not after a link
833                // that simply dropped, which is every time iOS suspends the
834                // app: re-asking then would put two round trips in front of
835                // every reconnect.
836                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
849/// Try to link again now, rather than when the backoff runs out.
850///
851/// For an app coming back to the foreground: iOS suspends a backgrounded app,
852/// its link dies with it, and the retry it was sleeping towards can be a minute
853/// away. Does nothing to a link that is up.
854pub fn nudge() {
855    *NUDGE.0.lock() = true;
856    NUDGE.1.notify_all();
857}
858
859/// Wait `d`, or less if nudged. True when nudged.
860fn 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
868/// Close the link, if it is up: the app has nothing to keep it open for.
869pub fn hang_up() {
870    if let Some(up) = LINK.lock().as_ref() {
871        up.waker.wake();
872    }
873}
874
875/// The link while it is up: what waits to go up it, and how to wake it.
876struct Up {
877    waker: Arc<Waker>,
878    outbox: Vec<LinkReport>,
879}
880
881static LINK: Mutex<Option<Up>> = Mutex::new(None);
882
883/// Send `report` up the link. False when the link is down.
884pub 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
903/// `/rest/koanLink` with the same credentials every other call carries.
904fn 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    // What this device proves itself with on the local network, registered
914    // against the API key the link signs in with. A server that predates
915    // device keys ignores it.
916    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        // Send this link the account's other devices. A server that predates
922        // them ignores it.
923        ("devices", "1"),
924        // Answer commands sent with an id. A server that predates it ignores
925        // it, and relays none.
926        ("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    /// Set while the server says a device of the account is watching this
975    /// one's levels. It counts the watchers; this link holds one watch.
976    levels: Option<crate::remote::levels::Watch>,
977    /// `SIGN_IN` when the config this link signed in with was read.
978    signed_in: u64,
979    /// The account it signed in as, which the device keys it is sent are
980    /// for (`proof::account_of`).
981    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        // Answered up this link, whichever thread finishes it.
1039        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                // Checked by the server, and again here: this device decides
1063                // what another account may have it do.
1064                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
1093/// This library's tracks for the server's ids, in the order given, and
1094/// whether a sync ran to find them.
1095///
1096/// A koan server names a track by its uid, which this library adopted when it
1097/// synced the track; another server by the id it issued. A server can name a
1098/// track added since the last sync; if any are missing and `may_sync`, an
1099/// sync runs first, and whatever is still missing after it is left
1100/// out. An id a sync already failed to find does not start another for a
1101/// while: a command naming a track deleted on the server would otherwise sync
1102/// every time it arrived.
1103pub 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
1145/// Ids a sync looked for and did not find, and when.
1146static 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
1156/// A sync from the configured server, as the app runs its own: the library,
1157/// then favourites and playlists.
1158pub 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
1174/// A random id kept in the config directory, and on iOS in the Keychain as
1175/// well: deleting an app empties its container but not its Keychain items, so
1176/// a reinstalled app keeps its id rather than appearing as a second device.
1177fn 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    // SAFETY: the buffer is valid for its whole length, and gethostname
1214    // writes at most that many bytes.
1215    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    /// Another account, or a device on the network under Full control, gets
1303    /// the playback set: more than a stranger (outputs, presets, volume,
1304    /// hand-off), and nothing of the library or the server's news.
1305    #[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        // Playback only: a stranger's set, run as a stranger.
1342        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    // Neither case may reach `sync`: a test has no business reading the
1353    // machine's config and syncing against the server it names.
1354    #[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    /// The keys are what a network peer's claims are checked against, so no
1525    /// peer may send them: not a stranger, not under Full control, not a
1526    /// shared account.
1527    #[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}