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}
34
35impl ClientInfo {
36    /// Where the playhead is now, from where it was reported to be.
37    pub fn position_ms(&self) -> u64 {
38        let pos = self.state.position_ms;
39        if !self.state.playing {
40            return pos;
41        }
42        let run = (chrono::Utc::now().timestamp_millis() - self.state_at).max(0) as u64;
43        let pos = pos + run;
44        if self.state.duration_ms > 0 {
45            pos.min(self.state.duration_ms)
46        } else {
47            pos
48        }
49    }
50}
51
52struct Entry {
53    info: ClientInfo,
54    /// The client's own id for itself, so a reconnect replaces its entry.
55    device: String,
56    tx: UnboundedSender<LinkCommand>,
57}
58
59/// "When this album is in the library, queue it on my device": a request made
60/// before the album exists, fulfilled by the scan that finds it.
61#[derive(Debug, Clone)]
62pub struct Order {
63    pub id: String,
64    /// Whose devices it may go to.
65    pub username: Option<String>,
66    /// A client id or name; `None` for whichever `send` would pick then.
67    pub client: Option<String>,
68    pub artist: String,
69    pub album: String,
70    /// Insert after the current track rather than at the end.
71    pub play_next: bool,
72    /// Unix seconds.
73    pub created_at: i64,
74}
75
76/// An order nobody's scan has fulfilled in this long is dropped.
77const ORDER_TTL: i64 = 24 * 60 * 60;
78
79#[derive(Default)]
80pub struct Registry {
81    entries: Mutex<Vec<Entry>>,
82    orders: Mutex<Vec<Order>>,
83}
84
85/// One registry per process: the WebSocket route and the GraphQL schema are
86/// built in different places and both need it.
87pub fn registry() -> &'static Registry {
88    static REGISTRY: LazyLock<Registry> = LazyLock::new(Registry::default);
89    &REGISTRY
90}
91
92impl Registry {
93    /// Add a client, replacing any earlier connection from the same device
94    /// and account. Returns its id.
95    pub fn register(
96        &self,
97        username: &str,
98        name: &str,
99        platform: &str,
100        device: &str,
101        tx: UnboundedSender<LinkCommand>,
102    ) -> String {
103        let id = uuid::Uuid::now_v7().to_string();
104        let mut entries = self.entries.lock();
105        entries.retain(|e| !(e.device == device && e.info.username == username));
106        entries.push(Entry {
107            info: ClientInfo {
108                id: id.clone(),
109                name: name.to_string(),
110                platform: platform.to_string(),
111                username: username.to_string(),
112                connected_at: chrono::Utc::now().timestamp(),
113                state: LinkState::default(),
114                last_played_at: None,
115                state_at: chrono::Utc::now().timestamp_millis(),
116                reports: false,
117            },
118            device: device.to_string(),
119            tx,
120        });
121        id
122    }
123
124    /// Record what a client says it is doing.
125    pub fn report(&self, id: &str, state: LinkState) {
126        let mut entries = self.entries.lock();
127        if let Some(e) = entries.iter_mut().find(|e| e.info.id == id) {
128            if state.playing || e.info.state.playing {
129                e.info.last_played_at = Some(chrono::Utc::now().timestamp());
130            }
131            e.info.state = state;
132            e.info.state_at = chrono::Utc::now().timestamp_millis();
133            e.info.reports = true;
134        }
135    }
136
137    pub fn unregister(&self, id: &str) {
138        self.entries.lock().retain(|e| e.info.id != id);
139    }
140
141    /// Clients `username` may command, newest first; every client for `None`.
142    pub fn list(&self, username: Option<&str>) -> Vec<ClientInfo> {
143        let mut out: Vec<ClientInfo> = self
144            .entries
145            .lock()
146            .iter()
147            .filter(|e| username.is_none_or(|u| e.info.username == u))
148            .map(|e| e.info.clone())
149            .collect();
150        out.sort_by_key(|c| std::cmp::Reverse(c.connected_at));
151        out
152    }
153
154    /// Send to `id` (an id or a name), or with none to the client the
155    /// command most likely means: the one playing, else the one that played
156    /// within `RECENT`, else the only one linked. `Err` names the choices when
157    /// there is no telling, so whoever asked can ask the person.
158    pub fn send(
159        &self,
160        username: Option<&str>,
161        id: Option<&str>,
162        cmd: LinkCommand,
163    ) -> Result<ClientInfo, String> {
164        let clients = self.list(username);
165        let target = match id {
166            Some(id) => clients
167                .iter()
168                .find(|c| c.id == id || c.name.eq_ignore_ascii_case(id))
169                .ok_or_else(|| format!("no linked client {id}; see `clients`"))?,
170            None => pick(&clients, chrono::Utc::now().timestamp())?,
171        };
172        let entries = self.entries.lock();
173        let entry = entries
174            .iter()
175            .find(|e| e.info.id == target.id)
176            .ok_or("that client has just gone")?;
177        entry
178            .tx
179            .send(cmd)
180            .map_err(|_| "that client has just gone".to_string())?;
181        Ok(target.clone())
182    }
183}
184
185impl Registry {
186    pub fn add_order(&self, order: Order) {
187        self.orders.lock().push(order);
188    }
189
190    pub fn orders(&self, username: Option<&str>) -> Vec<Order> {
191        self.orders
192            .lock()
193            .iter()
194            .filter(|o| username.is_none() || o.username.as_deref() == username)
195            .cloned()
196            .collect()
197    }
198
199    pub fn cancel_order(&self, username: Option<&str>, id: &str) -> bool {
200        let mut orders = self.orders.lock();
201        let before = orders.len();
202        orders
203            .retain(|o| !(o.id == id && (username.is_none() || o.username.as_deref() == username)));
204        orders.len() != before
205    }
206
207    /// Send every order whose album the library now holds, and drop it.
208    /// `find` answers an order with the album's track ids, in order.
209    pub fn fulfil_orders(&self, find: impl Fn(&Order) -> Option<Vec<i64>>) {
210        let now = chrono::Utc::now().timestamp();
211        let pending: Vec<Order> = {
212            let mut orders = self.orders.lock();
213            orders.retain(|o| now - o.created_at < ORDER_TTL);
214            orders.clone()
215        };
216        for order in pending {
217            let Some(ids) = find(&order).filter(|ids| !ids.is_empty()) else {
218                continue;
219            };
220            let track_ids = ids.iter().map(i64::to_string).collect();
221            let cmd = if order.play_next {
222                LinkCommand::PlayNext { track_ids }
223            } else {
224                LinkCommand::Enqueue { track_ids }
225            };
226            match self.send(order.username.as_deref(), order.client.as_deref(), cmd) {
227                Ok(c) => {
228                    log::info!(
229                        "link: {} — {} arrived; queued on {}",
230                        order.artist,
231                        order.album,
232                        c.name
233                    );
234                    self.orders.lock().retain(|o| o.id != order.id);
235                }
236                // No device to send to yet: kept, and tried after the next scan.
237                Err(e) => log::info!("link: {} — {} arrived but {e}", order.artist, order.album),
238            }
239        }
240    }
241}
242
243/// Fulfil standing orders against the library at `db_path`.
244pub fn fulfil_from(db_path: &std::path::Path) {
245    let registry = registry();
246    if registry.orders.lock().is_empty() {
247        return;
248    }
249    let Ok(db) = koan_core::db::connection::Database::open(db_path) else {
250        return;
251    };
252    registry.fulfil_orders(|order| album_tracks(&db.conn, &order.artist, &order.album));
253}
254
255/// The newest album whose artist and title contain these, as track ids in
256/// disc and track order.
257pub fn album_tracks(conn: &rusqlite::Connection, artist: &str, album: &str) -> Option<Vec<i64>> {
258    let like = |s: &str| format!("%{}%", s.replace(['%', '_'], ""));
259    let album_id: i64 = conn
260        .query_row(
261            "SELECT al.id FROM albums al JOIN artists a ON a.id = al.artist_id
262              WHERE a.name LIKE ?1 COLLATE NOCASE AND al.title LIKE ?2 COLLATE NOCASE
263              ORDER BY al.id DESC LIMIT 1",
264            [like(artist), like(album)],
265            |r| r.get(0),
266        )
267        .ok()?;
268    let mut stmt = conn
269        .prepare("SELECT id FROM tracks WHERE album_id = ?1 ORDER BY disc, track_number, id")
270        .ok()?;
271    let ids = stmt
272        .query_map([album_id], |r| r.get(0))
273        .ok()?
274        .filter_map(Result::ok)
275        .collect();
276    Some(ids)
277}
278
279/// How long ago a client can have stopped playing and still be the obvious
280/// one to send music to.
281const RECENT: i64 = 6 * 60 * 60;
282
283fn pick(clients: &[ClientInfo], now: i64) -> Result<&ClientInfo, String> {
284    if let Some(c) = clients.iter().find(|c| c.state.playing) {
285        return Ok(c);
286    }
287    if let Some(c) = clients
288        .iter()
289        .filter(|c| c.last_played_at.is_some_and(|t| now - t < RECENT))
290        .max_by_key(|c| c.last_played_at)
291    {
292        return Ok(c);
293    }
294    match clients {
295        [] => Err("no koan app is linked to this server; open koan on the device".into()),
296        [only] => Ok(only),
297        several => Err(format!(
298            "several koan apps are linked and none has played recently: {}. Ask which, then pass `client`",
299            several
300                .iter()
301                .map(|c| c.name.as_str())
302                .collect::<Vec<_>>()
303                .join(", ")
304        )),
305    }
306}
307
308#[cfg(test)]
309mod tests {
310    use super::*;
311
312    #[test]
313    fn a_reconnect_replaces_the_device_and_commands_reach_it() {
314        let reg = Registry::default();
315        let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
316        let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
317        let (tx3, _rx3) = tokio::sync::mpsc::unbounded_channel();
318        reg.register("j", "phone", "ios", "dev-1", tx1);
319        let id = reg.register("j", "phone", "ios", "dev-1", tx2);
320        reg.register("someone", "laptop", "macos", "dev-2", tx3);
321
322        assert_eq!(reg.list(Some("j")).len(), 1);
323        assert_eq!(reg.list(None).len(), 2);
324
325        let sent = reg.send(Some("j"), None, LinkCommand::Pause).unwrap();
326        assert_eq!(sent.id, id);
327        assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
328
329        // Another account's device is not this account's to command.
330        assert!(
331            reg.send(Some("j"), Some("laptop"), LinkCommand::Pause)
332                .is_err()
333        );
334
335        reg.unregister(&id);
336        assert!(reg.send(Some("j"), None, LinkCommand::Pause).is_err());
337    }
338
339    #[test]
340    fn the_device_playing_is_the_one_meant() {
341        let reg = Registry::default();
342        let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
343        let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
344        let mac = reg.register("j", "mac", "macos", "dev-1", tx1);
345        let phone = reg.register("j", "phone", "ios", "dev-2", tx2);
346
347        // Two idle devices: no telling, so the caller is told to ask.
348        let err = reg.send(Some("j"), None, LinkCommand::Pause).unwrap_err();
349        assert!(err.contains("mac") && err.contains("phone"), "{err}");
350
351        reg.report(
352            &phone,
353            LinkState {
354                playing: true,
355                ..Default::default()
356            },
357        );
358        assert_eq!(
359            reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
360            phone
361        );
362        assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
363
364        // Stopped a moment ago: still the one meant, over the Mac.
365        reg.report(&phone, LinkState::default());
366        assert_eq!(
367            reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
368            phone
369        );
370        let _ = mac;
371    }
372}