1use std::sync::LazyLock;
10
11use koan_core::remote::link::LinkCommand;
12use parking_lot::Mutex;
13use tokio::sync::mpsc::UnboundedSender;
14
15#[derive(Debug, Clone)]
17pub struct ClientInfo {
18 pub id: String,
19 pub name: String,
20 pub platform: String,
21 pub username: String,
22 pub connected_at: i64,
24}
25
26struct Entry {
27 info: ClientInfo,
28 device: String,
30 tx: UnboundedSender<LinkCommand>,
31}
32
33#[derive(Default)]
34pub struct Registry {
35 entries: Mutex<Vec<Entry>>,
36}
37
38pub fn registry() -> &'static Registry {
41 static REGISTRY: LazyLock<Registry> = LazyLock::new(Registry::default);
42 ®ISTRY
43}
44
45impl Registry {
46 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 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 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 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}