Skip to main content

koan_server/
clients.rs

1//! koan clients linked to this server, and the way to reach them.
2//!
3//! A client that syncs from this server holds a WebSocket open at
4//! `/rest/koanLink` (see `koan_core::remote::link`). Each is registered here
5//! under the account it signed in as, so GraphQL (and MCP through it) can hand
6//! one a list of tracks to play: build a playlist on the server, hear it on a
7//! phone.
8
9use std::sync::LazyLock;
10
11use koan_core::db::queries::{self, UidKind};
12use koan_core::remote::link::{LinkCommand, LinkDevice, LinkState};
13use outbox::Absent;
14use parking_lot::Mutex;
15use tokio::sync::mpsc::UnboundedSender;
16
17/// A linked client, as listed.
18#[derive(Debug, Clone)]
19pub struct ClientInfo {
20    pub id: String,
21    /// The device's own id, stable across its reconnects.
22    pub device: String,
23    pub name: String,
24    pub platform: String,
25    pub username: String,
26    /// Unix seconds.
27    pub connected_at: i64,
28    /// What the client last said it was doing.
29    pub state: LinkState,
30    /// Unix seconds; when it was last seen playing, if ever since linking.
31    pub last_played_at: Option<i64>,
32    /// When `state` was reported, in Unix milliseconds.
33    pub state_at: i64,
34    /// Whether the client has reported its state at all. An app older than
35    /// the reports never does, and its `state` then says nothing about it.
36    pub reports: bool,
37    /// Reached by a notification it shows rather than over its link: iOS had
38    /// suspended it. What was sent runs when someone taps the notification.
39    pub notified: bool,
40}
41
42impl ClientInfo {
43    /// Where the playhead is now, from where it was reported to be.
44    pub fn position_ms(&self) -> u64 {
45        let pos = self.state.position_ms;
46        if !self.state.playing {
47            return pos;
48        }
49        let run = (chrono::Utc::now().timestamp_millis() - self.state_at).max(0) as u64;
50        let pos = pos + run;
51        if self.state.duration_ms > 0 {
52            pos.min(self.state.duration_ms)
53        } else {
54            pos
55        }
56    }
57}
58
59struct Entry {
60    info: ClientInfo,
61    /// The client's own id for itself, so a reconnect replaces its entry.
62    device: String,
63    tx: UnboundedSender<LinkCommand>,
64    /// Sent the account's other devices whenever one changes. Asked for by
65    /// the client; one that predates them would log each as a bad command.
66    wants_devices: bool,
67}
68
69/// A Live Activity on a phone showing another device, and where to push its
70/// updates.
71struct Activity {
72    username: String,
73    /// The phone showing it.
74    watcher: String,
75    /// The device it shows.
76    target: String,
77    token: String,
78    sandbox: bool,
79    /// What it was last sent, so only a change goes out: Apple budgets these.
80    sent: Option<crate::push::ActivityState>,
81}
82
83/// "When this album is in the library, queue it on my device": a request made
84/// before the album exists, fulfilled by the scan that finds it.
85#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
86pub struct Order {
87    pub id: String,
88    /// Whose devices it may go to.
89    pub username: Option<String>,
90    /// A client id or name; `None` for whichever `send` would pick then.
91    pub client: Option<String>,
92    pub artist: String,
93    pub album: String,
94    /// Insert after the current track rather than at the end.
95    pub play_next: bool,
96    /// Add to this playlist rather than a device's queue.
97    #[serde(default)]
98    pub playlist: Option<i64>,
99    /// Only these tracks of the album (title substrings, in this order);
100    /// empty for all of it.
101    #[serde(default)]
102    pub titles: Vec<String>,
103    /// Unix seconds.
104    pub created_at: i64,
105}
106
107/// An order nobody's scan has fulfilled in this long is dropped.
108const ORDER_TTL: i64 = 24 * 60 * 60;
109
110#[derive(Default)]
111pub struct Registry {
112    entries: Mutex<Vec<Entry>>,
113    orders: Mutex<Vec<Order>>,
114    activities: Mutex<Vec<Activity>>,
115}
116
117/// One registry per process: the WebSocket route and the GraphQL schema are
118/// built in different places and both need it.
119pub fn registry() -> &'static Registry {
120    static REGISTRY: LazyLock<Registry> = LazyLock::new(|| {
121        let registry = Registry::default();
122        *registry.orders.lock() = outbox::load_orders();
123        registry
124    });
125    &REGISTRY
126}
127
128impl Registry {
129    /// Add a client, replacing any earlier connection from the same device
130    /// and account. Returns its id.
131    pub fn register(
132        &self,
133        username: &str,
134        name: &str,
135        platform: &str,
136        device: &str,
137        tx: UnboundedSender<LinkCommand>,
138        wants_devices: bool,
139    ) -> String {
140        let id = uuid::Uuid::now_v7().to_string();
141        // What waited for this device while it was away goes down the new
142        // link first.
143        for cmd in outbox::take_and_remember(username, device, name, platform) {
144            let _ = tx.send(cmd);
145        }
146        let mut entries = self.entries.lock();
147        entries.retain(|e| !(e.device == device && e.info.username == username));
148        entries.push(Entry {
149            info: ClientInfo {
150                id: id.clone(),
151                device: device.to_string(),
152                name: name.to_string(),
153                platform: platform.to_string(),
154                username: username.to_string(),
155                connected_at: chrono::Utc::now().timestamp(),
156                state: LinkState::default(),
157                last_played_at: None,
158                state_at: chrono::Utc::now().timestamp_millis(),
159                reports: false,
160                notified: false,
161            },
162            device: device.to_string(),
163            tx,
164            wants_devices,
165        });
166        drop(entries);
167        self.announce(username);
168        id
169    }
170
171    /// Record what a client says it is doing.
172    pub fn report(&self, id: &str, state: LinkState) {
173        let mut entries = self.entries.lock();
174        if let Some(e) = entries.iter_mut().find(|e| e.info.id == id) {
175            if state.playing || e.info.state.playing {
176                e.info.last_played_at = Some(chrono::Utc::now().timestamp());
177            }
178            e.info.state = state;
179            e.info.state_at = chrono::Utc::now().timestamp_millis();
180            e.info.reports = true;
181            let (username, device) = (e.info.username.clone(), e.device.clone());
182            drop(entries);
183            self.announce(&username);
184            self.update_activities(&username, &device);
185        }
186    }
187
188    /// Record where Apple's push service reaches this device.
189    pub fn set_push(&self, username: &str, device: &str, token: &str, sandbox: bool) {
190        outbox::save_push(username, device, token, sandbox);
191    }
192
193    /// Close every link `username` has open. Dropping an entry's sender ends
194    /// its session, and the client must then sign in again to reconnect.
195    pub fn disconnect(&self, username: &str) {
196        self.entries.lock().retain(|e| e.info.username != username);
197    }
198
199    pub fn unregister(&self, id: &str) {
200        let mut entries = self.entries.lock();
201        let username = entries
202            .iter()
203            .find(|e| e.info.id == id)
204            .map(|e| e.info.username.clone());
205        entries.retain(|e| e.info.id != id);
206        drop(entries);
207        if let Some(username) = username {
208            self.announce(&username);
209        }
210    }
211
212    /// Send each of `username`'s links that asked for them the account's
213    /// other devices: those linked, with what each is doing, and those a push
214    /// can wake.
215    fn announce(&self, username: &str) {
216        let listening = |e: &Entry| e.info.username == username && e.wants_devices;
217        if !self.entries.lock().iter().any(listening) {
218            return;
219        }
220        let asleep = outbox::push_targets(Some(username));
221        let entries = self.entries.lock();
222        let ours: Vec<&Entry> = entries
223            .iter()
224            .filter(|e| e.info.username == username)
225            .collect();
226        if !ours.iter().any(|e| e.wants_devices) {
227            return;
228        }
229        let mut all: Vec<LinkDevice> = ours
230            .iter()
231            .map(|e| LinkDevice {
232                id: e.device.clone(),
233                name: e.info.name.clone(),
234                platform: e.info.platform.clone(),
235                linked: true,
236                state: e.info.reports.then(|| LinkState {
237                    position_ms: e.info.position_ms(),
238                    ..e.info.state.clone()
239                }),
240            })
241            .collect();
242        for t in asleep {
243            if !all.iter().any(|d| d.id == t.device) {
244                all.push(LinkDevice {
245                    id: t.device,
246                    name: t.name,
247                    platform: t.platform,
248                    linked: false,
249                    state: None,
250                });
251            }
252        }
253        for e in ours.iter().filter(|e| e.wants_devices) {
254            let devices = all.iter().filter(|d| d.id != e.device).cloned().collect();
255            let _ = e.tx.send(LinkCommand::Devices { devices });
256        }
257    }
258
259    /// Relay `command` from one of `username`'s devices to another, `to`.
260    pub fn relay(
261        &self,
262        username: &str,
263        to: &str,
264        command: LinkCommand,
265    ) -> Result<ClientInfo, String> {
266        if matches!(command, LinkCommand::Devices { .. }) {
267            return Err("not a command".into());
268        }
269        self.send(Some(username), Some(to), command)
270    }
271
272    /// Where to push a Live Activity's updates: `watcher` shows `target`.
273    /// `None` ends it.
274    pub fn set_activity(
275        &self,
276        username: &str,
277        watcher: &str,
278        activity: Option<(String, String, bool)>,
279    ) {
280        let mut activities = self.activities.lock();
281        activities.retain(|a| !(a.username == username && a.watcher == watcher));
282        let Some((token, target, sandbox)) = activity else {
283            return;
284        };
285        activities.push(Activity {
286            username: username.to_string(),
287            watcher: watcher.to_string(),
288            target: target.clone(),
289            token,
290            sandbox,
291            sent: None,
292        });
293        drop(activities);
294        self.update_activities(username, &target);
295    }
296
297    /// Push `target`'s state to every Live Activity showing it, where it has
298    /// changed in a way the activity shows.
299    fn update_activities(&self, username: &str, target: &str) {
300        let Some(pusher) = crate::push::pusher() else {
301            return;
302        };
303        let Some(info) = self
304            .list(Some(username))
305            .into_iter()
306            .find(|c| c.device == target)
307        else {
308            return;
309        };
310        let state = crate::push::ActivityState::of(&info);
311        let mut due = Vec::new();
312        for a in self.activities.lock().iter_mut() {
313            if a.username == username
314                && a.target == target
315                && a.sent.as_ref().is_none_or(|s| s.differs(&state))
316            {
317                a.sent = Some(state.clone());
318                due.push((a.token.clone(), a.sandbox, a.watcher.clone()));
319            }
320        }
321        if due.is_empty() {
322            return;
323        }
324        let username = username.to_string();
325        std::thread::spawn(move || {
326            for (token, sandbox, watcher) in due {
327                let push = crate::push::Push::Activity(state.clone());
328                match pusher.send(&token, sandbox, &push) {
329                    crate::push::Outcome::Sent => {}
330                    crate::push::Outcome::Gone => {
331                        log::info!("push: a Live Activity on {watcher} has ended");
332                        registry().set_activity(&username, &watcher, None);
333                    }
334                    crate::push::Outcome::Failed(e) => {
335                        log::warn!("push: Live Activity on {watcher}: {e}");
336                    }
337                }
338            }
339        });
340    }
341
342    /// Clients `username` may command, newest first; every client for `None`.
343    pub fn list(&self, username: Option<&str>) -> Vec<ClientInfo> {
344        let mut out: Vec<ClientInfo> = self
345            .entries
346            .lock()
347            .iter()
348            .filter(|e| username.is_none_or(|u| e.info.username == u))
349            .map(|e| e.info.clone())
350            .collect();
351        out.sort_by_key(|c| std::cmp::Reverse(c.connected_at));
352        out
353    }
354
355    /// Send to `id` (an id or a name), or with none to the client the
356    /// command most likely means: the one playing, else the one that played
357    /// within `RECENT`, else the only one linked. `Err` names the choices when
358    /// there is no telling, so whoever asked can ask the person.
359    pub fn send(
360        &self,
361        username: Option<&str>,
362        id: Option<&str>,
363        cmd: LinkCommand,
364    ) -> Result<ClientInfo, String> {
365        let clients = self.list(username);
366        let target = match id {
367            Some(id) => clients
368                .iter()
369                .find(|c| c.id == id || c.device == id || c.name.eq_ignore_ascii_case(id)),
370            None if clients.is_empty() => None,
371            None => Some(pick(&clients, chrono::Utc::now().timestamp())?),
372        };
373        // Not linked: a phone iOS has suspended is woken to take it.
374        let Some(target) = target else {
375            return reach_absent(username, id, &cmd).unwrap_or_else(|| {
376                Err(match id {
377                    Some(id) => format!("no linked client {id}; see `clients`"),
378                    None => "no koan app is linked to this server; open koan on the device".into(),
379                })
380            });
381        };
382        let entries = self.entries.lock();
383        let entry = entries
384            .iter()
385            .find(|e| e.info.id == target.id)
386            .ok_or("that client has just gone")?;
387        entry
388            .tx
389            .send(cmd)
390            .map_err(|_| "that client has just gone".to_string())?;
391        Ok(target.clone())
392    }
393}
394
395impl Registry {
396    pub fn add_order(&self, order: Order) {
397        outbox::save_order(&order);
398        self.orders.lock().push(order);
399    }
400
401    pub fn orders(&self, username: Option<&str>) -> Vec<Order> {
402        self.orders
403            .lock()
404            .iter()
405            .filter(|o| username.is_none() || o.username.as_deref() == username)
406            .cloned()
407            .collect()
408    }
409
410    pub fn cancel_order(&self, username: Option<&str>, id: &str) -> bool {
411        let mut orders = self.orders.lock();
412        let before = orders.len();
413        orders
414            .retain(|o| !(o.id == id && (username.is_none() || o.username.as_deref() == username)));
415        let gone = orders.len() != before;
416        if gone {
417            outbox::drop_order(id);
418        }
419        gone
420    }
421
422    fn done(&self, id: &str) {
423        self.orders.lock().retain(|o| o.id != id);
424        outbox::drop_order(id);
425    }
426
427    /// Send every order whose album the library now holds, and drop it.
428    /// `find` answers an order with the album's track uids, in order.
429    pub fn fulfil_orders(&self, find: impl Fn(&Order) -> Option<Vec<String>>) {
430        let now = chrono::Utc::now().timestamp();
431        let pending: Vec<Order> = {
432            let mut orders = self.orders.lock();
433            for o in orders.iter().filter(|o| now - o.created_at >= ORDER_TTL) {
434                outbox::drop_order(&o.id);
435            }
436            orders.retain(|o| now - o.created_at < ORDER_TTL);
437            orders.clone()
438        };
439        for order in pending.into_iter().filter(|o| o.playlist.is_none()) {
440            let Some(ids) = find(&order).filter(|ids| !ids.is_empty()) else {
441                continue;
442            };
443            let track_ids = ids;
444            let cmd = if order.play_next {
445                LinkCommand::PlayNext { track_ids }
446            } else {
447                LinkCommand::Enqueue { track_ids }
448            };
449            match self.send(order.username.as_deref(), order.client.as_deref(), cmd) {
450                Ok(c) => {
451                    log::info!(
452                        "link: {} — {} arrived; queued on {}",
453                        order.artist,
454                        order.album,
455                        c.name
456                    );
457                    self.done(&order.id);
458                }
459                // No device to send to yet: kept, and tried after the next scan.
460                Err(e) => log::info!("link: {} — {} arrived but {e}", order.artist, order.album),
461            }
462        }
463    }
464}
465
466/// Fulfil standing orders against the library at `db_path`.
467pub fn fulfil_from(db_path: &std::path::Path) {
468    let registry = registry();
469    if registry.orders.lock().is_empty() {
470        return;
471    }
472    let Ok(db) = koan_core::db::connection::Database::open_existing(db_path) else {
473        return;
474    };
475    // Playlist orders are the server's own to carry out: add the tracks, and
476    // the playlist reaches every device like any other edit.
477    let for_playlists: Vec<Order> = registry
478        .orders
479        .lock()
480        .iter()
481        .filter(|o| o.playlist.is_some())
482        .cloned()
483        .collect();
484    let mut edited = false;
485    for order in for_playlists {
486        let Some(playlist) = order.playlist else {
487            continue;
488        };
489        let Some(ids) = order_tracks(&db.conn, &order) else {
490            continue;
491        };
492        match koan_core::db::queries::add_tracks(&db.conn, playlist, &ids) {
493            Ok(_) => {
494                // Tests would push to whatever server this machine signs in to.
495                if !cfg!(test) {
496                    koan_core::playlists::push_to_remote(playlist);
497                }
498                log::info!(
499                    "link: {} — {} arrived; added {} tracks to playlist {playlist}",
500                    order.artist,
501                    order.album,
502                    ids.len()
503                );
504                registry.done(&order.id);
505                edited = true;
506            }
507            Err(e) => log::warn!("link: could not add to playlist {playlist}: {e}"),
508        }
509    }
510    if edited {
511        changed();
512    }
513    registry.fulfil_orders(|order| {
514        let rows = order_tracks(&db.conn, order)?;
515        queries::uids_in_order(&db.conn, UidKind::Track, &rows).ok()
516    });
517}
518
519/// The tracks an order asks for, once its album is in the library: all of
520/// it, or the named ones in the order named. `None` until they are there.
521fn order_tracks(conn: &rusqlite::Connection, order: &Order) -> Option<Vec<i64>> {
522    let tracks = album_tracks(conn, &order.artist, &order.album)?;
523    if order.titles.is_empty() {
524        return Some(tracks.into_iter().map(|(id, _)| id).collect());
525    }
526    let picked: Vec<i64> = order
527        .titles
528        .iter()
529        .filter_map(|want| {
530            let want = want.to_lowercase();
531            tracks
532                .iter()
533                .find(|(_, t)| t.to_lowercase().contains(&want))
534                .map(|(id, _)| *id)
535        })
536        .collect();
537    (!picked.is_empty()).then_some(picked)
538}
539
540/// The newest album whose artist and title contain these, as its tracks in
541/// disc and track order.
542pub fn album_tracks(
543    conn: &rusqlite::Connection,
544    artist: &str,
545    album: &str,
546) -> Option<Vec<(i64, String)>> {
547    let like = |s: &str| format!("%{}%", s.replace(['%', '_'], ""));
548    let album_id: i64 = conn
549        .query_row(
550            "SELECT al.id FROM albums al JOIN artists a ON a.id = al.artist_id
551              WHERE a.name LIKE ?1 COLLATE NOCASE AND al.title LIKE ?2 COLLATE NOCASE
552              ORDER BY al.id DESC LIMIT 1",
553            [like(artist), like(album)],
554            |r| r.get(0),
555        )
556        .ok()?;
557    let mut stmt = conn
558        .prepare("SELECT id, title FROM tracks WHERE album_id = ?1 ORDER BY disc, track_number, id")
559        .ok()?;
560    let tracks = stmt
561        .query_map([album_id], |r| Ok((r.get(0)?, r.get(1)?)))
562        .ok()?
563        .filter_map(Result::ok)
564        .collect();
565    Some(tracks)
566}
567
568impl Registry {
569    /// Send to every client `username` may command. The names of those it
570    /// reached.
571    pub fn broadcast(&self, username: Option<&str>, cmd: LinkCommand) -> Vec<String> {
572        let ids: Vec<String> = self.list(username).into_iter().map(|c| c.id).collect();
573        let entries = self.entries.lock();
574        entries
575            .iter()
576            .filter(|e| ids.contains(&e.info.id) && e.tx.send(cmd.clone()).is_ok())
577            .map(|e| e.info.name.clone())
578            .collect()
579    }
580}
581
582impl Registry {
583    /// Send to every device `username` may command: at once to those linked,
584    /// and to those that have linked before but are away now, when they next
585    /// link. For commands still right hours later (`Sync`, `Evict`), never
586    /// playback. The names reached now, and the names it waits for.
587    pub fn deliver(&self, username: Option<&str>, cmd: LinkCommand) -> (Vec<String>, Vec<String>) {
588        let (sent, queued) = self.link_or_queue(username, &cmd);
589        wake(&queued.iter().map(Absent::key).collect::<Vec<_>>());
590        (sent, queued.into_iter().map(|q| q.name).collect())
591    }
592
593    fn link_or_queue(
594        &self,
595        username: Option<&str>,
596        cmd: &LinkCommand,
597    ) -> (Vec<String>, Vec<Absent>) {
598        let sent = self.broadcast(username, cmd.clone());
599        let queued = outbox::queue_for_absent(username, &self.live(), cmd);
600        (sent, queued)
601    }
602
603    /// Each linked device, as `(device, username)`.
604    fn live(&self) -> Vec<(String, String)> {
605        self.entries
606            .lock()
607            .iter()
608            .map(|e| (e.device.clone(), e.info.username.clone()))
609            .collect()
610    }
611}
612
613impl Absent {
614    fn key(&self) -> (String, String) {
615        (self.device.clone(), self.username.clone())
616    }
617}
618
619/// Wake these absent devices, each a `(device, username)`, where a push can:
620/// each links, and takes what was queued for it.
621fn wake(devices: &[(String, String)]) {
622    let Some(pusher) = crate::push::pusher() else {
623        return;
624    };
625    let targets: Vec<outbox::PushTarget> = outbox::push_targets(None)
626        .into_iter()
627        .filter(|t| {
628            devices
629                .iter()
630                .any(|(device, username)| *device == t.device && *username == t.username)
631        })
632        .collect();
633    if targets.is_empty() {
634        return;
635    }
636    std::thread::spawn(move || {
637        for t in targets {
638            deliver_push(pusher, &t, &crate::push::Push::Wake);
639        }
640    });
641}
642
643/// Reach a device that is not linked.
644///
645/// Music is sent as a notification to tap, at once: iOS does not let an app it
646/// woke start audio, so waking it for that only delays the notification. Every
647/// other command goes into the device's outbox with a background push to wake
648/// it, and runs as if it had been linked. `None` without a push key, or with no
649/// device in scope that has given a push token.
650fn reach_absent(
651    username: Option<&str>,
652    id: Option<&str>,
653    cmd: &LinkCommand,
654) -> Option<Result<ClientInfo, String>> {
655    let pusher = crate::push::pusher()?;
656    let targets: Vec<outbox::PushTarget> = outbox::push_targets(username)
657        .into_iter()
658        .filter(|t| id.is_none_or(|id| t.device == id || t.name.eq_ignore_ascii_case(id)))
659        .collect();
660    let target = match targets.as_slice() {
661        [] => return None,
662        [only] => only.clone(),
663        several => {
664            return Some(Err(format!(
665                "no koan app is linked, and several can be reached: {}. Ask which, then pass `client`",
666                several
667                    .iter()
668                    .map(|t| t.name.as_str())
669                    .collect::<Vec<_>>()
670                    .join(", ")
671            )));
672        }
673    };
674    let info = ClientInfo {
675        id: target.device.clone(),
676        device: target.device.clone(),
677        name: target.name.clone(),
678        platform: target.platform.clone(),
679        username: target.username.clone(),
680        connected_at: 0,
681        state: LinkState::default(),
682        last_played_at: None,
683        state_at: 0,
684        reports: false,
685        notified: true,
686    };
687    let verb = match cmd {
688        LinkCommand::Play { .. } | LinkCommand::JumpTo { .. } | LinkCommand::Resume => Some("Play"),
689        _ => None,
690    };
691    let push = match verb {
692        Some(verb) => crate::push::Push::Notify {
693            title: format!("{verb} on {}", target.name),
694            body: outbox::describe(cmd).unwrap_or_else(|| "From your koan server".into()),
695            command: serde_json::to_value(cmd).ok()?,
696            image: cover_track(cmd)
697                .and_then(outbox::track_row)
698                .and_then(|t| pusher.cover_link(t)),
699        },
700        None => {
701            outbox::queue_for(&target.device, &target.username, cmd);
702            crate::push::Push::Wake
703        }
704    };
705    std::thread::spawn(move || deliver_push(pusher, &target, &push));
706    Some(Ok(info))
707}
708
709/// The track whose album cover a notification for `cmd` shows. Commands carry
710/// uids, so this is the id as sent; see `outbox::track_row`.
711fn cover_track(cmd: &LinkCommand) -> Option<&str> {
712    match cmd {
713        LinkCommand::Play {
714            track_ids,
715            start_at,
716            ..
717        } => track_ids
718            .get(*start_at as usize)
719            .or(track_ids.first())
720            .map(String::as_str),
721        LinkCommand::JumpTo { track_id } => Some(track_id),
722        _ => None,
723    }
724}
725
726/// Send one push, forgetting a token Apple says is no longer good.
727fn deliver_push(
728    pusher: &crate::push::Pusher,
729    target: &outbox::PushTarget,
730    push: &crate::push::Push,
731) {
732    use crate::push::Outcome;
733    match pusher.send(&target.token, target.sandbox, push) {
734        Outcome::Sent => log::info!("push: sent to {}", target.name),
735        Outcome::Gone => {
736            log::info!(
737                "push: {}'s token is no longer valid; forgotten",
738                target.name
739            );
740            outbox::forget_push(&target.username, &target.device);
741        }
742        Outcome::Failed(e) => log::warn!("push: to {} failed: {e}", target.name),
743    }
744}
745
746/// Have every device pull what the server just changed (a playlist edited,
747/// albums added): at once where linked, on next link where not. Syncs waiting
748/// for a device collapse into one, and so do the pushes that wake it: see
749/// `Wakes`.
750pub fn changed() {
751    let (_, queued) = registry().link_or_queue(None, &LinkCommand::Sync { full: false });
752    if queued.is_empty() || crate::push::pusher().is_none() {
753        return;
754    }
755    let mut wakes = WAKES.lock();
756    wakes.add(queued.iter().map(Absent::key), std::time::Instant::now());
757    if !wakes.timer {
758        wakes.timer = true;
759        std::thread::spawn(send_wakes);
760    }
761}
762
763/// How long the library has to stay still before a suspended device is woken
764/// to sync. A download of several albums scans after each one; iOS rations
765/// background pushes, and the device needs waking once, at the end.
766const QUIET: std::time::Duration = std::time::Duration::from_secs(30);
767
768/// However busy the library stays, a device waits no longer than this.
769const LONGEST_WAIT: std::time::Duration = std::time::Duration::from_secs(5 * 60);
770
771static WAKES: Mutex<Wakes> = Mutex::new(Wakes {
772    pending: Vec::new(),
773    timer: false,
774});
775
776/// Background pushes held back until the library is quiet: at most one
777/// pending per device.
778struct Wakes {
779    /// `(device, username)`, when first asked for, when last asked for.
780    pending: Vec<((String, String), std::time::Instant, std::time::Instant)>,
781    /// Whether a thread is waiting to send them.
782    timer: bool,
783}
784
785impl Wakes {
786    fn add(
787        &mut self,
788        devices: impl IntoIterator<Item = (String, String)>,
789        now: std::time::Instant,
790    ) {
791        for key in devices {
792            match self.pending.iter_mut().find(|(k, _, _)| *k == key) {
793                Some((_, _, last)) => *last = now,
794                None => self.pending.push((key, now, now)),
795            }
796        }
797    }
798
799    fn due_at(first: std::time::Instant, last: std::time::Instant) -> std::time::Instant {
800        (last + QUIET).min(first + LONGEST_WAIT)
801    }
802
803    /// The earliest a pending wake is due.
804    fn next(&self) -> Option<std::time::Instant> {
805        self.pending
806            .iter()
807            .map(|(_, first, last)| Self::due_at(*first, *last))
808            .min()
809    }
810
811    /// Take the wakes due by `now`.
812    fn take_due(&mut self, now: std::time::Instant) -> Vec<(String, String)> {
813        let (due, waiting) = std::mem::take(&mut self.pending)
814            .into_iter()
815            .partition(|(_, first, last)| Self::due_at(*first, *last) <= now);
816        self.pending = waiting;
817        due.into_iter().map(|(key, _, _)| key).collect()
818    }
819}
820
821/// Send each wake once it is due, until none is pending. A device that has
822/// linked meanwhile took its sync down the link and is not pushed.
823fn send_wakes() {
824    loop {
825        let (due, next) = {
826            let mut wakes = WAKES.lock();
827            let due = wakes.take_due(std::time::Instant::now());
828            let next = wakes.next();
829            if due.is_empty() && next.is_none() {
830                wakes.timer = false;
831                return;
832            }
833            (due, next)
834        };
835        let live = registry().live();
836        let absent: Vec<(String, String)> = due.into_iter().filter(|d| !live.contains(d)).collect();
837        if !absent.is_empty() {
838            wake(&absent);
839        }
840        if let Some(next) = next {
841            std::thread::sleep(next.saturating_duration_since(std::time::Instant::now()));
842        }
843    }
844}
845
846/// After a library scan or sync: if the library holds different tracks or
847/// albums from when this was last asked, tell every device. A scan that found
848/// nothing new, which is most of them, sends nothing.
849pub fn changed_if_library_moved(conn: &rusqlite::Connection) {
850    static LAST: parking_lot::Mutex<Option<(i64, i64, i64)>> = parking_lot::Mutex::new(None);
851    let Ok(now) = conn.query_row(
852        "SELECT (SELECT COUNT(*) FROM tracks), (SELECT COALESCE(MAX(id), 0) FROM tracks),
853                (SELECT COUNT(*) FROM albums)",
854        [],
855        |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
856    ) else {
857        return;
858    };
859    let before = LAST.lock().replace(now);
860    if before.is_some_and(|b| b != now) {
861        changed();
862    }
863}
864
865/// The server-side record of devices and their waiting commands, in the
866/// library database so it outlives a restart.
867mod outbox {
868    use koan_core::db::queries;
869    use koan_core::remote::link::LinkCommand;
870
871    /// Dropped undelivered after this long: a device away a month re-syncs
872    /// on its own when opened.
873    const KEEP_SECS: i64 = 30 * 24 * 60 * 60;
874
875    /// Tests keep to memory: the configured database is whoever ran them.
876    ///
877    /// The shared pool, because a state report from every linked device lands
878    /// here: `Database::open` would run the schema and a checkpoint each time.
879    fn db() -> Option<koan_core::db::pool::Handle<'static>> {
880        if cfg!(test) {
881            return None;
882        }
883        koan_core::db::pool::shared().get().ok()
884    }
885
886    pub fn load_orders() -> Vec<super::Order> {
887        let Some(db) = db() else { return Vec::new() };
888        db.conn
889            .prepare("SELECT body FROM link_orders ORDER BY created_at")
890            .and_then(|mut s| {
891                s.query_map([], |r| r.get::<_, String>(0))?
892                    .collect::<Result<Vec<_>, _>>()
893            })
894            .unwrap_or_default()
895            .into_iter()
896            .filter_map(|b| serde_json::from_str(&b).ok())
897            .collect()
898    }
899
900    pub fn save_order(order: &super::Order) {
901        let (Some(db), Ok(body)) = (db(), serde_json::to_string(order)) else {
902            return;
903        };
904        let _ = db.conn.execute(
905            "INSERT OR REPLACE INTO link_orders (id, body, created_at) VALUES (?1, ?2, ?3)",
906            rusqlite::params![order.id, body, order.created_at],
907        );
908    }
909
910    pub fn drop_order(id: &str) {
911        if let Some(db) = db() {
912            let _ = db
913                .conn
914                .execute("DELETE FROM link_orders WHERE id = ?1", [id]);
915        }
916    }
917
918    pub fn take_and_remember(
919        username: &str,
920        device: &str,
921        name: &str,
922        platform: &str,
923    ) -> Vec<LinkCommand> {
924        let Some(db) = db() else { return Vec::new() };
925        let now = chrono::Utc::now().timestamp();
926        let _ = db.conn.execute(
927            "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES (?1, ?2, ?3, ?4, ?5)
928             ON CONFLICT (device, username) DO UPDATE SET name = ?3, platform = ?4, last_seen = ?5",
929            rusqlite::params![device, username, name, platform, now],
930        );
931        let _ = db.conn.execute(
932            "DELETE FROM link_outbox WHERE created_at < ?1",
933            [now - KEEP_SECS],
934        );
935        let waiting: Vec<(i64, String)> = db
936            .conn
937            .prepare("SELECT id, command FROM link_outbox WHERE device = ?1 AND username = ?2 ORDER BY id")
938            .and_then(|mut s| {
939                s.query_map([device, username], |r| Ok((r.get(0)?, r.get(1)?)))?
940                    .collect()
941            })
942            .unwrap_or_default();
943        let _ = db.conn.execute(
944            "DELETE FROM link_outbox WHERE device = ?1 AND username = ?2",
945            [device, username],
946        );
947        if !waiting.is_empty() {
948            log::info!("link: {} waiting commands for {name}", waiting.len());
949        }
950        waiting
951            .into_iter()
952            .filter_map(|(_, c)| serde_json::from_str(&c).ok())
953            .collect()
954    }
955
956    /// A known device that was not linked when something was queued for it.
957    pub struct Absent {
958        pub device: String,
959        pub username: String,
960        pub name: String,
961    }
962
963    /// Queue `cmd` for each known device in scope that is not in `live`.
964    pub fn queue_for_absent(
965        username: Option<&str>,
966        live: &[(String, String)],
967        cmd: &LinkCommand,
968    ) -> Vec<Absent> {
969        let Some(db) = db() else { return Vec::new() };
970        let known: Vec<(String, String, String)> = db
971            .conn
972            .prepare("SELECT device, username, name FROM link_devices")
973            .and_then(|mut s| {
974                s.query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))?
975                    .collect()
976            })
977            .unwrap_or_default();
978        let Ok(text) = serde_json::to_string(cmd) else {
979            return Vec::new();
980        };
981        let absent: Vec<_> = known
982            .into_iter()
983            .filter(|(device, user, _)| {
984                !username.is_some_and(|u| u != user)
985                    && !live.iter().any(|(d, u)| d == device && u == user)
986            })
987            .collect();
988        if absent.is_empty() {
989            return Vec::new();
990        }
991        let is_sync = matches!(cmd, LinkCommand::Sync { .. });
992        let now = chrono::Utc::now().timestamp();
993        // Every playlist edit and scan lands here: one transaction, not two
994        // per device.
995        koan_core::db::queries::atomically(&db.conn, || {
996            let mut queued = Vec::new();
997            for (device, user, name) in absent {
998                if is_sync {
999                    // One pending sync is enough; a full one covers an incremental.
1000                    let _ = db.conn.execute(
1001                        "DELETE FROM link_outbox WHERE device = ?1 AND username = ?2 AND command LIKE '{\"type\":\"sync\"%'",
1002                        [&device, &user],
1003                    );
1004                }
1005                if db
1006                    .conn
1007                    .execute(
1008                        "INSERT INTO link_outbox (device, username, command, created_at) VALUES (?1, ?2, ?3, ?4)",
1009                        rusqlite::params![device, user, text, now],
1010                    )
1011                    .is_ok()
1012                {
1013                    queued.push(Absent {
1014                        device,
1015                        username: user,
1016                        name,
1017                    });
1018                }
1019            }
1020            Ok::<_, rusqlite::Error>(queued)
1021        })
1022        .unwrap_or_default()
1023    }
1024
1025    /// A device Apple's push service can reach.
1026    #[derive(Clone)]
1027    pub struct PushTarget {
1028        pub device: String,
1029        pub username: String,
1030        pub name: String,
1031        pub platform: String,
1032        pub token: String,
1033        pub sandbox: bool,
1034    }
1035
1036    /// Queue `cmd` for one device, to go down its next link.
1037    pub fn queue_for(device: &str, username: &str, cmd: &LinkCommand) {
1038        let (Some(db), Ok(text)) = (db(), serde_json::to_string(cmd)) else {
1039            return;
1040        };
1041        let _ = db.conn.execute(
1042            "INSERT INTO link_outbox (device, username, command, created_at) VALUES (?1, ?2, ?3, ?4)",
1043            rusqlite::params![device, username, text, chrono::Utc::now().timestamp()],
1044        );
1045    }
1046
1047    pub fn save_push(username: &str, device: &str, token: &str, sandbox: bool) {
1048        let Some(db) = db() else { return };
1049        let _ = db.conn.execute(
1050            "INSERT INTO link_push (device, username, token, sandbox, updated_at) VALUES (?1, ?2, ?3, ?4, ?5)
1051             ON CONFLICT (device, username) DO UPDATE SET token = ?3, sandbox = ?4, updated_at = ?5",
1052            rusqlite::params![device, username, token, sandbox, chrono::Utc::now().timestamp()],
1053        );
1054    }
1055
1056    pub fn forget_push(username: &str, device: &str) {
1057        if let Some(db) = db() {
1058            let _ = db.conn.execute(
1059                "DELETE FROM link_push WHERE device = ?1 AND username = ?2",
1060                [device, username],
1061            );
1062        }
1063    }
1064
1065    /// Devices in scope with a push token, most recently seen first.
1066    pub fn push_targets(username: Option<&str>) -> Vec<PushTarget> {
1067        let Some(db) = db() else { return Vec::new() };
1068        db.conn
1069            .prepare(
1070                "SELECT p.device, p.username, d.name, d.platform, p.token, p.sandbox
1071                   FROM link_push p JOIN link_devices d ON d.device = p.device AND d.username = p.username
1072                  WHERE ?1 IS NULL OR p.username = ?1
1073                  ORDER BY d.last_seen DESC",
1074            )
1075            .and_then(|mut s| {
1076                s.query_map([username], |r| {
1077                    Ok(PushTarget {
1078                        device: r.get(0)?,
1079                        username: r.get(1)?,
1080                        name: r.get(2)?,
1081                        platform: r.get(3)?,
1082                        token: r.get(4)?,
1083                        sandbox: r.get(5)?,
1084                    })
1085                })?
1086                .collect()
1087            })
1088            .unwrap_or_default()
1089    }
1090
1091    /// The row id of a track a command names by uid or row id.
1092    pub fn track_row(id: &str) -> Option<i64> {
1093        let db = db()?;
1094        queries::resolve_id(&db.conn, queries::UidKind::Track, id)
1095            .ok()
1096            .flatten()
1097    }
1098
1099    /// What a playback command would play, for a notification to say:
1100    /// "Golden Standard — Tony Petersen", or a track and how many follow.
1101    pub fn describe(cmd: &LinkCommand) -> Option<String> {
1102        let ids: Vec<&String> = match cmd {
1103            LinkCommand::Play { track_ids, .. }
1104            | LinkCommand::Enqueue { track_ids }
1105            | LinkCommand::PlayNext { track_ids } => track_ids.iter().collect(),
1106            LinkCommand::JumpTo { track_id } => vec![track_id],
1107            _ => return None,
1108        };
1109        let db = db()?;
1110        let ids: Vec<i64> = ids
1111            .into_iter()
1112            .filter_map(|t| {
1113                queries::resolve_id(&db.conn, queries::UidKind::Track, t)
1114                    .ok()
1115                    .flatten()
1116            })
1117            .collect();
1118        let row = |id: i64| {
1119            db.conn
1120                .query_row(
1121                    "SELECT t.title, COALESCE(a.name, ''), COALESCE(al.title, ''), t.album_id
1122                       FROM tracks t LEFT JOIN artists a ON a.id = t.artist_id
1123                       LEFT JOIN albums al ON al.id = t.album_id WHERE t.id = ?1",
1124                    [id],
1125                    |r| {
1126                        Ok((
1127                            r.get::<_, String>(0)?,
1128                            r.get::<_, String>(1)?,
1129                            r.get::<_, String>(2)?,
1130                            r.get::<_, Option<i64>>(3)?,
1131                        ))
1132                    },
1133                )
1134                .ok()
1135        };
1136        let (title, artist, album, album_id) = row(*ids.first()?)?;
1137        let one_album = ids.len() > 1
1138            && album_id.is_some()
1139            && ids
1140                .iter()
1141                .all(|id| row(*id).is_some_and(|r| r.3 == album_id));
1142        Some(match (one_album, ids.len()) {
1143            (true, _) => format!("{album} — {artist}"),
1144            (false, 1) => format!("{title} — {artist}"),
1145            (false, n) => format!("{title} — {artist}, and {} more", n - 1),
1146        })
1147    }
1148}
1149
1150/// How long ago a client can have stopped playing and still be the obvious
1151/// one to send music to.
1152const RECENT: i64 = 6 * 60 * 60;
1153
1154fn pick(clients: &[ClientInfo], now: i64) -> Result<&ClientInfo, String> {
1155    if let Some(c) = clients.iter().find(|c| c.state.playing) {
1156        return Ok(c);
1157    }
1158    if let Some(c) = clients
1159        .iter()
1160        .filter(|c| c.last_played_at.is_some_and(|t| now - t < RECENT))
1161        .max_by_key(|c| c.last_played_at)
1162    {
1163        return Ok(c);
1164    }
1165    match clients {
1166        [] => Err("no koan app is linked to this server; open koan on the device".into()),
1167        [only] => Ok(only),
1168        several => Err(format!(
1169            "several koan apps are linked and none has played recently: {}. Ask which, then pass `client`",
1170            several
1171                .iter()
1172                .map(|c| c.name.as_str())
1173                .collect::<Vec<_>>()
1174                .join(", ")
1175        )),
1176    }
1177}
1178
1179#[cfg(test)]
1180mod tests {
1181
1182    #[test]
1183    fn a_playlist_order_adds_the_named_tracks_once_they_arrive() {
1184        let dir = tempfile::tempdir().unwrap();
1185        let path = dir.path().join("koan.db");
1186        let db = koan_core::db::connection::Database::open(&path).unwrap();
1187        let playlist = koan_core::db::queries::create_playlist(
1188            &db.conn,
1189            koan_core::db::queries::LOCAL_USER,
1190            "cyberpunk",
1191            None,
1192        )
1193        .unwrap();
1194        let order = Order {
1195            id: "o1".into(),
1196            username: None,
1197            client: None,
1198            artist: "Perturbator".into(),
1199            album: "Dangerous Days".into(),
1200            play_next: false,
1201            playlist: Some(playlist),
1202            titles: vec!["Future Club".into()],
1203            created_at: chrono::Utc::now().timestamp(),
1204        };
1205        registry().add_order(order);
1206
1207        // Not in the library yet: nothing happens, and the order waits.
1208        fulfil_from(&path);
1209        assert!(registry().orders(None).iter().any(|o| o.id == "o1"));
1210
1211        db.conn
1212            .execute_batch(
1213                "INSERT INTO artists (id, name) VALUES (1, 'Perturbator');
1214                 INSERT INTO albums (id, title, artist_id) VALUES (1, 'Dangerous Days', 1);
1215                 INSERT INTO tracks (id, title, album_id, artist_id, track_number, path) VALUES
1216                   (1, 'Welcome Back', 1, 1, 1, '/1.flac'), (2, 'Future Club', 1, 1, 2, '/2.flac');",
1217            )
1218            .unwrap();
1219        fulfil_from(&path);
1220        let held: Vec<i64> = db
1221            .conn
1222            .prepare("SELECT track_id FROM playlist_tracks WHERE playlist_id = ?1")
1223            .unwrap()
1224            .query_map([playlist], |r| r.get(0))
1225            .unwrap()
1226            .collect::<Result<_, _>>()
1227            .unwrap();
1228        assert_eq!(held, [2]);
1229        assert!(!registry().orders(None).iter().any(|o| o.id == "o1"));
1230    }
1231
1232    use super::*;
1233
1234    #[test]
1235    fn a_notification_cover_is_named_by_the_uid_the_command_carries() {
1236        // Commands carry uids; parsing them as row ids found no cover.
1237        let uid = "0190a5b2-7c3d-7e4f-8a1b-2c3d4e5f6a7b".to_string();
1238        let play = LinkCommand::Play {
1239            track_ids: vec!["x".into(), uid.clone()],
1240            start_at: 1,
1241            position_ms: 0,
1242        };
1243        assert_eq!(cover_track(&play), Some(uid.as_str()));
1244    }
1245
1246    #[test]
1247    fn a_reconnect_replaces_the_device_and_commands_reach_it() {
1248        let reg = Registry::default();
1249        let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
1250        let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
1251        let (tx3, _rx3) = tokio::sync::mpsc::unbounded_channel();
1252        reg.register("j", "phone", "ios", "dev-1", tx1, false);
1253        let id = reg.register("j", "phone", "ios", "dev-1", tx2, false);
1254        reg.register("someone", "laptop", "macos", "dev-2", tx3, false);
1255
1256        assert_eq!(reg.list(Some("j")).len(), 1);
1257        assert_eq!(reg.list(None).len(), 2);
1258
1259        let sent = reg.send(Some("j"), None, LinkCommand::Pause).unwrap();
1260        assert_eq!(sent.id, id);
1261        assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
1262
1263        // Another account's device is not this account's to command.
1264        assert!(
1265            reg.send(Some("j"), Some("laptop"), LinkCommand::Pause)
1266                .is_err()
1267        );
1268
1269        reg.unregister(&id);
1270        assert!(reg.send(Some("j"), None, LinkCommand::Pause).is_err());
1271    }
1272
1273    #[test]
1274    fn disconnecting_an_account_closes_only_its_links() {
1275        let reg = Registry::default();
1276        let (tx1, mut rx1) = tokio::sync::mpsc::unbounded_channel();
1277        let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
1278        reg.register("j", "phone", "ios", "dev-1", tx1, false);
1279        reg.register("someone", "laptop", "macos", "dev-2", tx2, false);
1280
1281        reg.disconnect("j");
1282        assert!(reg.list(Some("j")).is_empty());
1283        // The session sees its channel close, which is what ends it.
1284        assert!(matches!(
1285            rx1.try_recv(),
1286            Err(tokio::sync::mpsc::error::TryRecvError::Disconnected)
1287        ));
1288        assert_eq!(reg.list(Some("someone")).len(), 1);
1289        assert!(matches!(
1290            rx2.try_recv(),
1291            Err(tokio::sync::mpsc::error::TryRecvError::Empty)
1292        ));
1293    }
1294
1295    #[test]
1296    fn a_burst_of_changes_wakes_each_device_once_when_it_goes_quiet() {
1297        let t0 = std::time::Instant::now();
1298        let s = std::time::Duration::from_secs;
1299        let phone = || ("dev-1".to_string(), "j".to_string());
1300        let ipad = || ("dev-2".to_string(), "j".to_string());
1301        let mut wakes = Wakes {
1302            pending: Vec::new(),
1303            timer: false,
1304        };
1305
1306        wakes.add([phone()], t0);
1307        wakes.add([phone(), ipad()], t0 + s(10));
1308        wakes.add([phone()], t0 + s(20));
1309        assert_eq!(wakes.pending.len(), 2);
1310
1311        // Quiet is measured from each device's last change.
1312        assert!(wakes.take_due(t0 + s(39)).is_empty());
1313        assert_eq!(wakes.take_due(t0 + s(40)), [ipad()]);
1314        assert_eq!(wakes.next(), Some(t0 + s(50)));
1315        assert_eq!(wakes.take_due(t0 + s(50)), [phone()]);
1316        assert_eq!(wakes.next(), None);
1317    }
1318
1319    #[test]
1320    fn a_library_that_never_goes_quiet_still_wakes_devices() {
1321        let t0 = std::time::Instant::now();
1322        let phone = || ("dev-1".to_string(), "j".to_string());
1323        let mut wakes = Wakes {
1324            pending: Vec::new(),
1325            timer: false,
1326        };
1327        let mut sent = 0;
1328        for i in 0..40 {
1329            let now = t0 + std::time::Duration::from_secs(i * 10);
1330            sent += wakes.take_due(now).len();
1331            wakes.add([phone()], now);
1332        }
1333        // Six and a half minutes of changes every ten seconds: one push at
1334        // the five-minute mark, and one pending.
1335        assert_eq!(sent, 1);
1336        assert_eq!(wakes.pending.len(), 1);
1337    }
1338
1339    #[test]
1340    fn the_device_playing_is_the_one_meant() {
1341        let reg = Registry::default();
1342        let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
1343        let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
1344        let mac = reg.register("j", "mac", "macos", "dev-1", tx1, false);
1345        let phone = reg.register("j", "phone", "ios", "dev-2", tx2, false);
1346
1347        // Two idle devices: no telling, so the caller is told to ask.
1348        let err = reg.send(Some("j"), None, LinkCommand::Pause).unwrap_err();
1349        assert!(err.contains("mac") && err.contains("phone"), "{err}");
1350
1351        reg.report(
1352            &phone,
1353            LinkState {
1354                playing: true,
1355                ..Default::default()
1356            },
1357        );
1358        assert_eq!(
1359            reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
1360            phone
1361        );
1362        assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
1363
1364        // Stopped a moment ago: still the one meant, over the Mac.
1365        reg.report(&phone, LinkState::default());
1366        assert_eq!(
1367            reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
1368            phone
1369        );
1370        let _ = mac;
1371    }
1372
1373    #[test]
1374    fn a_device_hears_its_peers_and_commands_reach_them_by_device_id() {
1375        let reg = Registry::default();
1376        let (tx1, mut rx1) = tokio::sync::mpsc::unbounded_channel();
1377        let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
1378        let (tx3, mut rx3) = tokio::sync::mpsc::unbounded_channel();
1379        reg.register("j", "mac", "macos", "dev-mac", tx1, true);
1380        let phone = reg.register("j", "phone", "ios", "dev-phone", tx2, true);
1381        reg.register("someone", "laptop", "macos", "dev-other", tx3, true);
1382
1383        reg.report(
1384            &phone,
1385            LinkState {
1386                playing: true,
1387                title: Some("Roygbiv".into()),
1388                ..Default::default()
1389            },
1390        );
1391        let mut last = None;
1392        while let Ok(LinkCommand::Devices { devices }) = rx1.try_recv() {
1393            last = Some(devices);
1394        }
1395        let devices = last.expect("the Mac is told of the phone");
1396        assert_eq!(devices.len(), 1, "not itself, not another account's");
1397        assert_eq!(devices[0].id, "dev-phone");
1398        assert_eq!(
1399            devices[0].state.as_ref().and_then(|s| s.title.as_deref()),
1400            Some("Roygbiv")
1401        );
1402        while rx3.try_recv().is_ok() {}
1403        assert!(rx3.try_recv().is_err());
1404
1405        while rx2.try_recv().is_ok() {}
1406        reg.relay("j", "dev-phone", LinkCommand::Pause).unwrap();
1407        assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
1408        assert!(
1409            reg.relay("someone", "dev-phone", LinkCommand::Pause)
1410                .is_err(),
1411            "another account cannot reach it"
1412        );
1413        assert!(
1414            reg.relay("j", "dev-phone", LinkCommand::Devices { devices: vec![] })
1415                .is_err()
1416        );
1417    }
1418}