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;
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}
25
26struct Entry {
27    info: ClientInfo,
28    /// The client's own id for itself, so a reconnect replaces its entry.
29    device: String,
30    tx: UnboundedSender<LinkCommand>,
31}
32
33#[derive(Default)]
34pub struct Registry {
35    entries: Mutex<Vec<Entry>>,
36}
37
38/// One registry per process: the WebSocket route and the GraphQL schema are
39/// built in different places and both need it.
40pub fn registry() -> &'static Registry {
41    static REGISTRY: LazyLock<Registry> = LazyLock::new(Registry::default);
42    &REGISTRY
43}
44
45impl Registry {
46    /// Add a client, replacing any earlier connection from the same device
47    /// and account. Returns its id.
48    pub fn register(
49        &self,
50        username: &str,
51        name: &str,
52        platform: &str,
53        device: &str,
54        tx: UnboundedSender<LinkCommand>,
55    ) -> String {
56        let id = uuid::Uuid::now_v7().to_string();
57        let mut entries = self.entries.lock();
58        entries.retain(|e| !(e.device == device && e.info.username == username));
59        entries.push(Entry {
60            info: ClientInfo {
61                id: id.clone(),
62                name: name.to_string(),
63                platform: platform.to_string(),
64                username: username.to_string(),
65                connected_at: chrono::Utc::now().timestamp(),
66            },
67            device: device.to_string(),
68            tx,
69        });
70        id
71    }
72
73    pub fn unregister(&self, id: &str) {
74        self.entries.lock().retain(|e| e.info.id != id);
75    }
76
77    /// Clients `username` may command, newest first; every client for `None`.
78    pub fn list(&self, username: Option<&str>) -> Vec<ClientInfo> {
79        let mut out: Vec<ClientInfo> = self
80            .entries
81            .lock()
82            .iter()
83            .filter(|e| username.is_none_or(|u| e.info.username == u))
84            .map(|e| e.info.clone())
85            .collect();
86        out.sort_by_key(|c| std::cmp::Reverse(c.connected_at));
87        out
88    }
89
90    /// Send to `id`, or with no id to the newest client `username` may
91    /// command. The client that received it, or `None` when there was none to
92    /// send to.
93    pub fn send(
94        &self,
95        username: Option<&str>,
96        id: Option<&str>,
97        cmd: LinkCommand,
98    ) -> Option<ClientInfo> {
99        let target = match id {
100            Some(id) => self
101                .list(username)
102                .into_iter()
103                .find(|c| c.id == id || c.name.eq_ignore_ascii_case(id))?,
104            None => self.list(username).into_iter().next()?,
105        };
106        let entries = self.entries.lock();
107        let entry = entries.iter().find(|e| e.info.id == target.id)?;
108        entry.tx.send(cmd).ok()?;
109        Some(target)
110    }
111}
112
113#[cfg(test)]
114mod tests {
115    use super::*;
116
117    #[test]
118    fn a_reconnect_replaces_the_device_and_commands_reach_it() {
119        let reg = Registry::default();
120        let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
121        let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
122        let (tx3, _rx3) = tokio::sync::mpsc::unbounded_channel();
123        reg.register("j", "phone", "ios", "dev-1", tx1);
124        let id = reg.register("j", "phone", "ios", "dev-1", tx2);
125        reg.register("someone", "laptop", "macos", "dev-2", tx3);
126
127        assert_eq!(reg.list(Some("j")).len(), 1);
128        assert_eq!(reg.list(None).len(), 2);
129
130        let sent = reg.send(Some("j"), None, LinkCommand::Pause).unwrap();
131        assert_eq!(sent.id, id);
132        assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
133
134        // Another account's device is not this account's to command.
135        assert!(
136            reg.send(Some("j"), Some("laptop"), LinkCommand::Pause)
137                .is_none()
138        );
139
140        reg.unregister(&id);
141        assert!(reg.send(Some("j"), None, LinkCommand::Pause).is_none());
142    }
143}