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