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