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::remote::link::{LinkCommand, LinkState};
12use parking_lot::Mutex;
13use tokio::sync::mpsc::UnboundedSender;
14
15/// A linked client, as listed.
16#[derive(Debug, Clone)]
17pub struct ClientInfo {
18    pub id: String,
19    pub name: String,
20    pub platform: String,
21    pub username: String,
22    /// Unix seconds.
23    pub connected_at: i64,
24    /// What the client last said it was doing.
25    pub state: LinkState,
26    /// Unix seconds; when it was last seen playing, if ever since linking.
27    pub last_played_at: Option<i64>,
28    /// When `state` was reported, in Unix milliseconds.
29    pub state_at: i64,
30    /// Whether the client has reported its state at all. An app older than
31    /// the reports never does, and its `state` then says nothing about it.
32    pub reports: bool,
33    /// Reached by a notification it shows rather than over its link: iOS had
34    /// suspended it. What was sent runs when someone taps the notification.
35    pub notified: bool,
36}
37
38impl ClientInfo {
39    /// Where the playhead is now, from where it was reported to be.
40    pub fn position_ms(&self) -> u64 {
41        let pos = self.state.position_ms;
42        if !self.state.playing {
43            return pos;
44        }
45        let run = (chrono::Utc::now().timestamp_millis() - self.state_at).max(0) as u64;
46        let pos = pos + run;
47        if self.state.duration_ms > 0 {
48            pos.min(self.state.duration_ms)
49        } else {
50            pos
51        }
52    }
53}
54
55struct Entry {
56    info: ClientInfo,
57    /// The client's own id for itself, so a reconnect replaces its entry.
58    device: String,
59    tx: UnboundedSender<LinkCommand>,
60}
61
62/// "When this album is in the library, queue it on my device": a request made
63/// before the album exists, fulfilled by the scan that finds it.
64#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
65pub struct Order {
66    pub id: String,
67    /// Whose devices it may go to.
68    pub username: Option<String>,
69    /// A client id or name; `None` for whichever `send` would pick then.
70    pub client: Option<String>,
71    pub artist: String,
72    pub album: String,
73    /// Insert after the current track rather than at the end.
74    pub play_next: bool,
75    /// Add to this playlist rather than a device's queue.
76    #[serde(default)]
77    pub playlist: Option<i64>,
78    /// Only these tracks of the album (title substrings, in this order);
79    /// empty for all of it.
80    #[serde(default)]
81    pub titles: Vec<String>,
82    /// Unix seconds.
83    pub created_at: i64,
84}
85
86/// An order nobody's scan has fulfilled in this long is dropped.
87const ORDER_TTL: i64 = 24 * 60 * 60;
88
89#[derive(Default)]
90pub struct Registry {
91    entries: Mutex<Vec<Entry>>,
92    orders: Mutex<Vec<Order>>,
93}
94
95/// One registry per process: the WebSocket route and the GraphQL schema are
96/// built in different places and both need it.
97pub fn registry() -> &'static Registry {
98    static REGISTRY: LazyLock<Registry> = LazyLock::new(|| {
99        let registry = Registry::default();
100        *registry.orders.lock() = outbox::load_orders();
101        registry
102    });
103    &REGISTRY
104}
105
106impl Registry {
107    /// Add a client, replacing any earlier connection from the same device
108    /// and account. Returns its id.
109    pub fn register(
110        &self,
111        username: &str,
112        name: &str,
113        platform: &str,
114        device: &str,
115        tx: UnboundedSender<LinkCommand>,
116    ) -> String {
117        let id = uuid::Uuid::now_v7().to_string();
118        // What waited for this device while it was away goes down the new
119        // link first.
120        for cmd in outbox::take_and_remember(username, device, name, platform) {
121            let _ = tx.send(cmd);
122        }
123        let mut entries = self.entries.lock();
124        entries.retain(|e| !(e.device == device && e.info.username == username));
125        entries.push(Entry {
126            info: ClientInfo {
127                id: id.clone(),
128                name: name.to_string(),
129                platform: platform.to_string(),
130                username: username.to_string(),
131                connected_at: chrono::Utc::now().timestamp(),
132                state: LinkState::default(),
133                last_played_at: None,
134                state_at: chrono::Utc::now().timestamp_millis(),
135                reports: false,
136                notified: false,
137            },
138            device: device.to_string(),
139            tx,
140        });
141        id
142    }
143
144    /// Record what a client says it is doing.
145    pub fn report(&self, id: &str, state: LinkState) {
146        let mut entries = self.entries.lock();
147        if let Some(e) = entries.iter_mut().find(|e| e.info.id == id) {
148            if state.playing || e.info.state.playing {
149                e.info.last_played_at = Some(chrono::Utc::now().timestamp());
150            }
151            e.info.state = state;
152            e.info.state_at = chrono::Utc::now().timestamp_millis();
153            e.info.reports = true;
154        }
155    }
156
157    /// Record where Apple's push service reaches this device.
158    pub fn set_push(&self, username: &str, device: &str, token: &str, sandbox: bool) {
159        outbox::save_push(username, device, token, sandbox);
160    }
161
162    pub fn unregister(&self, id: &str) {
163        self.entries.lock().retain(|e| e.info.id != id);
164    }
165
166    /// Clients `username` may command, newest first; every client for `None`.
167    pub fn list(&self, username: Option<&str>) -> Vec<ClientInfo> {
168        let mut out: Vec<ClientInfo> = self
169            .entries
170            .lock()
171            .iter()
172            .filter(|e| username.is_none_or(|u| e.info.username == u))
173            .map(|e| e.info.clone())
174            .collect();
175        out.sort_by_key(|c| std::cmp::Reverse(c.connected_at));
176        out
177    }
178
179    /// Send to `id` (an id or a name), or with none to the client the
180    /// command most likely means: the one playing, else the one that played
181    /// within `RECENT`, else the only one linked. `Err` names the choices when
182    /// there is no telling, so whoever asked can ask the person.
183    pub fn send(
184        &self,
185        username: Option<&str>,
186        id: Option<&str>,
187        cmd: LinkCommand,
188    ) -> Result<ClientInfo, String> {
189        let clients = self.list(username);
190        let target = match id {
191            Some(id) => clients
192                .iter()
193                .find(|c| c.id == id || c.name.eq_ignore_ascii_case(id)),
194            None if clients.is_empty() => None,
195            None => Some(pick(&clients, chrono::Utc::now().timestamp())?),
196        };
197        // Not linked: a phone iOS has suspended can still be asked, by a
198        // notification it shows.
199        let Some(target) = target else {
200            return notify(username, id, &cmd).unwrap_or_else(|| {
201                Err(match id {
202                    Some(id) => format!("no linked client {id}; see `clients`"),
203                    None => "no koan app is linked to this server; open koan on the device".into(),
204                })
205            });
206        };
207        let entries = self.entries.lock();
208        let entry = entries
209            .iter()
210            .find(|e| e.info.id == target.id)
211            .ok_or("that client has just gone")?;
212        entry
213            .tx
214            .send(cmd)
215            .map_err(|_| "that client has just gone".to_string())?;
216        Ok(target.clone())
217    }
218}
219
220impl Registry {
221    pub fn add_order(&self, order: Order) {
222        outbox::save_order(&order);
223        self.orders.lock().push(order);
224    }
225
226    pub fn orders(&self, username: Option<&str>) -> Vec<Order> {
227        self.orders
228            .lock()
229            .iter()
230            .filter(|o| username.is_none() || o.username.as_deref() == username)
231            .cloned()
232            .collect()
233    }
234
235    pub fn cancel_order(&self, username: Option<&str>, id: &str) -> bool {
236        let mut orders = self.orders.lock();
237        let before = orders.len();
238        orders
239            .retain(|o| !(o.id == id && (username.is_none() || o.username.as_deref() == username)));
240        let gone = orders.len() != before;
241        if gone {
242            outbox::drop_order(id);
243        }
244        gone
245    }
246
247    fn done(&self, id: &str) {
248        self.orders.lock().retain(|o| o.id != id);
249        outbox::drop_order(id);
250    }
251
252    /// Send every order whose album the library now holds, and drop it.
253    /// `find` answers an order with the album's track ids, in order.
254    pub fn fulfil_orders(&self, find: impl Fn(&Order) -> Option<Vec<i64>>) {
255        let now = chrono::Utc::now().timestamp();
256        let pending: Vec<Order> = {
257            let mut orders = self.orders.lock();
258            for o in orders.iter().filter(|o| now - o.created_at >= ORDER_TTL) {
259                outbox::drop_order(&o.id);
260            }
261            orders.retain(|o| now - o.created_at < ORDER_TTL);
262            orders.clone()
263        };
264        for order in pending.into_iter().filter(|o| o.playlist.is_none()) {
265            let Some(ids) = find(&order).filter(|ids| !ids.is_empty()) else {
266                continue;
267            };
268            let track_ids = ids.iter().map(i64::to_string).collect();
269            let cmd = if order.play_next {
270                LinkCommand::PlayNext { track_ids }
271            } else {
272                LinkCommand::Enqueue { track_ids }
273            };
274            match self.send(order.username.as_deref(), order.client.as_deref(), cmd) {
275                Ok(c) => {
276                    log::info!(
277                        "link: {} — {} arrived; queued on {}",
278                        order.artist,
279                        order.album,
280                        c.name
281                    );
282                    self.done(&order.id);
283                }
284                // No device to send to yet: kept, and tried after the next scan.
285                Err(e) => log::info!("link: {} — {} arrived but {e}", order.artist, order.album),
286            }
287        }
288    }
289}
290
291/// Fulfil standing orders against the library at `db_path`.
292pub fn fulfil_from(db_path: &std::path::Path) {
293    let registry = registry();
294    if registry.orders.lock().is_empty() {
295        return;
296    }
297    let Ok(db) = koan_core::db::connection::Database::open(db_path) else {
298        return;
299    };
300    // Playlist orders are the server's own to carry out: add the tracks, and
301    // the playlist reaches every device like any other edit.
302    let for_playlists: Vec<Order> = registry
303        .orders
304        .lock()
305        .iter()
306        .filter(|o| o.playlist.is_some())
307        .cloned()
308        .collect();
309    let mut edited = false;
310    for order in for_playlists {
311        let Some(playlist) = order.playlist else {
312            continue;
313        };
314        let Some(ids) = order_tracks(&db.conn, &order) else {
315            continue;
316        };
317        match koan_core::db::queries::add_tracks(&db.conn, playlist, &ids) {
318            Ok(_) => {
319                // Tests would push to whatever server this machine signs in to.
320                if !cfg!(test) {
321                    koan_core::playlists::push_to_remote(playlist);
322                }
323                log::info!(
324                    "link: {} — {} arrived; added {} tracks to playlist {playlist}",
325                    order.artist,
326                    order.album,
327                    ids.len()
328                );
329                registry.done(&order.id);
330                edited = true;
331            }
332            Err(e) => log::warn!("link: could not add to playlist {playlist}: {e}"),
333        }
334    }
335    if edited {
336        changed();
337    }
338    registry.fulfil_orders(|order| order_tracks(&db.conn, order));
339}
340
341/// The tracks an order asks for, once its album is in the library: all of
342/// it, or the named ones in the order named. `None` until they are there.
343fn order_tracks(conn: &rusqlite::Connection, order: &Order) -> Option<Vec<i64>> {
344    let tracks = album_tracks(conn, &order.artist, &order.album)?;
345    if order.titles.is_empty() {
346        return Some(tracks.into_iter().map(|(id, _)| id).collect());
347    }
348    let picked: Vec<i64> = order
349        .titles
350        .iter()
351        .filter_map(|want| {
352            let want = want.to_lowercase();
353            tracks
354                .iter()
355                .find(|(_, t)| t.to_lowercase().contains(&want))
356                .map(|(id, _)| *id)
357        })
358        .collect();
359    (!picked.is_empty()).then_some(picked)
360}
361
362/// The newest album whose artist and title contain these, as its tracks in
363/// disc and track order.
364pub fn album_tracks(
365    conn: &rusqlite::Connection,
366    artist: &str,
367    album: &str,
368) -> Option<Vec<(i64, String)>> {
369    let like = |s: &str| format!("%{}%", s.replace(['%', '_'], ""));
370    let album_id: i64 = conn
371        .query_row(
372            "SELECT al.id FROM albums al JOIN artists a ON a.id = al.artist_id
373              WHERE a.name LIKE ?1 COLLATE NOCASE AND al.title LIKE ?2 COLLATE NOCASE
374              ORDER BY al.id DESC LIMIT 1",
375            [like(artist), like(album)],
376            |r| r.get(0),
377        )
378        .ok()?;
379    let mut stmt = conn
380        .prepare("SELECT id, title FROM tracks WHERE album_id = ?1 ORDER BY disc, track_number, id")
381        .ok()?;
382    let tracks = stmt
383        .query_map([album_id], |r| Ok((r.get(0)?, r.get(1)?)))
384        .ok()?
385        .filter_map(Result::ok)
386        .collect();
387    Some(tracks)
388}
389
390impl Registry {
391    /// Send to every client `username` may command. The names of those it
392    /// reached.
393    pub fn broadcast(&self, username: Option<&str>, cmd: LinkCommand) -> Vec<String> {
394        let ids: Vec<String> = self.list(username).into_iter().map(|c| c.id).collect();
395        let entries = self.entries.lock();
396        entries
397            .iter()
398            .filter(|e| ids.contains(&e.info.id) && e.tx.send(cmd.clone()).is_ok())
399            .map(|e| e.info.name.clone())
400            .collect()
401    }
402}
403
404impl Registry {
405    /// Send to every device `username` may command: at once to those linked,
406    /// and to those that have linked before but are away now, when they next
407    /// link. For commands still right hours later (`Sync`, `Evict`), never
408    /// playback. The names reached now, and the names it waits for.
409    pub fn deliver(&self, username: Option<&str>, cmd: LinkCommand) -> (Vec<String>, Vec<String>) {
410        let sent = self.broadcast(username, cmd.clone());
411        let live: Vec<(String, String)> = self
412            .entries
413            .lock()
414            .iter()
415            .map(|e| (e.device.clone(), e.info.username.clone()))
416            .collect();
417        let queued = outbox::queue_for_absent(username, &live, &cmd);
418        wake(&queued);
419        (sent, queued.into_iter().map(|q| q.name).collect())
420    }
421}
422
423/// Wake these absent devices, where a push can: each links, and takes what
424/// was just queued for it.
425fn wake(devices: &[outbox::Absent]) {
426    let Some(pusher) = crate::push::pusher() else {
427        return;
428    };
429    let targets: Vec<outbox::PushTarget> = outbox::push_targets(None)
430        .into_iter()
431        .filter(|t| {
432            devices
433                .iter()
434                .any(|d| d.device == t.device && d.username == t.username)
435        })
436        .collect();
437    if targets.is_empty() {
438        return;
439    }
440    std::thread::spawn(move || {
441        for t in targets {
442            deliver_push(pusher, &t, &crate::push::Push::Wake);
443        }
444    });
445}
446
447/// Ask a device that is not linked to run `cmd`, by a notification. `None`
448/// when a notification cannot carry it: not a request to play, no push key,
449/// or no device in scope that has given a push token.
450fn notify(
451    username: Option<&str>,
452    id: Option<&str>,
453    cmd: &LinkCommand,
454) -> Option<Result<ClientInfo, String>> {
455    let verb = match cmd {
456        LinkCommand::Play { .. } | LinkCommand::JumpTo { .. } => "Play",
457        LinkCommand::PlayNext { .. } => "Play next",
458        LinkCommand::Enqueue { .. } => "Add to queue",
459        _ => return None,
460    };
461    let pusher = crate::push::pusher()?;
462    let targets: Vec<outbox::PushTarget> = outbox::push_targets(username)
463        .into_iter()
464        .filter(|t| id.is_none_or(|id| t.device == id || t.name.eq_ignore_ascii_case(id)))
465        .collect();
466    let target = match targets.as_slice() {
467        [] => return None,
468        [only] => only.clone(),
469        several => {
470            return Some(Err(format!(
471                "no koan app is linked, and several can be sent a notification: {}. Ask which, then pass `client`",
472                several
473                    .iter()
474                    .map(|t| t.name.as_str())
475                    .collect::<Vec<_>>()
476                    .join(", ")
477            )));
478        }
479    };
480    let command = serde_json::to_value(cmd).ok()?;
481    let push = crate::push::Push::Notify {
482        title: format!("{verb} on {}", target.name),
483        body: outbox::describe(cmd).unwrap_or_else(|| "From your koan server".into()),
484        command,
485    };
486    let info = ClientInfo {
487        id: target.device.clone(),
488        name: target.name.clone(),
489        platform: target.platform.clone(),
490        username: target.username.clone(),
491        connected_at: 0,
492        state: LinkState::default(),
493        last_played_at: None,
494        state_at: 0,
495        reports: false,
496        notified: true,
497    };
498    std::thread::spawn(move || deliver_push(pusher, &target, &push));
499    Some(Ok(info))
500}
501
502/// Send one push, forgetting a token Apple says is no longer good.
503fn deliver_push(
504    pusher: &crate::push::Pusher,
505    target: &outbox::PushTarget,
506    push: &crate::push::Push,
507) {
508    use crate::push::Outcome;
509    match pusher.send(&target.token, target.sandbox, push) {
510        Outcome::Sent => log::info!("push: sent to {}", target.name),
511        Outcome::Gone => {
512            log::info!(
513                "push: {}'s token is no longer valid; forgotten",
514                target.name
515            );
516            outbox::forget_push(&target.username, &target.device);
517        }
518        Outcome::Failed(e) => log::warn!("push: to {} failed: {e}", target.name),
519    }
520}
521
522/// Have every device pull what the server just changed (a playlist edited,
523/// albums added): at once where linked, on next link where not. Syncs waiting
524/// for a device collapse into one.
525pub fn changed() {
526    registry().deliver(None, LinkCommand::Sync { full: false });
527}
528
529/// After a library scan: if the library holds different tracks or albums from
530/// the last scan, tell every device. A scan that found nothing new, which is
531/// most of them, sends nothing.
532pub fn changed_if_library_moved(db_path: &std::path::Path) {
533    static LAST: parking_lot::Mutex<Option<(i64, i64, i64)>> = parking_lot::Mutex::new(None);
534    let Ok(db) = koan_core::db::connection::Database::open(db_path) else {
535        return;
536    };
537    let Ok(now) = db.conn.query_row(
538        "SELECT (SELECT COUNT(*) FROM tracks), (SELECT COALESCE(MAX(id), 0) FROM tracks),
539                (SELECT COUNT(*) FROM albums)",
540        [],
541        |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
542    ) else {
543        return;
544    };
545    let before = LAST.lock().replace(now);
546    if before.is_some_and(|b| b != now) {
547        changed();
548    }
549}
550
551/// The server-side record of devices and their waiting commands, in the
552/// library database so it outlives a restart.
553mod outbox {
554    use koan_core::remote::link::LinkCommand;
555
556    /// Dropped undelivered after this long: a device away a month re-syncs
557    /// on its own when opened.
558    const KEEP_SECS: i64 = 30 * 24 * 60 * 60;
559
560    /// Tests keep to memory: the configured database is whoever ran them.
561    fn db() -> Option<koan_core::db::connection::Database> {
562        if cfg!(test) {
563            return None;
564        }
565        koan_core::db::connection::Database::open(&koan_core::config::db_path()).ok()
566    }
567
568    pub fn load_orders() -> Vec<super::Order> {
569        let Some(db) = db() else { return Vec::new() };
570        db.conn
571            .prepare("SELECT body FROM link_orders ORDER BY created_at")
572            .and_then(|mut s| {
573                s.query_map([], |r| r.get::<_, String>(0))?
574                    .collect::<Result<Vec<_>, _>>()
575            })
576            .unwrap_or_default()
577            .into_iter()
578            .filter_map(|b| serde_json::from_str(&b).ok())
579            .collect()
580    }
581
582    pub fn save_order(order: &super::Order) {
583        let (Some(db), Ok(body)) = (db(), serde_json::to_string(order)) else {
584            return;
585        };
586        let _ = db.conn.execute(
587            "INSERT OR REPLACE INTO link_orders (id, body, created_at) VALUES (?1, ?2, ?3)",
588            rusqlite::params![order.id, body, order.created_at],
589        );
590    }
591
592    pub fn drop_order(id: &str) {
593        if let Some(db) = db() {
594            let _ = db
595                .conn
596                .execute("DELETE FROM link_orders WHERE id = ?1", [id]);
597        }
598    }
599
600    pub fn take_and_remember(
601        username: &str,
602        device: &str,
603        name: &str,
604        platform: &str,
605    ) -> Vec<LinkCommand> {
606        let Some(db) = db() else { return Vec::new() };
607        let now = chrono::Utc::now().timestamp();
608        let _ = db.conn.execute(
609            "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES (?1, ?2, ?3, ?4, ?5)
610             ON CONFLICT (device, username) DO UPDATE SET name = ?3, platform = ?4, last_seen = ?5",
611            rusqlite::params![device, username, name, platform, now],
612        );
613        let _ = db.conn.execute(
614            "DELETE FROM link_outbox WHERE created_at < ?1",
615            [now - KEEP_SECS],
616        );
617        let waiting: Vec<(i64, String)> = db
618            .conn
619            .prepare("SELECT id, command FROM link_outbox WHERE device = ?1 AND username = ?2 ORDER BY id")
620            .and_then(|mut s| {
621                s.query_map([device, username], |r| Ok((r.get(0)?, r.get(1)?)))?
622                    .collect()
623            })
624            .unwrap_or_default();
625        let _ = db.conn.execute(
626            "DELETE FROM link_outbox WHERE device = ?1 AND username = ?2",
627            [device, username],
628        );
629        if !waiting.is_empty() {
630            log::info!("link: {} waiting commands for {name}", waiting.len());
631        }
632        waiting
633            .into_iter()
634            .filter_map(|(_, c)| serde_json::from_str(&c).ok())
635            .collect()
636    }
637
638    /// A known device that was not linked when something was queued for it.
639    pub struct Absent {
640        pub device: String,
641        pub username: String,
642        pub name: String,
643    }
644
645    /// Queue `cmd` for each known device in scope that is not in `live`.
646    pub fn queue_for_absent(
647        username: Option<&str>,
648        live: &[(String, String)],
649        cmd: &LinkCommand,
650    ) -> Vec<Absent> {
651        let Some(db) = db() else { return Vec::new() };
652        let known: Vec<(String, String, String)> = db
653            .conn
654            .prepare("SELECT device, username, name FROM link_devices")
655            .and_then(|mut s| {
656                s.query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))?
657                    .collect()
658            })
659            .unwrap_or_default();
660        let Ok(text) = serde_json::to_string(cmd) else {
661            return Vec::new();
662        };
663        let is_sync = matches!(cmd, LinkCommand::Sync { .. });
664        let now = chrono::Utc::now().timestamp();
665        let mut queued = Vec::new();
666        for (device, user, name) in known {
667            if username.is_some_and(|u| u != user)
668                || live.iter().any(|(d, u)| *d == device && *u == user)
669            {
670                continue;
671            }
672            if is_sync {
673                // One pending sync is enough; a full one covers an incremental.
674                let _ = db.conn.execute(
675                    "DELETE FROM link_outbox WHERE device = ?1 AND username = ?2 AND command LIKE '{\"type\":\"sync\"%'",
676                    [&device, &user],
677                );
678            }
679            if db
680                .conn
681                .execute(
682                    "INSERT INTO link_outbox (device, username, command, created_at) VALUES (?1, ?2, ?3, ?4)",
683                    rusqlite::params![device, user, text, now],
684                )
685                .is_ok()
686            {
687                queued.push(Absent {
688                    device,
689                    username: user,
690                    name,
691                });
692            }
693        }
694        queued
695    }
696
697    /// A device Apple's push service can reach.
698    #[derive(Clone)]
699    pub struct PushTarget {
700        pub device: String,
701        pub username: String,
702        pub name: String,
703        pub platform: String,
704        pub token: String,
705        pub sandbox: bool,
706    }
707
708    pub fn save_push(username: &str, device: &str, token: &str, sandbox: bool) {
709        let Some(db) = db() else { return };
710        let _ = db.conn.execute(
711            "INSERT INTO link_push (device, username, token, sandbox, updated_at) VALUES (?1, ?2, ?3, ?4, ?5)
712             ON CONFLICT (device, username) DO UPDATE SET token = ?3, sandbox = ?4, updated_at = ?5",
713            rusqlite::params![device, username, token, sandbox, chrono::Utc::now().timestamp()],
714        );
715    }
716
717    pub fn forget_push(username: &str, device: &str) {
718        if let Some(db) = db() {
719            let _ = db.conn.execute(
720                "DELETE FROM link_push WHERE device = ?1 AND username = ?2",
721                [device, username],
722            );
723        }
724    }
725
726    /// Devices in scope with a push token, most recently seen first.
727    pub fn push_targets(username: Option<&str>) -> Vec<PushTarget> {
728        let Some(db) = db() else { return Vec::new() };
729        db.conn
730            .prepare(
731                "SELECT p.device, p.username, d.name, d.platform, p.token, p.sandbox
732                   FROM link_push p JOIN link_devices d ON d.device = p.device AND d.username = p.username
733                  WHERE ?1 IS NULL OR p.username = ?1
734                  ORDER BY d.last_seen DESC",
735            )
736            .and_then(|mut s| {
737                s.query_map([username], |r| {
738                    Ok(PushTarget {
739                        device: r.get(0)?,
740                        username: r.get(1)?,
741                        name: r.get(2)?,
742                        platform: r.get(3)?,
743                        token: r.get(4)?,
744                        sandbox: r.get(5)?,
745                    })
746                })?
747                .collect()
748            })
749            .unwrap_or_default()
750    }
751
752    /// What a playback command would play, for a notification to say:
753    /// "Golden Standard — Tony Petersen", or a track and how many follow.
754    pub fn describe(cmd: &LinkCommand) -> Option<String> {
755        let ids: Vec<i64> = match cmd {
756            LinkCommand::Play { track_ids, .. }
757            | LinkCommand::Enqueue { track_ids }
758            | LinkCommand::PlayNext { track_ids } => {
759                track_ids.iter().filter_map(|t| t.parse().ok()).collect()
760            }
761            LinkCommand::JumpTo { track_id } => vec![track_id.parse().ok()?],
762            _ => return None,
763        };
764        let db = db()?;
765        let row = |id: i64| {
766            db.conn
767                .query_row(
768                    "SELECT t.title, COALESCE(a.name, ''), COALESCE(al.title, ''), t.album_id
769                       FROM tracks t LEFT JOIN artists a ON a.id = t.artist_id
770                       LEFT JOIN albums al ON al.id = t.album_id WHERE t.id = ?1",
771                    [id],
772                    |r| {
773                        Ok((
774                            r.get::<_, String>(0)?,
775                            r.get::<_, String>(1)?,
776                            r.get::<_, String>(2)?,
777                            r.get::<_, Option<i64>>(3)?,
778                        ))
779                    },
780                )
781                .ok()
782        };
783        let (title, artist, album, album_id) = row(*ids.first()?)?;
784        let one_album = ids.len() > 1
785            && album_id.is_some()
786            && ids
787                .iter()
788                .all(|id| row(*id).is_some_and(|r| r.3 == album_id));
789        Some(match (one_album, ids.len()) {
790            (true, _) => format!("{album} — {artist}"),
791            (false, 1) => format!("{title} — {artist}"),
792            (false, n) => format!("{title} — {artist}, and {} more", n - 1),
793        })
794    }
795}
796
797/// How long ago a client can have stopped playing and still be the obvious
798/// one to send music to.
799const RECENT: i64 = 6 * 60 * 60;
800
801fn pick(clients: &[ClientInfo], now: i64) -> Result<&ClientInfo, String> {
802    if let Some(c) = clients.iter().find(|c| c.state.playing) {
803        return Ok(c);
804    }
805    if let Some(c) = clients
806        .iter()
807        .filter(|c| c.last_played_at.is_some_and(|t| now - t < RECENT))
808        .max_by_key(|c| c.last_played_at)
809    {
810        return Ok(c);
811    }
812    match clients {
813        [] => Err("no koan app is linked to this server; open koan on the device".into()),
814        [only] => Ok(only),
815        several => Err(format!(
816            "several koan apps are linked and none has played recently: {}. Ask which, then pass `client`",
817            several
818                .iter()
819                .map(|c| c.name.as_str())
820                .collect::<Vec<_>>()
821                .join(", ")
822        )),
823    }
824}
825
826#[cfg(test)]
827mod tests {
828
829    #[test]
830    fn a_playlist_order_adds_the_named_tracks_once_they_arrive() {
831        let dir = tempfile::tempdir().unwrap();
832        let path = dir.path().join("koan.db");
833        let db = koan_core::db::connection::Database::open(&path).unwrap();
834        let playlist =
835            koan_core::db::queries::create_playlist(&db.conn, "cyberpunk", None).unwrap();
836        let order = Order {
837            id: "o1".into(),
838            username: None,
839            client: None,
840            artist: "Perturbator".into(),
841            album: "Dangerous Days".into(),
842            play_next: false,
843            playlist: Some(playlist),
844            titles: vec!["Future Club".into()],
845            created_at: chrono::Utc::now().timestamp(),
846        };
847        registry().add_order(order);
848
849        // Not in the library yet: nothing happens, and the order waits.
850        fulfil_from(&path);
851        assert!(registry().orders(None).iter().any(|o| o.id == "o1"));
852
853        db.conn
854            .execute_batch(
855                "INSERT INTO artists (id, name) VALUES (1, 'Perturbator');
856                 INSERT INTO albums (id, title, artist_id) VALUES (1, 'Dangerous Days', 1);
857                 INSERT INTO tracks (id, title, album_id, artist_id, track_number, path) VALUES
858                   (1, 'Welcome Back', 1, 1, 1, '/1.flac'), (2, 'Future Club', 1, 1, 2, '/2.flac');",
859            )
860            .unwrap();
861        fulfil_from(&path);
862        let held: Vec<i64> = db
863            .conn
864            .prepare("SELECT track_id FROM playlist_tracks WHERE playlist_id = ?1")
865            .unwrap()
866            .query_map([playlist], |r| r.get(0))
867            .unwrap()
868            .collect::<Result<_, _>>()
869            .unwrap();
870        assert_eq!(held, [2]);
871        assert!(!registry().orders(None).iter().any(|o| o.id == "o1"));
872    }
873
874    use super::*;
875
876    #[test]
877    fn a_reconnect_replaces_the_device_and_commands_reach_it() {
878        let reg = Registry::default();
879        let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
880        let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
881        let (tx3, _rx3) = tokio::sync::mpsc::unbounded_channel();
882        reg.register("j", "phone", "ios", "dev-1", tx1);
883        let id = reg.register("j", "phone", "ios", "dev-1", tx2);
884        reg.register("someone", "laptop", "macos", "dev-2", tx3);
885
886        assert_eq!(reg.list(Some("j")).len(), 1);
887        assert_eq!(reg.list(None).len(), 2);
888
889        let sent = reg.send(Some("j"), None, LinkCommand::Pause).unwrap();
890        assert_eq!(sent.id, id);
891        assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
892
893        // Another account's device is not this account's to command.
894        assert!(
895            reg.send(Some("j"), Some("laptop"), LinkCommand::Pause)
896                .is_err()
897        );
898
899        reg.unregister(&id);
900        assert!(reg.send(Some("j"), None, LinkCommand::Pause).is_err());
901    }
902
903    #[test]
904    fn the_device_playing_is_the_one_meant() {
905        let reg = Registry::default();
906        let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
907        let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
908        let mac = reg.register("j", "mac", "macos", "dev-1", tx1);
909        let phone = reg.register("j", "phone", "ios", "dev-2", tx2);
910
911        // Two idle devices: no telling, so the caller is told to ask.
912        let err = reg.send(Some("j"), None, LinkCommand::Pause).unwrap_err();
913        assert!(err.contains("mac") && err.contains("phone"), "{err}");
914
915        reg.report(
916            &phone,
917            LinkState {
918                playing: true,
919                ..Default::default()
920            },
921        );
922        assert_eq!(
923            reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
924            phone
925        );
926        assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
927
928        // Stopped a moment ago: still the one meant, over the Mac.
929        reg.report(&phone, LinkState::default());
930        assert_eq!(
931            reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
932            phone
933        );
934        let _ = mac;
935    }
936}