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