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