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