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