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 outbox::Absent;
13use parking_lot::Mutex;
14use tokio::sync::mpsc::UnboundedSender;
15
16/// A linked client, as listed.
17#[derive(Debug, Clone)]
18pub struct ClientInfo {
19    pub id: String,
20    pub name: String,
21    pub platform: String,
22    pub username: String,
23    /// Unix seconds.
24    pub connected_at: i64,
25    /// What the client last said it was doing.
26    pub state: LinkState,
27    /// Unix seconds; when it was last seen playing, if ever since linking.
28    pub last_played_at: Option<i64>,
29    /// When `state` was reported, in Unix milliseconds.
30    pub state_at: i64,
31    /// Whether the client has reported its state at all. An app older than
32    /// the reports never does, and its `state` then says nothing about it.
33    pub reports: bool,
34    /// Reached by a notification it shows rather than over its link: iOS had
35    /// suspended it. What was sent runs when someone taps the notification.
36    pub notified: bool,
37}
38
39impl ClientInfo {
40    /// Where the playhead is now, from where it was reported to be.
41    pub fn position_ms(&self) -> u64 {
42        let pos = self.state.position_ms;
43        if !self.state.playing {
44            return pos;
45        }
46        let run = (chrono::Utc::now().timestamp_millis() - self.state_at).max(0) as u64;
47        let pos = pos + run;
48        if self.state.duration_ms > 0 {
49            pos.min(self.state.duration_ms)
50        } else {
51            pos
52        }
53    }
54}
55
56struct Entry {
57    info: ClientInfo,
58    /// The client's own id for itself, so a reconnect replaces its entry.
59    device: String,
60    tx: UnboundedSender<LinkCommand>,
61}
62
63/// "When this album is in the library, queue it on my device": a request made
64/// before the album exists, fulfilled by the scan that finds it.
65#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
66pub struct Order {
67    pub id: String,
68    /// Whose devices it may go to.
69    pub username: Option<String>,
70    /// A client id or name; `None` for whichever `send` would pick then.
71    pub client: Option<String>,
72    pub artist: String,
73    pub album: String,
74    /// Insert after the current track rather than at the end.
75    pub play_next: bool,
76    /// Add to this playlist rather than a device's queue.
77    #[serde(default)]
78    pub playlist: Option<i64>,
79    /// Only these tracks of the album (title substrings, in this order);
80    /// empty for all of it.
81    #[serde(default)]
82    pub titles: Vec<String>,
83    /// Unix seconds.
84    pub created_at: i64,
85}
86
87/// An order nobody's scan has fulfilled in this long is dropped.
88const ORDER_TTL: i64 = 24 * 60 * 60;
89
90#[derive(Default)]
91pub struct Registry {
92    entries: Mutex<Vec<Entry>>,
93    orders: Mutex<Vec<Order>>,
94}
95
96/// One registry per process: the WebSocket route and the GraphQL schema are
97/// built in different places and both need it.
98pub fn registry() -> &'static Registry {
99    static REGISTRY: LazyLock<Registry> = LazyLock::new(|| {
100        let registry = Registry::default();
101        *registry.orders.lock() = outbox::load_orders();
102        registry
103    });
104    &REGISTRY
105}
106
107impl Registry {
108    /// Add a client, replacing any earlier connection from the same device
109    /// and account. Returns its id.
110    pub fn register(
111        &self,
112        username: &str,
113        name: &str,
114        platform: &str,
115        device: &str,
116        tx: UnboundedSender<LinkCommand>,
117    ) -> String {
118        let id = uuid::Uuid::now_v7().to_string();
119        // What waited for this device while it was away goes down the new
120        // link first.
121        for cmd in outbox::take_and_remember(username, device, name, platform) {
122            let _ = tx.send(cmd);
123        }
124        let mut entries = self.entries.lock();
125        entries.retain(|e| !(e.device == device && e.info.username == username));
126        entries.push(Entry {
127            info: ClientInfo {
128                id: id.clone(),
129                name: name.to_string(),
130                platform: platform.to_string(),
131                username: username.to_string(),
132                connected_at: chrono::Utc::now().timestamp(),
133                state: LinkState::default(),
134                last_played_at: None,
135                state_at: chrono::Utc::now().timestamp_millis(),
136                reports: false,
137                notified: false,
138            },
139            device: device.to_string(),
140            tx,
141        });
142        id
143    }
144
145    /// Record what a client says it is doing.
146    pub fn report(&self, id: &str, state: LinkState) {
147        let mut entries = self.entries.lock();
148        if let Some(e) = entries.iter_mut().find(|e| e.info.id == id) {
149            if state.playing || e.info.state.playing {
150                e.info.last_played_at = Some(chrono::Utc::now().timestamp());
151            }
152            e.info.state = state;
153            e.info.state_at = chrono::Utc::now().timestamp_millis();
154            e.info.reports = true;
155        }
156    }
157
158    /// Record where Apple's push service reaches this device.
159    pub fn set_push(&self, username: &str, device: &str, token: &str, sandbox: bool) {
160        outbox::save_push(username, device, token, sandbox);
161    }
162
163    pub fn unregister(&self, id: &str) {
164        self.entries.lock().retain(|e| e.info.id != id);
165    }
166
167    /// Clients `username` may command, newest first; every client for `None`.
168    pub fn list(&self, username: Option<&str>) -> Vec<ClientInfo> {
169        let mut out: Vec<ClientInfo> = self
170            .entries
171            .lock()
172            .iter()
173            .filter(|e| username.is_none_or(|u| e.info.username == u))
174            .map(|e| e.info.clone())
175            .collect();
176        out.sort_by_key(|c| std::cmp::Reverse(c.connected_at));
177        out
178    }
179
180    /// Send to `id` (an id or a name), or with none to the client the
181    /// command most likely means: the one playing, else the one that played
182    /// within `RECENT`, else the only one linked. `Err` names the choices when
183    /// there is no telling, so whoever asked can ask the person.
184    pub fn send(
185        &self,
186        username: Option<&str>,
187        id: Option<&str>,
188        cmd: LinkCommand,
189    ) -> Result<ClientInfo, String> {
190        let clients = self.list(username);
191        let target = match id {
192            Some(id) => clients
193                .iter()
194                .find(|c| c.id == id || c.name.eq_ignore_ascii_case(id)),
195            None if clients.is_empty() => None,
196            None => Some(pick(&clients, chrono::Utc::now().timestamp())?),
197        };
198        // Not linked: a phone iOS has suspended is woken to take it.
199        let Some(target) = target else {
200            return reach_absent(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, queued) = self.link_or_queue(username, &cmd);
411        wake(&queued.iter().map(Absent::key).collect::<Vec<_>>());
412        (sent, queued.into_iter().map(|q| q.name).collect())
413    }
414
415    fn link_or_queue(
416        &self,
417        username: Option<&str>,
418        cmd: &LinkCommand,
419    ) -> (Vec<String>, Vec<Absent>) {
420        let sent = self.broadcast(username, cmd.clone());
421        let queued = outbox::queue_for_absent(username, &self.live(), cmd);
422        (sent, queued)
423    }
424
425    /// Each linked device, as `(device, username)`.
426    fn live(&self) -> Vec<(String, String)> {
427        self.entries
428            .lock()
429            .iter()
430            .map(|e| (e.device.clone(), e.info.username.clone()))
431            .collect()
432    }
433}
434
435impl Absent {
436    fn key(&self) -> (String, String) {
437        (self.device.clone(), self.username.clone())
438    }
439}
440
441/// Wake these absent devices, each a `(device, username)`, where a push can:
442/// each links, and takes what was queued for it.
443fn wake(devices: &[(String, String)]) {
444    let Some(pusher) = crate::push::pusher() else {
445        return;
446    };
447    let targets: Vec<outbox::PushTarget> = outbox::push_targets(None)
448        .into_iter()
449        .filter(|t| {
450            devices
451                .iter()
452                .any(|(device, username)| *device == t.device && *username == t.username)
453        })
454        .collect();
455    if targets.is_empty() {
456        return;
457    }
458    std::thread::spawn(move || {
459        for t in targets {
460            deliver_push(pusher, &t, &crate::push::Push::Wake);
461        }
462    });
463}
464
465/// Reach a device that is not linked.
466///
467/// Music is sent as a notification to tap, at once: iOS does not let an app it
468/// woke start audio, so waking it for that only delays the notification. Every
469/// other command goes into the device's outbox with a background push to wake
470/// it, and runs as if it had been linked. `None` without a push key, or with no
471/// device in scope that has given a push token.
472fn reach_absent(
473    username: Option<&str>,
474    id: Option<&str>,
475    cmd: &LinkCommand,
476) -> Option<Result<ClientInfo, String>> {
477    let pusher = crate::push::pusher()?;
478    let targets: Vec<outbox::PushTarget> = outbox::push_targets(username)
479        .into_iter()
480        .filter(|t| id.is_none_or(|id| t.device == id || t.name.eq_ignore_ascii_case(id)))
481        .collect();
482    let target = match targets.as_slice() {
483        [] => return None,
484        [only] => only.clone(),
485        several => {
486            return Some(Err(format!(
487                "no koan app is linked, and several can be reached: {}. Ask which, then pass `client`",
488                several
489                    .iter()
490                    .map(|t| t.name.as_str())
491                    .collect::<Vec<_>>()
492                    .join(", ")
493            )));
494        }
495    };
496    let info = ClientInfo {
497        id: target.device.clone(),
498        name: target.name.clone(),
499        platform: target.platform.clone(),
500        username: target.username.clone(),
501        connected_at: 0,
502        state: LinkState::default(),
503        last_played_at: None,
504        state_at: 0,
505        reports: false,
506        notified: true,
507    };
508    let verb = match cmd {
509        LinkCommand::Play { .. } | LinkCommand::JumpTo { .. } | LinkCommand::Resume => Some("Play"),
510        _ => None,
511    };
512    let push = match verb {
513        Some(verb) => crate::push::Push::Notify {
514            title: format!("{verb} on {}", target.name),
515            body: outbox::describe(cmd).unwrap_or_else(|| "From your koan server".into()),
516            command: serde_json::to_value(cmd).ok()?,
517            image: cover_track(cmd).and_then(|t| pusher.cover_link(t)),
518        },
519        None => {
520            outbox::queue_for(&target.device, &target.username, cmd);
521            crate::push::Push::Wake
522        }
523    };
524    std::thread::spawn(move || deliver_push(pusher, &target, &push));
525    Some(Ok(info))
526}
527
528/// The track whose album cover a notification for `cmd` shows.
529fn cover_track(cmd: &LinkCommand) -> Option<i64> {
530    match cmd {
531        LinkCommand::Play {
532            track_ids,
533            start_at,
534        } => track_ids
535            .get(*start_at as usize)
536            .or(track_ids.first())?
537            .parse()
538            .ok(),
539        LinkCommand::JumpTo { track_id } => track_id.parse().ok(),
540        _ => None,
541    }
542}
543
544/// Send one push, forgetting a token Apple says is no longer good.
545fn deliver_push(
546    pusher: &crate::push::Pusher,
547    target: &outbox::PushTarget,
548    push: &crate::push::Push,
549) {
550    use crate::push::Outcome;
551    match pusher.send(&target.token, target.sandbox, push) {
552        Outcome::Sent => log::info!("push: sent to {}", target.name),
553        Outcome::Gone => {
554            log::info!(
555                "push: {}'s token is no longer valid; forgotten",
556                target.name
557            );
558            outbox::forget_push(&target.username, &target.device);
559        }
560        Outcome::Failed(e) => log::warn!("push: to {} failed: {e}", target.name),
561    }
562}
563
564/// Have every device pull what the server just changed (a playlist edited,
565/// albums added): at once where linked, on next link where not. Syncs waiting
566/// for a device collapse into one, and so do the pushes that wake it: see
567/// `Wakes`.
568pub fn changed() {
569    let (_, queued) = registry().link_or_queue(None, &LinkCommand::Sync { full: false });
570    if queued.is_empty() || crate::push::pusher().is_none() {
571        return;
572    }
573    let mut wakes = WAKES.lock();
574    wakes.add(queued.iter().map(Absent::key), std::time::Instant::now());
575    if !wakes.timer {
576        wakes.timer = true;
577        std::thread::spawn(send_wakes);
578    }
579}
580
581/// How long the library has to stay still before a suspended device is woken
582/// to sync. A download of several albums scans after each one; iOS rations
583/// background pushes, and the device needs waking once, at the end.
584const QUIET: std::time::Duration = std::time::Duration::from_secs(30);
585
586/// However busy the library stays, a device waits no longer than this.
587const LONGEST_WAIT: std::time::Duration = std::time::Duration::from_secs(5 * 60);
588
589static WAKES: Mutex<Wakes> = Mutex::new(Wakes {
590    pending: Vec::new(),
591    timer: false,
592});
593
594/// Background pushes held back until the library is quiet: at most one
595/// pending per device.
596struct Wakes {
597    /// `(device, username)`, when first asked for, when last asked for.
598    pending: Vec<((String, String), std::time::Instant, std::time::Instant)>,
599    /// Whether a thread is waiting to send them.
600    timer: bool,
601}
602
603impl Wakes {
604    fn add(
605        &mut self,
606        devices: impl IntoIterator<Item = (String, String)>,
607        now: std::time::Instant,
608    ) {
609        for key in devices {
610            match self.pending.iter_mut().find(|(k, _, _)| *k == key) {
611                Some((_, _, last)) => *last = now,
612                None => self.pending.push((key, now, now)),
613            }
614        }
615    }
616
617    fn due_at(first: std::time::Instant, last: std::time::Instant) -> std::time::Instant {
618        (last + QUIET).min(first + LONGEST_WAIT)
619    }
620
621    /// The earliest a pending wake is due.
622    fn next(&self) -> Option<std::time::Instant> {
623        self.pending
624            .iter()
625            .map(|(_, first, last)| Self::due_at(*first, *last))
626            .min()
627    }
628
629    /// Take the wakes due by `now`.
630    fn take_due(&mut self, now: std::time::Instant) -> Vec<(String, String)> {
631        let (due, waiting) = std::mem::take(&mut self.pending)
632            .into_iter()
633            .partition(|(_, first, last)| Self::due_at(*first, *last) <= now);
634        self.pending = waiting;
635        due.into_iter().map(|(key, _, _)| key).collect()
636    }
637}
638
639/// Send each wake once it is due, until none is pending. A device that has
640/// linked meanwhile took its sync down the link and is not pushed.
641fn send_wakes() {
642    loop {
643        let (due, next) = {
644            let mut wakes = WAKES.lock();
645            let due = wakes.take_due(std::time::Instant::now());
646            let next = wakes.next();
647            if due.is_empty() && next.is_none() {
648                wakes.timer = false;
649                return;
650            }
651            (due, next)
652        };
653        let live = registry().live();
654        let absent: Vec<(String, String)> = due.into_iter().filter(|d| !live.contains(d)).collect();
655        if !absent.is_empty() {
656            wake(&absent);
657        }
658        if let Some(next) = next {
659            std::thread::sleep(next.saturating_duration_since(std::time::Instant::now()));
660        }
661    }
662}
663
664/// After a library scan or sync: if the library holds different tracks or
665/// albums from when this was last asked, tell every device. A scan that found
666/// nothing new, which is most of them, sends nothing.
667pub fn changed_if_library_moved(conn: &rusqlite::Connection) {
668    static LAST: parking_lot::Mutex<Option<(i64, i64, i64)>> = parking_lot::Mutex::new(None);
669    let Ok(now) = conn.query_row(
670        "SELECT (SELECT COUNT(*) FROM tracks), (SELECT COALESCE(MAX(id), 0) FROM tracks),
671                (SELECT COUNT(*) FROM albums)",
672        [],
673        |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
674    ) else {
675        return;
676    };
677    let before = LAST.lock().replace(now);
678    if before.is_some_and(|b| b != now) {
679        changed();
680    }
681}
682
683/// The server-side record of devices and their waiting commands, in the
684/// library database so it outlives a restart.
685mod outbox {
686    use koan_core::remote::link::LinkCommand;
687
688    /// Dropped undelivered after this long: a device away a month re-syncs
689    /// on its own when opened.
690    const KEEP_SECS: i64 = 30 * 24 * 60 * 60;
691
692    /// Tests keep to memory: the configured database is whoever ran them.
693    fn db() -> Option<koan_core::db::connection::Database> {
694        if cfg!(test) {
695            return None;
696        }
697        koan_core::db::connection::Database::open(&koan_core::config::db_path()).ok()
698    }
699
700    pub fn load_orders() -> Vec<super::Order> {
701        let Some(db) = db() else { return Vec::new() };
702        db.conn
703            .prepare("SELECT body FROM link_orders ORDER BY created_at")
704            .and_then(|mut s| {
705                s.query_map([], |r| r.get::<_, String>(0))?
706                    .collect::<Result<Vec<_>, _>>()
707            })
708            .unwrap_or_default()
709            .into_iter()
710            .filter_map(|b| serde_json::from_str(&b).ok())
711            .collect()
712    }
713
714    pub fn save_order(order: &super::Order) {
715        let (Some(db), Ok(body)) = (db(), serde_json::to_string(order)) else {
716            return;
717        };
718        let _ = db.conn.execute(
719            "INSERT OR REPLACE INTO link_orders (id, body, created_at) VALUES (?1, ?2, ?3)",
720            rusqlite::params![order.id, body, order.created_at],
721        );
722    }
723
724    pub fn drop_order(id: &str) {
725        if let Some(db) = db() {
726            let _ = db
727                .conn
728                .execute("DELETE FROM link_orders WHERE id = ?1", [id]);
729        }
730    }
731
732    pub fn take_and_remember(
733        username: &str,
734        device: &str,
735        name: &str,
736        platform: &str,
737    ) -> Vec<LinkCommand> {
738        let Some(db) = db() else { return Vec::new() };
739        let now = chrono::Utc::now().timestamp();
740        let _ = db.conn.execute(
741            "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES (?1, ?2, ?3, ?4, ?5)
742             ON CONFLICT (device, username) DO UPDATE SET name = ?3, platform = ?4, last_seen = ?5",
743            rusqlite::params![device, username, name, platform, now],
744        );
745        let _ = db.conn.execute(
746            "DELETE FROM link_outbox WHERE created_at < ?1",
747            [now - KEEP_SECS],
748        );
749        let waiting: Vec<(i64, String)> = db
750            .conn
751            .prepare("SELECT id, command FROM link_outbox WHERE device = ?1 AND username = ?2 ORDER BY id")
752            .and_then(|mut s| {
753                s.query_map([device, username], |r| Ok((r.get(0)?, r.get(1)?)))?
754                    .collect()
755            })
756            .unwrap_or_default();
757        let _ = db.conn.execute(
758            "DELETE FROM link_outbox WHERE device = ?1 AND username = ?2",
759            [device, username],
760        );
761        if !waiting.is_empty() {
762            log::info!("link: {} waiting commands for {name}", waiting.len());
763        }
764        waiting
765            .into_iter()
766            .filter_map(|(_, c)| serde_json::from_str(&c).ok())
767            .collect()
768    }
769
770    /// A known device that was not linked when something was queued for it.
771    pub struct Absent {
772        pub device: String,
773        pub username: String,
774        pub name: String,
775    }
776
777    /// Queue `cmd` for each known device in scope that is not in `live`.
778    pub fn queue_for_absent(
779        username: Option<&str>,
780        live: &[(String, String)],
781        cmd: &LinkCommand,
782    ) -> Vec<Absent> {
783        let Some(db) = db() else { return Vec::new() };
784        let known: Vec<(String, String, String)> = db
785            .conn
786            .prepare("SELECT device, username, name FROM link_devices")
787            .and_then(|mut s| {
788                s.query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))?
789                    .collect()
790            })
791            .unwrap_or_default();
792        let Ok(text) = serde_json::to_string(cmd) else {
793            return Vec::new();
794        };
795        let is_sync = matches!(cmd, LinkCommand::Sync { .. });
796        let now = chrono::Utc::now().timestamp();
797        let mut queued = Vec::new();
798        for (device, user, name) in known {
799            if username.is_some_and(|u| u != user)
800                || live.iter().any(|(d, u)| *d == device && *u == user)
801            {
802                continue;
803            }
804            if is_sync {
805                // One pending sync is enough; a full one covers an incremental.
806                let _ = db.conn.execute(
807                    "DELETE FROM link_outbox WHERE device = ?1 AND username = ?2 AND command LIKE '{\"type\":\"sync\"%'",
808                    [&device, &user],
809                );
810            }
811            if db
812                .conn
813                .execute(
814                    "INSERT INTO link_outbox (device, username, command, created_at) VALUES (?1, ?2, ?3, ?4)",
815                    rusqlite::params![device, user, text, now],
816                )
817                .is_ok()
818            {
819                queued.push(Absent {
820                    device,
821                    username: user,
822                    name,
823                });
824            }
825        }
826        queued
827    }
828
829    /// A device Apple's push service can reach.
830    #[derive(Clone)]
831    pub struct PushTarget {
832        pub device: String,
833        pub username: String,
834        pub name: String,
835        pub platform: String,
836        pub token: String,
837        pub sandbox: bool,
838    }
839
840    /// Queue `cmd` for one device, to go down its next link.
841    pub fn queue_for(device: &str, username: &str, cmd: &LinkCommand) {
842        let (Some(db), Ok(text)) = (db(), serde_json::to_string(cmd)) else {
843            return;
844        };
845        let _ = db.conn.execute(
846            "INSERT INTO link_outbox (device, username, command, created_at) VALUES (?1, ?2, ?3, ?4)",
847            rusqlite::params![device, username, text, chrono::Utc::now().timestamp()],
848        );
849    }
850
851    pub fn save_push(username: &str, device: &str, token: &str, sandbox: bool) {
852        let Some(db) = db() else { return };
853        let _ = db.conn.execute(
854            "INSERT INTO link_push (device, username, token, sandbox, updated_at) VALUES (?1, ?2, ?3, ?4, ?5)
855             ON CONFLICT (device, username) DO UPDATE SET token = ?3, sandbox = ?4, updated_at = ?5",
856            rusqlite::params![device, username, token, sandbox, chrono::Utc::now().timestamp()],
857        );
858    }
859
860    pub fn forget_push(username: &str, device: &str) {
861        if let Some(db) = db() {
862            let _ = db.conn.execute(
863                "DELETE FROM link_push WHERE device = ?1 AND username = ?2",
864                [device, username],
865            );
866        }
867    }
868
869    /// Devices in scope with a push token, most recently seen first.
870    pub fn push_targets(username: Option<&str>) -> Vec<PushTarget> {
871        let Some(db) = db() else { return Vec::new() };
872        db.conn
873            .prepare(
874                "SELECT p.device, p.username, d.name, d.platform, p.token, p.sandbox
875                   FROM link_push p JOIN link_devices d ON d.device = p.device AND d.username = p.username
876                  WHERE ?1 IS NULL OR p.username = ?1
877                  ORDER BY d.last_seen DESC",
878            )
879            .and_then(|mut s| {
880                s.query_map([username], |r| {
881                    Ok(PushTarget {
882                        device: r.get(0)?,
883                        username: r.get(1)?,
884                        name: r.get(2)?,
885                        platform: r.get(3)?,
886                        token: r.get(4)?,
887                        sandbox: r.get(5)?,
888                    })
889                })?
890                .collect()
891            })
892            .unwrap_or_default()
893    }
894
895    /// What a playback command would play, for a notification to say:
896    /// "Golden Standard — Tony Petersen", or a track and how many follow.
897    pub fn describe(cmd: &LinkCommand) -> Option<String> {
898        let ids: Vec<i64> = match cmd {
899            LinkCommand::Play { track_ids, .. }
900            | LinkCommand::Enqueue { track_ids }
901            | LinkCommand::PlayNext { track_ids } => {
902                track_ids.iter().filter_map(|t| t.parse().ok()).collect()
903            }
904            LinkCommand::JumpTo { track_id } => vec![track_id.parse().ok()?],
905            _ => return None,
906        };
907        let db = db()?;
908        let row = |id: i64| {
909            db.conn
910                .query_row(
911                    "SELECT t.title, COALESCE(a.name, ''), COALESCE(al.title, ''), t.album_id
912                       FROM tracks t LEFT JOIN artists a ON a.id = t.artist_id
913                       LEFT JOIN albums al ON al.id = t.album_id WHERE t.id = ?1",
914                    [id],
915                    |r| {
916                        Ok((
917                            r.get::<_, String>(0)?,
918                            r.get::<_, String>(1)?,
919                            r.get::<_, String>(2)?,
920                            r.get::<_, Option<i64>>(3)?,
921                        ))
922                    },
923                )
924                .ok()
925        };
926        let (title, artist, album, album_id) = row(*ids.first()?)?;
927        let one_album = ids.len() > 1
928            && album_id.is_some()
929            && ids
930                .iter()
931                .all(|id| row(*id).is_some_and(|r| r.3 == album_id));
932        Some(match (one_album, ids.len()) {
933            (true, _) => format!("{album} — {artist}"),
934            (false, 1) => format!("{title} — {artist}"),
935            (false, n) => format!("{title} — {artist}, and {} more", n - 1),
936        })
937    }
938}
939
940/// How long ago a client can have stopped playing and still be the obvious
941/// one to send music to.
942const RECENT: i64 = 6 * 60 * 60;
943
944fn pick(clients: &[ClientInfo], now: i64) -> Result<&ClientInfo, String> {
945    if let Some(c) = clients.iter().find(|c| c.state.playing) {
946        return Ok(c);
947    }
948    if let Some(c) = clients
949        .iter()
950        .filter(|c| c.last_played_at.is_some_and(|t| now - t < RECENT))
951        .max_by_key(|c| c.last_played_at)
952    {
953        return Ok(c);
954    }
955    match clients {
956        [] => Err("no koan app is linked to this server; open koan on the device".into()),
957        [only] => Ok(only),
958        several => Err(format!(
959            "several koan apps are linked and none has played recently: {}. Ask which, then pass `client`",
960            several
961                .iter()
962                .map(|c| c.name.as_str())
963                .collect::<Vec<_>>()
964                .join(", ")
965        )),
966    }
967}
968
969#[cfg(test)]
970mod tests {
971
972    #[test]
973    fn a_playlist_order_adds_the_named_tracks_once_they_arrive() {
974        let dir = tempfile::tempdir().unwrap();
975        let path = dir.path().join("koan.db");
976        let db = koan_core::db::connection::Database::open(&path).unwrap();
977        let playlist = koan_core::db::queries::create_playlist(
978            &db.conn,
979            koan_core::db::queries::LOCAL_USER,
980            "cyberpunk",
981            None,
982        )
983        .unwrap();
984        let order = Order {
985            id: "o1".into(),
986            username: None,
987            client: None,
988            artist: "Perturbator".into(),
989            album: "Dangerous Days".into(),
990            play_next: false,
991            playlist: Some(playlist),
992            titles: vec!["Future Club".into()],
993            created_at: chrono::Utc::now().timestamp(),
994        };
995        registry().add_order(order);
996
997        // Not in the library yet: nothing happens, and the order waits.
998        fulfil_from(&path);
999        assert!(registry().orders(None).iter().any(|o| o.id == "o1"));
1000
1001        db.conn
1002            .execute_batch(
1003                "INSERT INTO artists (id, name) VALUES (1, 'Perturbator');
1004                 INSERT INTO albums (id, title, artist_id) VALUES (1, 'Dangerous Days', 1);
1005                 INSERT INTO tracks (id, title, album_id, artist_id, track_number, path) VALUES
1006                   (1, 'Welcome Back', 1, 1, 1, '/1.flac'), (2, 'Future Club', 1, 1, 2, '/2.flac');",
1007            )
1008            .unwrap();
1009        fulfil_from(&path);
1010        let held: Vec<i64> = db
1011            .conn
1012            .prepare("SELECT track_id FROM playlist_tracks WHERE playlist_id = ?1")
1013            .unwrap()
1014            .query_map([playlist], |r| r.get(0))
1015            .unwrap()
1016            .collect::<Result<_, _>>()
1017            .unwrap();
1018        assert_eq!(held, [2]);
1019        assert!(!registry().orders(None).iter().any(|o| o.id == "o1"));
1020    }
1021
1022    use super::*;
1023
1024    #[test]
1025    fn a_reconnect_replaces_the_device_and_commands_reach_it() {
1026        let reg = Registry::default();
1027        let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
1028        let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
1029        let (tx3, _rx3) = tokio::sync::mpsc::unbounded_channel();
1030        reg.register("j", "phone", "ios", "dev-1", tx1);
1031        let id = reg.register("j", "phone", "ios", "dev-1", tx2);
1032        reg.register("someone", "laptop", "macos", "dev-2", tx3);
1033
1034        assert_eq!(reg.list(Some("j")).len(), 1);
1035        assert_eq!(reg.list(None).len(), 2);
1036
1037        let sent = reg.send(Some("j"), None, LinkCommand::Pause).unwrap();
1038        assert_eq!(sent.id, id);
1039        assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
1040
1041        // Another account's device is not this account's to command.
1042        assert!(
1043            reg.send(Some("j"), Some("laptop"), LinkCommand::Pause)
1044                .is_err()
1045        );
1046
1047        reg.unregister(&id);
1048        assert!(reg.send(Some("j"), None, LinkCommand::Pause).is_err());
1049    }
1050
1051    #[test]
1052    fn a_burst_of_changes_wakes_each_device_once_when_it_goes_quiet() {
1053        let t0 = std::time::Instant::now();
1054        let s = std::time::Duration::from_secs;
1055        let phone = || ("dev-1".to_string(), "j".to_string());
1056        let ipad = || ("dev-2".to_string(), "j".to_string());
1057        let mut wakes = Wakes {
1058            pending: Vec::new(),
1059            timer: false,
1060        };
1061
1062        wakes.add([phone()], t0);
1063        wakes.add([phone(), ipad()], t0 + s(10));
1064        wakes.add([phone()], t0 + s(20));
1065        assert_eq!(wakes.pending.len(), 2);
1066
1067        // Quiet is measured from each device's last change.
1068        assert!(wakes.take_due(t0 + s(39)).is_empty());
1069        assert_eq!(wakes.take_due(t0 + s(40)), [ipad()]);
1070        assert_eq!(wakes.next(), Some(t0 + s(50)));
1071        assert_eq!(wakes.take_due(t0 + s(50)), [phone()]);
1072        assert_eq!(wakes.next(), None);
1073    }
1074
1075    #[test]
1076    fn a_library_that_never_goes_quiet_still_wakes_devices() {
1077        let t0 = std::time::Instant::now();
1078        let phone = || ("dev-1".to_string(), "j".to_string());
1079        let mut wakes = Wakes {
1080            pending: Vec::new(),
1081            timer: false,
1082        };
1083        let mut sent = 0;
1084        for i in 0..40 {
1085            let now = t0 + std::time::Duration::from_secs(i * 10);
1086            sent += wakes.take_due(now).len();
1087            wakes.add([phone()], now);
1088        }
1089        // Six and a half minutes of changes every ten seconds: one push at
1090        // the five-minute mark, and one pending.
1091        assert_eq!(sent, 1);
1092        assert_eq!(wakes.pending.len(), 1);
1093    }
1094
1095    #[test]
1096    fn the_device_playing_is_the_one_meant() {
1097        let reg = Registry::default();
1098        let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
1099        let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
1100        let mac = reg.register("j", "mac", "macos", "dev-1", tx1);
1101        let phone = reg.register("j", "phone", "ios", "dev-2", tx2);
1102
1103        // Two idle devices: no telling, so the caller is told to ask.
1104        let err = reg.send(Some("j"), None, LinkCommand::Pause).unwrap_err();
1105        assert!(err.contains("mac") && err.contains("phone"), "{err}");
1106
1107        reg.report(
1108            &phone,
1109            LinkState {
1110                playing: true,
1111                ..Default::default()
1112            },
1113        );
1114        assert_eq!(
1115            reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
1116            phone
1117        );
1118        assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
1119
1120        // Stopped a moment ago: still the one meant, over the Mac.
1121        reg.report(&phone, LinkState::default());
1122        assert_eq!(
1123            reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
1124            phone
1125        );
1126        let _ = mac;
1127    }
1128}