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