Skip to main content

koan_server/
clients.rs

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