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::{Arc, LazyLock};
15use std::time::{Duration, Instant};
16
17use parking_lot::{Condvar, Mutex};
18use serde::{Deserialize, Serialize};
19use tungstenite::stream::MaybeTlsStream;
20
21use crate::config::{self, Config};
22use crate::helpers::{subsonic_auth, subsonic_client};
23use crate::remote::client::SubsonicAuth;
24pub use crate::remote::outputs::{LinkOutput, LinkOutputs, OutputChoice};
25use crate::remote::profile;
26use crate::remote::wire::{self, Waker};
27
28/// What a server asks a linked client to do.
29#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
30#[serde(tag = "type", rename_all = "camelCase")]
31pub enum LinkCommand {
32    /// Replace the queue with these tracks and play from `start_at`,
33    /// `position_ms` into it — or load it there paused, when `paused`.
34    #[serde(rename_all = "camelCase")]
35    Play {
36        track_ids: Vec<String>,
37        #[serde(default)]
38        start_at: u32,
39        #[serde(default, skip_serializing_if = "is_zero")]
40        position_ms: u64,
41        #[serde(default, skip_serializing_if = "std::ops::Not::not")]
42        paused: bool,
43        /// Sent by a device handing its music over. One that receives it
44        /// while controlling another device takes control back: the music
45        /// is here now, so this device is what the person is listening to.
46        /// A plain play leaves control alone, so two devices can still
47        /// control each other on purpose.
48        #[serde(default, skip_serializing_if = "std::ops::Not::not")]
49        handoff: bool,
50    },
51    /// Append these tracks to the queue.
52    #[serde(rename_all = "camelCase")]
53    Enqueue {
54        track_ids: Vec<String>,
55    },
56    /// Insert these tracks after the current one.
57    #[serde(rename_all = "camelCase")]
58    PlayNext {
59        track_ids: Vec<String>,
60    },
61    /// Take every queue entry for these tracks out of the queue.
62    #[serde(rename_all = "camelCase")]
63    Remove {
64        track_ids: Vec<String>,
65    },
66    Clear,
67    /// Pull what the server has changed: library, favourites and playlists.
68    /// The library is walked only if the server says it moved, unless `full`,
69    /// which walks it regardless.
70    Sync {
71        #[serde(default)]
72        full: bool,
73    },
74    /// Delete the downloaded copies of these tracks, so the next play fetches
75    /// them again: for a copy that was cached while the server's was bad.
76    #[serde(rename_all = "camelCase")]
77    Evict {
78        track_ids: Vec<String>,
79    },
80    /// Play this track: from where it sits in the queue, or slotted in after
81    /// the current one when the queue does not hold it.
82    #[serde(rename_all = "camelCase")]
83    JumpTo {
84        track_id: String,
85    },
86    #[serde(rename_all = "camelCase")]
87    Seek {
88        position_ms: u64,
89    },
90    Pause,
91    Resume,
92    Next,
93    Previous,
94    /// Play this queue entry, by the id the device reported it under. Unlike
95    /// `JumpTo` it names one entry of a track queued twice, and reaches a file
96    /// only that device has.
97    PlayItem {
98        id: String,
99    },
100    RemoveItems {
101        ids: Vec<String>,
102    },
103    /// Move these queue entries before or after `target`, in the order given.
104    MoveItems {
105        ids: Vec<String>,
106        target: String,
107        after: bool,
108    },
109    /// Insert these tracks after the queue entry `after`.
110    #[serde(rename_all = "camelCase")]
111    Insert {
112        track_ids: Vec<String>,
113        after: String,
114    },
115    Undo,
116    Redo,
117    /// Turn shuffle on or off: the rest of the queue reordered at random, or
118    /// put back as it was.
119    Shuffle {
120        on: bool,
121    },
122    /// What follows a track at its end: `off`, `queue` or `one`.
123    Repeat {
124        mode: crate::player::state::Repeat,
125    },
126    /// Send this device's queue and playhead to the device `to`, as a `play`,
127    /// and pause here. The device holding the queue does it, so taking music
128    /// from another device and sending it there are one command.
129    HandOff {
130        to: String,
131    },
132    /// The account's other devices, as they are now. News rather than a
133    /// command: sent whenever one of them changes, to links that asked for it.
134    Devices {
135        devices: Vec<LinkDevice>,
136    },
137    /// Play through this output from now on, carrying on from where the music
138    /// is, as the device's own output menu would.
139    SetOutput {
140        output: OutputChoice,
141    },
142    /// The volume of the renderer the device plays to, 0–100.
143    SetRendererVolume {
144        volume: u8,
145    },
146    /// Play the output `device` through the DSP profile `profile`, or
147    /// untouched with `None`. A renderer is named by its UDN.
148    SetPreset {
149        device: String,
150        profile: Option<String>,
151    },
152}
153
154fn is_zero(n: &u64) -> bool {
155    *n == 0
156}
157
158impl LinkCommand {
159    /// Whether a device on the same network, which may belong to anyone, may
160    /// send this. Playback and the queue; nothing that touches the library or
161    /// the files on disk.
162    ///
163    /// Where the sound goes is not playback: an output switch reaches into
164    /// the room, and a preset into the config. Those are the account's own.
165    pub fn allowed_nearby(&self) -> bool {
166        !matches!(
167            self,
168            Self::Sync { .. }
169                | Self::Evict { .. }
170                | Self::Devices { .. }
171                | Self::SetOutput { .. }
172                | Self::SetRendererVolume { .. }
173                | Self::SetPreset { .. }
174        )
175    }
176
177    /// The server's ids for the tracks this command names, if any.
178    pub fn track_ids(&self) -> &[String] {
179        match self {
180            Self::Play { track_ids, .. }
181            | Self::Enqueue { track_ids }
182            | Self::PlayNext { track_ids }
183            | Self::Remove { track_ids }
184            | Self::Evict { track_ids }
185            | Self::Insert { track_ids, .. } => track_ids,
186            Self::JumpTo { track_id } => std::slice::from_ref(track_id),
187            _ => &[],
188        }
189    }
190}
191
192/// Where a command came from, which decides what it may cost.
193#[derive(Debug, Clone, Copy, PartialEq, Eq)]
194pub enum CommandSource {
195    /// The signed-in server, or this person's own devices through it.
196    Account,
197    /// A device on the same network, which may belong to anyone.
198    Nearby,
199}
200
201/// Another device on the same account, as the server sends it.
202#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
203#[serde(rename_all = "camelCase")]
204pub struct LinkDevice {
205    /// The device's own id: stable across its reconnects.
206    pub id: String,
207    pub name: String,
208    pub platform: String,
209    /// Linked now. A device that is not is one iOS has suspended: a command
210    /// wakes it, and music reaches it as a notification to tap.
211    pub linked: bool,
212    /// What it last reported, with the playhead placed as of sending. `None`
213    /// until it has reported at all.
214    pub state: Option<LinkState>,
215}
216
217impl LinkCommand {
218    /// Every track id the command carries, for translating between the
219    /// server's row ids and the uids it publishes.
220    pub fn track_ids_mut(&mut self) -> Vec<&mut String> {
221        match self {
222            Self::Play { track_ids, .. }
223            | Self::Enqueue { track_ids }
224            | Self::PlayNext { track_ids }
225            | Self::Remove { track_ids }
226            | Self::Evict { track_ids }
227            | Self::Insert { track_ids, .. } => track_ids.iter_mut().collect(),
228            Self::JumpTo { track_id } => vec![track_id],
229            Self::PlayItem { .. }
230            | Self::RemoveItems { .. }
231            | Self::MoveItems { .. }
232            | Self::Undo
233            | Self::Redo
234            | Self::Shuffle { .. }
235            | Self::Repeat { .. }
236            | Self::HandOff { .. }
237            | Self::Devices { .. }
238            | Self::SetOutput { .. }
239            | Self::SetRendererVolume { .. }
240            | Self::SetPreset { .. }
241            | Self::Clear
242            | Self::Sync { .. }
243            | Self::Seek { .. }
244            | Self::Pause
245            | Self::Resume
246            | Self::Next
247            | Self::Previous => Vec::new(),
248        }
249    }
250}
251
252/// What a linked client tells the server about itself, as it changes.
253#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
254#[serde(rename_all = "camelCase")]
255pub struct LinkState {
256    pub playing: bool,
257    pub title: Option<String>,
258    pub artist: Option<String>,
259    #[serde(default)]
260    pub album: Option<String>,
261    /// Into the current track, as of when this was sent.
262    #[serde(default)]
263    pub position_ms: u64,
264    #[serde(default)]
265    pub duration_ms: u64,
266    /// The queue, or the part of it around the current track when it is long.
267    #[serde(default)]
268    pub queue: Vec<LinkQueueEntry>,
269    /// What the device can play through, for the device controlling it.
270    #[serde(default, skip_serializing_if = "Option::is_none")]
271    pub outputs: Option<LinkOutputs>,
272    #[serde(default, skip_serializing_if = "std::ops::Not::not")]
273    pub shuffle: bool,
274    #[serde(default, skip_serializing_if = "crate::player::state::Repeat::is_off")]
275    pub repeat: crate::player::state::Repeat,
276}
277
278#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
279#[serde(rename_all = "camelCase")]
280pub struct LinkQueueEntry {
281    /// The queue entry's own id on that device, for `playItem` and the rest.
282    #[serde(default)]
283    pub id: Option<String>,
284    /// The server's id for the track; `None` for a file only this device has.
285    pub track_id: Option<String>,
286    pub title: String,
287    pub artist: String,
288    #[serde(default)]
289    pub album: String,
290    #[serde(default)]
291    pub duration_ms: u64,
292    pub current: bool,
293}
294
295impl LinkState {
296    /// Whether `self` says something `sent`, reported `elapsed` ago, did not.
297    /// A playhead moving at one second per second is not news; a seek, a
298    /// pause, another track or an edited queue is.
299    pub fn differs(&self, sent: &LinkState, elapsed: Duration) -> bool {
300        let strip = |s: &LinkState| LinkState {
301            position_ms: 0,
302            ..s.clone()
303        };
304        if strip(self) != strip(sent) {
305            return true;
306        }
307        let expected = if sent.playing {
308            sent.position_ms + elapsed.as_millis() as u64
309        } else {
310            sent.position_ms
311        };
312        self.position_ms.abs_diff(expected) > 3000
313    }
314}
315
316/// A message from a client, up the same socket.
317#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
318#[serde(tag = "type", rename_all = "camelCase")]
319pub enum LinkReport {
320    State(LinkState),
321    /// Where Apple's push service reaches this device, so the server can wake
322    /// it once iOS has suspended it and the socket is gone. `sandbox` for a
323    /// development build, whose tokens only the sandbox gateway accepts.
324    Push {
325        token: String,
326        sandbox: bool,
327    },
328    /// Send `command` to the device `to` on the same account.
329    Command {
330        to: String,
331        command: LinkCommand,
332    },
333    /// Where to push updates to a Live Activity showing the device `device`;
334    /// `None` for both when the activity has ended.
335    Activity {
336        token: Option<String>,
337        device: Option<String>,
338        #[serde(default)]
339        sandbox: bool,
340    },
341    /// Who is at the other end of a connection made on the local network.
342    Hello(LinkHello),
343}
344
345/// How a device introduces itself to one that connected to it over the local
346/// network, before anything else.
347#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
348#[serde(rename_all = "camelCase")]
349pub struct LinkHello {
350    pub id: String,
351    pub name: String,
352    pub platform: String,
353    /// Which library it plays from, as `library_fingerprint` gives it; `None`
354    /// when it is signed in to none. Two devices with the same one share track
355    /// ids, so music can be handed between them.
356    pub library: Option<String>,
357}
358
359/// The server this client plays from, as something two devices can compare
360/// without either saying its address to the network.
361pub fn library_fingerprint(cfg: &Config) -> Option<String> {
362    let auth = subsonic_auth(cfg)?;
363    let url = auth.base_url.trim_end_matches('/').to_ascii_lowercase();
364    Some(format!("{:x}", md5::compute(url.as_bytes())))
365}
366
367/// This device's push token, once the OS has issued one. Set by the app; sent
368/// up each link as it opens, and again if it changes.
369static PUSH_TOKEN: Mutex<Option<(String, bool)>> = Mutex::new(None);
370
371/// A Live Activity on this device showing another: its push token, the device
372/// it shows, and whether the token is the sandbox's. Sent up each link as it
373/// opens, and again when it changes; `None` once the activity has ended.
374type ActivityToken = (String, String, bool);
375static ACTIVITY: Mutex<Option<Option<ActivityToken>>> = Mutex::new(None);
376
377/// Record where the server should push updates to this device's Live
378/// Activity, or that there is none now.
379pub fn set_activity(activity: Option<ActivityToken>) {
380    *ACTIVITY.lock() = Some(activity);
381    if let Some(up) = LINK.lock().as_ref() {
382        up.waker.wake();
383    }
384}
385
386/// A command as a push notification carries it: the same JSON as over the
387/// link.
388pub fn parse_command(json: &str) -> Result<LinkCommand, String> {
389    serde_json::from_str(json).map_err(|e| e.to_string())
390}
391
392/// Record the push token the OS issued this app, and link now to send it.
393pub fn set_push_token(token: String, sandbox: bool) {
394    *PUSH_TOKEN.lock() = Some((token, sandbox));
395    nudge();
396    if let Some(up) = LINK.lock().as_ref() {
397        up.waker.wake();
398    }
399}
400
401/// How a client describes itself when it links.
402#[derive(Debug, Clone)]
403pub struct LinkIdentity {
404    /// Shown to whoever picks a client to play on: "James's iPhone".
405    pub name: String,
406    /// `ios`, `macos` or `linux`.
407    pub platform: String,
408    /// Stable across restarts, so a reconnect replaces its own entry on the
409    /// server rather than listing the device twice.
410    pub device_id: String,
411}
412
413impl LinkIdentity {
414    /// This machine, named `name` or else by its hostname.
415    pub fn this_device(name: Option<String>) -> Self {
416        let (platform, label) = if cfg!(target_os = "ios") {
417            ("ios", "iPhone")
418        } else if cfg!(target_os = "macos") {
419            ("macos", "Mac")
420        } else {
421            ("linux", "Linux")
422        };
423        Self {
424            name: name
425                .filter(|n| !n.trim().is_empty())
426                .or_else(hostname)
427                .unwrap_or_else(|| label.to_string()),
428            platform: platform.to_string(),
429            device_id: device_id(&config::config_dir()),
430        }
431    }
432}
433
434/// What this device offers whatever controls it: who it is, what it is doing,
435/// and what to do with a command. The link to the server and the connections
436/// made on the local network all serve the same one.
437#[derive(Clone)]
438pub struct Local {
439    pub identity: LinkIdentity,
440    pub state: Arc<dyn Fn() -> LinkState + Send + Sync>,
441    pub on_command: Arc<dyn Fn(LinkCommand, CommandSource) + Send + Sync>,
442}
443
444/// Keep a link open to the configured server for as long as the process runs,
445/// handing each command to `on_command` on the link's own thread, and telling
446/// the server what `state` says whenever it changes: which of a person's
447/// devices is the one playing is how the server picks where to send music.
448///
449/// Reads the config before every attempt, so signing in later links without a
450/// restart. Links only to a server whose profile says it can: Navidrome has no
451/// such endpoint.
452pub fn spawn(local: Local) {
453    std::thread::Builder::new()
454        .name("koan-link".into())
455        .spawn(move || run(local))
456        .expect("failed to spawn the link thread");
457}
458
459const RETRY_MIN: Duration = Duration::from_secs(2);
460const RETRY_MAX: Duration = Duration::from_secs(60);
461
462fn run(local: Local) {
463    let mut wait = RETRY_MIN;
464    loop {
465        crate::quiet::wait_until_awake();
466        let cfg = Config::load().unwrap_or_default();
467        let Some(auth) = subsonic_auth(&cfg) else {
468            rest(RETRY_MAX);
469            continue;
470        };
471        match profile::for_auth(&auth) {
472            Some(p) if p.links() => {}
473            // Not a server that links; asked again when the sign-in changes.
474            Some(_) => {
475                rest(RETRY_MAX);
476                continue;
477            }
478            None => {
479                rest(wait);
480                wait = (wait * 2).min(RETRY_MAX);
481                continue;
482            }
483        }
484
485        match connect(&auth, &local.identity) {
486            Ok((mut socket, fd)) => {
487                log::info!("link: connected to {}", auth.base_url);
488                // The server is answering: downloads waiting out an outage
489                // against it need not wait for their backoff to find out.
490                if let Some(client) = crate::helpers::subsonic_client(&cfg) {
491                    client.outage().retry_now();
492                }
493                wait = RETRY_MIN;
494                if let Err(e) = serve(&mut socket, fd, &local) {
495                    log::info!("link: closed: {e}");
496                }
497                *LINK.lock() = None;
498                crate::remote::devices::set_linked(false);
499            }
500            Err(e) => {
501                log::warn!("link: {e}");
502                // Asked again before the next attempt: the server may have
503                // been replaced by one that does not link. Not after a link
504                // that simply dropped, which is every time iOS suspends the
505                // app: re-asking then would put two round trips in front of
506                // every reconnect.
507                profile::forget();
508            }
509        }
510        if rest(wait) {
511            wait = RETRY_MIN;
512            continue;
513        }
514        wait = (wait * 2).min(RETRY_MAX);
515    }
516}
517
518static NUDGE: (Mutex<bool>, Condvar) = (Mutex::new(false), Condvar::new());
519
520/// Try to link again now, rather than when the backoff runs out.
521///
522/// For an app coming back to the foreground: iOS suspends a backgrounded app,
523/// its link dies with it, and the retry it was sleeping towards can be a minute
524/// away. Does nothing to a link that is up.
525pub fn nudge() {
526    *NUDGE.0.lock() = true;
527    NUDGE.1.notify_all();
528}
529
530/// Wait `d`, or less if nudged. True when nudged.
531fn rest(d: Duration) -> bool {
532    let mut nudged = NUDGE.0.lock();
533    if !*nudged {
534        NUDGE.1.wait_for(&mut nudged, d);
535    }
536    std::mem::replace(&mut *nudged, false)
537}
538
539/// Close the link, if it is up: the app has nothing to keep it open for.
540pub fn hang_up() {
541    if let Some(up) = LINK.lock().as_ref() {
542        up.waker.wake();
543    }
544}
545
546/// The link while it is up: what waits to go up it, and how to wake it.
547struct Up {
548    waker: Arc<Waker>,
549    outbox: Vec<LinkReport>,
550}
551
552static LINK: Mutex<Option<Up>> = Mutex::new(None);
553
554/// Send `report` up the link. False when the link is down.
555pub fn report(report: LinkReport) -> bool {
556    let mut link = LINK.lock();
557    let Some(up) = link.as_mut() else {
558        return false;
559    };
560    up.outbox.push(report);
561    up.waker.wake();
562    true
563}
564
565type Socket = tungstenite::WebSocket<MaybeTlsStream<TcpStream>>;
566
567fn connect(auth: &SubsonicAuth, identity: &LinkIdentity) -> Result<(Socket, RawFd), String> {
568    let url = link_url(auth, identity)?;
569    let (socket, _) = tungstenite::connect(url).map_err(|e| e.to_string())?;
570    let fd = wire::prepare(socket.get_ref())?;
571    Ok((socket, fd))
572}
573
574/// `/rest/koanLink` with the same credentials every other call carries.
575fn link_url(auth: &SubsonicAuth, identity: &LinkIdentity) -> Result<String, String> {
576    let base = if let Some(rest) = auth.base_url.strip_prefix("https://") {
577        format!("wss://{rest}")
578    } else if let Some(rest) = auth.base_url.strip_prefix("http://") {
579        format!("ws://{rest}")
580    } else {
581        return Err(format!("not an http(s) server: {}", auth.base_url));
582    };
583    let mut query = auth.query().map_err(|e| e.to_string())?;
584    for (k, v) in [
585        ("client", identity.name.as_str()),
586        ("platform", identity.platform.as_str()),
587        ("device", identity.device_id.as_str()),
588        // Send this link the account's other devices. A server that predates
589        // them ignores it.
590        ("devices", "1"),
591    ] {
592        query.push('&');
593        query.push_str(k);
594        query.push('=');
595        query.push_str(&percent_encode(v));
596    }
597    Ok(format!("{base}/rest/koanLink?{query}"))
598}
599
600fn serve(socket: &mut Socket, fd: RawFd, local: &Local) -> Result<(), String> {
601    let waker = Waker::new().map_err(|e| e.to_string())?;
602    wire::wake_on_engine_change(&waker);
603    *LINK.lock() = Some(Up {
604        waker: waker.clone(),
605        outbox: Vec::new(),
606    });
607    crate::remote::devices::set_linked(true);
608    let mut session = LinkSession {
609        local,
610        sent: None,
611        sent_push: None,
612        sent_activity: None,
613    };
614    wire::drive(socket, fd, &waker, &mut session)
615}
616
617struct LinkSession<'a> {
618    local: &'a Local,
619    sent: Option<(LinkState, Instant)>,
620    sent_push: Option<(String, bool)>,
621    sent_activity: Option<Option<ActivityToken>>,
622}
623
624impl wire::Session for LinkSession<'_> {
625    fn outgoing(&mut self) -> Vec<String> {
626        let mut out = Vec::new();
627        let push = PUSH_TOKEN.lock().clone();
628        if let Some((token, sandbox)) = push.clone()
629            && push != self.sent_push
630        {
631            out.push(LinkReport::Push { token, sandbox });
632            self.sent_push = push;
633        }
634        let activity = ACTIVITY.lock().clone();
635        if activity.is_some() && activity != self.sent_activity {
636            let (token, device, sandbox) = match activity.clone().flatten() {
637                Some((t, d, s)) => (Some(t), Some(d), s),
638                None => (None, None, false),
639            };
640            out.push(LinkReport::Activity {
641                token,
642                device,
643                sandbox,
644            });
645            self.sent_activity = activity;
646        }
647        if let Some(up) = LINK.lock().as_mut() {
648            out.append(&mut up.outbox);
649        }
650        let now = (self.local.state)();
651        if self
652            .sent
653            .as_ref()
654            .is_none_or(|(s, at)| now.differs(s, at.elapsed()))
655        {
656            out.push(LinkReport::State(now.clone()));
657            self.sent = Some((now, Instant::now()));
658        }
659        out.iter()
660            .filter_map(|r| serde_json::to_string(r).ok())
661            .collect()
662    }
663
664    fn incoming(&mut self, text: &str) {
665        match serde_json::from_str::<LinkCommand>(text) {
666            Ok(LinkCommand::Devices { devices }) => {
667                crate::remote::devices::set_account(devices);
668            }
669            Ok(cmd) => (self.local.on_command)(cmd, CommandSource::Account),
670            Err(e) => log::warn!("link: not a command ({e}): {text}"),
671        }
672    }
673
674    fn done(&self) -> bool {
675        !crate::quiet::awake()
676    }
677}
678
679/// This library's tracks for the server's ids, in the order given, and
680/// whether a sync ran to find them.
681///
682/// A koan server names a track by its uid, which this library adopted when it
683/// synced the track; another server by the id it issued. A server can name a
684/// track added since the last sync; if any are missing and `may_sync`, an
685/// sync runs first, and whatever is still missing after it is left
686/// out. An id a sync already failed to find does not start another for a
687/// while: a command naming a track deleted on the server would otherwise sync
688/// every time it arrived.
689pub fn resolve_tracks(
690    db: &crate::db::connection::Database,
691    remote_ids: &[String],
692    may_sync: bool,
693) -> (Vec<i64>, bool) {
694    let lookup = |db: &crate::db::connection::Database| {
695        let mut stmt = db
696            .conn
697            .prepare_cached(
698                "SELECT id FROM tracks WHERE uid = ?1
699                 UNION ALL SELECT id FROM tracks WHERE remote_id = ?1 LIMIT 1",
700            )
701            .ok();
702        remote_ids
703            .iter()
704            .map(|rid| {
705                stmt.as_mut()
706                    .and_then(|s| s.query_row([rid], |r| r.get::<_, i64>(0)).ok())
707            })
708            .collect::<Vec<_>>()
709    };
710    let found = lookup(db);
711    let missing: Vec<&String> = remote_ids
712        .iter()
713        .zip(&found)
714        .filter(|(_, f)| f.is_none())
715        .map(|(id, _)| id)
716        .collect();
717    if missing.is_empty() || !may_sync || missing.iter().all(|id| recently_missed(id)) {
718        return (found.into_iter().flatten().collect(), false);
719    }
720    sync(db, crate::helpers::Walk::IfChanged);
721    let found = lookup(db);
722    let mut missed = MISSED.lock();
723    let now = Instant::now();
724    missed.retain(|_, at| now.duration_since(*at) < MISS_TTL);
725    for (id, _) in remote_ids.iter().zip(&found).filter(|(_, f)| f.is_none()) {
726        missed.insert(id.clone(), now);
727    }
728    (found.into_iter().flatten().collect(), true)
729}
730
731/// Ids a sync looked for and did not find, and when.
732static MISSED: LazyLock<Mutex<HashMap<String, Instant>>> = LazyLock::new(Default::default);
733const MISS_TTL: Duration = Duration::from_secs(300);
734
735fn recently_missed(id: &str) -> bool {
736    MISSED
737        .lock()
738        .get(id)
739        .is_some_and(|at| at.elapsed() < MISS_TTL)
740}
741
742/// A sync from the configured server, as the app runs its own: the library,
743/// then favourites and playlists.
744pub fn sync(db: &crate::db::connection::Database, walk: crate::helpers::Walk) {
745    let cfg = Config::load().unwrap_or_default();
746    if let Some(client) = subsonic_client(&cfg)
747        && let Err(e) = crate::helpers::sync_remote(
748            db,
749            &client,
750            walk,
751            &cfg.remote.url,
752            &cfg.remote.username,
753            &|_| {},
754        )
755    {
756        log::warn!("link: sync failed: {e}");
757    }
758}
759
760/// A random id kept in the config directory, and on iOS in the Keychain as
761/// well: deleting an app empties its container but not its Keychain items, so
762/// a reinstalled app keeps its id rather than appearing as a second device.
763fn device_id(dir: &Path) -> String {
764    #[cfg(target_os = "ios")]
765    {
766        use security_framework::passwords::{get_generic_password, set_generic_password};
767        const SERVICE: &str = "cc.blit.koan.link";
768        if let Some(id) = get_generic_password(SERVICE, "device-id")
769            .ok()
770            .and_then(|b| String::from_utf8(b).ok())
771            .filter(|id| !id.trim().is_empty())
772        {
773            return id;
774        }
775        let id = file_device_id(dir);
776        let _ = set_generic_password(SERVICE, "device-id", id.as_bytes());
777        id
778    }
779    #[cfg(not(target_os = "ios"))]
780    file_device_id(dir)
781}
782
783fn file_device_id(dir: &Path) -> String {
784    let path = dir.join("device-id");
785    if let Ok(id) = std::fs::read_to_string(&path) {
786        let id = id.trim();
787        if !id.is_empty() {
788            return id.to_string();
789        }
790    }
791    let id = uuid::Uuid::now_v7().to_string();
792    let _ = std::fs::create_dir_all(dir);
793    let _ = std::fs::write(&path, &id);
794    id
795}
796
797fn hostname() -> Option<String> {
798    let mut buf = [0u8; 256];
799    // SAFETY: the buffer is valid for its whole length, and gethostname
800    // writes at most that many bytes.
801    let ok = unsafe { libc::gethostname(buf.as_mut_ptr().cast(), buf.len()) } == 0;
802    if !ok {
803        return None;
804    }
805    let end = buf.iter().position(|&b| b == 0).unwrap_or(buf.len());
806    let name = String::from_utf8_lossy(&buf[..end]);
807    let name = name.trim_end_matches(".local").trim();
808    (!name.is_empty() && name != "localhost").then(|| name.to_string())
809}
810
811fn percent_encode(s: &str) -> String {
812    let mut out = String::with_capacity(s.len());
813    for b in s.bytes() {
814        if b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_' | b'.' | b'~') {
815            out.push(b as char);
816        } else {
817            out.push_str(&format!("%{b:02X}"));
818        }
819    }
820    out
821}
822
823#[cfg(test)]
824mod tests {
825    use super::*;
826
827    // Neither case may reach `sync`: a test has no business reading the
828    // machine's config and syncing against the server it names.
829    #[test]
830    fn unknown_ids_sync_only_when_allowed_and_not_recently_missed() {
831        let dir = tempfile::tempdir().unwrap();
832        let db = crate::db::connection::Database::open(&dir.path().join("koan.db")).unwrap();
833
834        let (found, synced) = resolve_tracks(&db, &["from-a-stranger".into()], false);
835        assert!(found.is_empty());
836        assert!(!synced, "a nearby peer's unknown id must not start a sync");
837
838        MISSED
839            .lock()
840            .insert("deleted-on-server".into(), Instant::now());
841        let (_, synced) = resolve_tracks(&db, &["deleted-on-server".into()], true);
842        assert!(
843            !synced,
844            "an id a sync just failed to find does not start another"
845        );
846    }
847
848    #[test]
849    fn commands_name_their_tracks() {
850        let play: LinkCommand =
851            serde_json::from_str(r#"{"type":"play","trackIds":["a","b"]}"#).unwrap();
852        assert_eq!(play.track_ids(), ["a", "b"]);
853        let jump: LinkCommand = serde_json::from_str(r#"{"type":"jumpTo","trackId":"c"}"#).unwrap();
854        assert_eq!(jump.track_ids(), ["c"]);
855        assert!(LinkCommand::Pause.track_ids().is_empty());
856    }
857
858    #[test]
859    fn a_push_token_is_a_tagged_report() {
860        let report = LinkReport::Push {
861            token: "ab12".into(),
862            sandbox: true,
863        };
864        let text = serde_json::to_string(&report).unwrap();
865        assert_eq!(text, r#"{"type":"push","token":"ab12","sandbox":true}"#);
866        assert_eq!(serde_json::from_str::<LinkReport>(&text).unwrap(), report);
867    }
868
869    #[test]
870    fn commands_are_tagged_json() {
871        let play = LinkCommand::Play {
872            track_ids: vec!["12".into(), "34".into()],
873            start_at: 1,
874            position_ms: 0,
875            paused: false,
876            handoff: false,
877        };
878        let json = serde_json::to_string(&play).unwrap();
879        assert_eq!(
880            json,
881            r#"{"type":"play","trackIds":["12","34"],"startAt":1}"#
882        );
883        assert_eq!(serde_json::from_str::<LinkCommand>(&json).unwrap(), play);
884        let held = LinkCommand::Play {
885            track_ids: vec!["12".into()],
886            start_at: 0,
887            position_ms: 61_250,
888            paused: true,
889            handoff: false,
890        };
891        let json = serde_json::to_string(&held).unwrap();
892        assert_eq!(
893            json,
894            r#"{"type":"play","trackIds":["12"],"startAt":0,"positionMs":61250,"paused":true}"#
895        );
896        assert_eq!(serde_json::from_str::<LinkCommand>(&json).unwrap(), held);
897        assert_eq!(
898            serde_json::from_str::<LinkCommand>(r#"{"type":"pause"}"#).unwrap(),
899            LinkCommand::Pause
900        );
901        let report = LinkReport::State(LinkState {
902            playing: true,
903            title: Some("Portions for Foxes".into()),
904            ..Default::default()
905        });
906        let json = serde_json::to_string(&report).unwrap();
907        assert!(json.starts_with(r#"{"type":"state","playing":true,"title":"Portions for Foxes""#));
908        assert_eq!(serde_json::from_str::<LinkReport>(&json).unwrap(), report);
909    }
910
911    #[test]
912    fn a_playhead_moving_on_time_is_not_news() {
913        let sent = LinkState {
914            playing: true,
915            position_ms: 10_000,
916            ..Default::default()
917        };
918        let later = |pos| LinkState {
919            position_ms: pos,
920            ..sent.clone()
921        };
922        let five = Duration::from_secs(5);
923        assert!(!later(15_000).differs(&sent, five));
924        assert!(later(60_000).differs(&sent, five), "a seek");
925        let paused = LinkState {
926            playing: false,
927            ..later(15_000)
928        };
929        assert!(paused.differs(&sent, five));
930    }
931
932    #[test]
933    fn the_url_follows_the_scheme_and_names_the_device() {
934        let identity = LinkIdentity {
935            name: "J's iPhone".into(),
936            platform: "ios".into(),
937            device_id: "abc".into(),
938        };
939        let url = link_url(
940            &SubsonicAuth::new("https://music.example.com", "j", "pw"),
941            &identity,
942        )
943        .unwrap();
944        assert!(url.starts_with("wss://music.example.com/rest/koanLink?"));
945        assert!(url.contains("client=J%27s%20iPhone"));
946        assert!(url.contains("device=abc"));
947        assert!(url.contains("devices=1"));
948        assert!(
949            link_url(&SubsonicAuth::new("http://h:4000", "j", "pw"), &identity)
950                .unwrap()
951                .starts_with("ws://h:4000/")
952        );
953    }
954
955    #[test]
956    fn a_relayed_command_nests_the_command() {
957        let report = LinkReport::Command {
958            to: "phone".into(),
959            command: LinkCommand::HandOff { to: "mac".into() },
960        };
961        let json = serde_json::to_string(&report).unwrap();
962        assert_eq!(
963            json,
964            r#"{"type":"command","to":"phone","command":{"type":"handOff","to":"mac"}}"#
965        );
966        assert_eq!(serde_json::from_str::<LinkReport>(&json).unwrap(), report);
967    }
968
969    #[test]
970    fn an_older_queue_entry_still_reads() {
971        let e: LinkQueueEntry =
972            serde_json::from_str(r#"{"trackId":"7","title":"t","artist":"a","current":true}"#)
973                .unwrap();
974        assert_eq!(e.id, None);
975        assert_eq!(e.duration_ms, 0);
976    }
977
978    #[test]
979    fn strangers_cannot_touch_the_library() {
980        assert!(LinkCommand::Pause.allowed_nearby());
981        assert!(LinkCommand::HandOff { to: "x".into() }.allowed_nearby());
982        assert!(!LinkCommand::Sync { full: true }.allowed_nearby());
983        assert!(!LinkCommand::Evict { track_ids: vec![] }.allowed_nearby());
984    }
985
986    #[test]
987    fn the_device_id_is_kept() {
988        let dir = tempfile::tempdir().unwrap();
989        let first = device_id(dir.path());
990        assert_eq!(device_id(dir.path()), first);
991    }
992}