1use std::sync::LazyLock;
10
11use koan_core::remote::link::{LinkCommand, LinkState};
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 pub state: LinkState,
26 pub last_played_at: Option<i64>,
28 pub state_at: i64,
30}
31
32impl ClientInfo {
33 pub fn position_ms(&self) -> u64 {
35 let pos = self.state.position_ms;
36 if !self.state.playing {
37 return pos;
38 }
39 let run = (chrono::Utc::now().timestamp_millis() - self.state_at).max(0) as u64;
40 let pos = pos + run;
41 if self.state.duration_ms > 0 {
42 pos.min(self.state.duration_ms)
43 } else {
44 pos
45 }
46 }
47}
48
49struct Entry {
50 info: ClientInfo,
51 device: String,
53 tx: UnboundedSender<LinkCommand>,
54}
55
56#[derive(Debug, Clone)]
59pub struct Order {
60 pub id: String,
61 pub username: Option<String>,
63 pub client: Option<String>,
65 pub artist: String,
66 pub album: String,
67 pub play_next: bool,
69 pub created_at: i64,
71}
72
73const ORDER_TTL: i64 = 24 * 60 * 60;
75
76#[derive(Default)]
77pub struct Registry {
78 entries: Mutex<Vec<Entry>>,
79 orders: Mutex<Vec<Order>>,
80}
81
82pub fn registry() -> &'static Registry {
85 static REGISTRY: LazyLock<Registry> = LazyLock::new(Registry::default);
86 ®ISTRY
87}
88
89impl Registry {
90 pub fn register(
93 &self,
94 username: &str,
95 name: &str,
96 platform: &str,
97 device: &str,
98 tx: UnboundedSender<LinkCommand>,
99 ) -> String {
100 let id = uuid::Uuid::now_v7().to_string();
101 let mut entries = self.entries.lock();
102 entries.retain(|e| !(e.device == device && e.info.username == username));
103 entries.push(Entry {
104 info: ClientInfo {
105 id: id.clone(),
106 name: name.to_string(),
107 platform: platform.to_string(),
108 username: username.to_string(),
109 connected_at: chrono::Utc::now().timestamp(),
110 state: LinkState::default(),
111 last_played_at: None,
112 state_at: chrono::Utc::now().timestamp_millis(),
113 },
114 device: device.to_string(),
115 tx,
116 });
117 id
118 }
119
120 pub fn report(&self, id: &str, state: LinkState) {
122 let mut entries = self.entries.lock();
123 if let Some(e) = entries.iter_mut().find(|e| e.info.id == id) {
124 if state.playing || e.info.state.playing {
125 e.info.last_played_at = Some(chrono::Utc::now().timestamp());
126 }
127 e.info.state = state;
128 e.info.state_at = chrono::Utc::now().timestamp_millis();
129 }
130 }
131
132 pub fn unregister(&self, id: &str) {
133 self.entries.lock().retain(|e| e.info.id != id);
134 }
135
136 pub fn list(&self, username: Option<&str>) -> Vec<ClientInfo> {
138 let mut out: Vec<ClientInfo> = self
139 .entries
140 .lock()
141 .iter()
142 .filter(|e| username.is_none_or(|u| e.info.username == u))
143 .map(|e| e.info.clone())
144 .collect();
145 out.sort_by_key(|c| std::cmp::Reverse(c.connected_at));
146 out
147 }
148
149 pub fn send(
154 &self,
155 username: Option<&str>,
156 id: Option<&str>,
157 cmd: LinkCommand,
158 ) -> Result<ClientInfo, String> {
159 let clients = self.list(username);
160 let target = match id {
161 Some(id) => clients
162 .iter()
163 .find(|c| c.id == id || c.name.eq_ignore_ascii_case(id))
164 .ok_or_else(|| format!("no linked client {id}; see `clients`"))?,
165 None => pick(&clients, chrono::Utc::now().timestamp())?,
166 };
167 let entries = self.entries.lock();
168 let entry = entries
169 .iter()
170 .find(|e| e.info.id == target.id)
171 .ok_or("that client has just gone")?;
172 entry
173 .tx
174 .send(cmd)
175 .map_err(|_| "that client has just gone".to_string())?;
176 Ok(target.clone())
177 }
178}
179
180impl Registry {
181 pub fn add_order(&self, order: Order) {
182 self.orders.lock().push(order);
183 }
184
185 pub fn orders(&self, username: Option<&str>) -> Vec<Order> {
186 self.orders
187 .lock()
188 .iter()
189 .filter(|o| username.is_none() || o.username.as_deref() == username)
190 .cloned()
191 .collect()
192 }
193
194 pub fn cancel_order(&self, username: Option<&str>, id: &str) -> bool {
195 let mut orders = self.orders.lock();
196 let before = orders.len();
197 orders
198 .retain(|o| !(o.id == id && (username.is_none() || o.username.as_deref() == username)));
199 orders.len() != before
200 }
201
202 pub fn fulfil_orders(&self, find: impl Fn(&Order) -> Option<Vec<i64>>) {
205 let now = chrono::Utc::now().timestamp();
206 let pending: Vec<Order> = {
207 let mut orders = self.orders.lock();
208 orders.retain(|o| now - o.created_at < ORDER_TTL);
209 orders.clone()
210 };
211 for order in pending {
212 let Some(ids) = find(&order).filter(|ids| !ids.is_empty()) else {
213 continue;
214 };
215 let track_ids = ids.iter().map(i64::to_string).collect();
216 let cmd = if order.play_next {
217 LinkCommand::PlayNext { track_ids }
218 } else {
219 LinkCommand::Enqueue { track_ids }
220 };
221 match self.send(order.username.as_deref(), order.client.as_deref(), cmd) {
222 Ok(c) => {
223 log::info!(
224 "link: {} — {} arrived; queued on {}",
225 order.artist,
226 order.album,
227 c.name
228 );
229 self.orders.lock().retain(|o| o.id != order.id);
230 }
231 Err(e) => log::info!("link: {} — {} arrived but {e}", order.artist, order.album),
233 }
234 }
235 }
236}
237
238pub fn fulfil_from(db_path: &std::path::Path) {
240 let registry = registry();
241 if registry.orders.lock().is_empty() {
242 return;
243 }
244 let Ok(db) = koan_core::db::connection::Database::open(db_path) else {
245 return;
246 };
247 registry.fulfil_orders(|order| album_tracks(&db.conn, &order.artist, &order.album));
248}
249
250pub fn album_tracks(conn: &rusqlite::Connection, artist: &str, album: &str) -> Option<Vec<i64>> {
253 let like = |s: &str| format!("%{}%", s.replace(['%', '_'], ""));
254 let album_id: i64 = conn
255 .query_row(
256 "SELECT al.id FROM albums al JOIN artists a ON a.id = al.artist_id
257 WHERE a.name LIKE ?1 COLLATE NOCASE AND al.title LIKE ?2 COLLATE NOCASE
258 ORDER BY al.id DESC LIMIT 1",
259 [like(artist), like(album)],
260 |r| r.get(0),
261 )
262 .ok()?;
263 let mut stmt = conn
264 .prepare("SELECT id FROM tracks WHERE album_id = ?1 ORDER BY disc, track_number, id")
265 .ok()?;
266 let ids = stmt
267 .query_map([album_id], |r| r.get(0))
268 .ok()?
269 .filter_map(Result::ok)
270 .collect();
271 Some(ids)
272}
273
274const RECENT: i64 = 6 * 60 * 60;
277
278fn pick(clients: &[ClientInfo], now: i64) -> Result<&ClientInfo, String> {
279 if let Some(c) = clients.iter().find(|c| c.state.playing) {
280 return Ok(c);
281 }
282 if let Some(c) = clients
283 .iter()
284 .filter(|c| c.last_played_at.is_some_and(|t| now - t < RECENT))
285 .max_by_key(|c| c.last_played_at)
286 {
287 return Ok(c);
288 }
289 match clients {
290 [] => Err("no koan app is linked to this server; open koan on the device".into()),
291 [only] => Ok(only),
292 several => Err(format!(
293 "several koan apps are linked and none has played recently: {}. Ask which, then pass `client`",
294 several
295 .iter()
296 .map(|c| c.name.as_str())
297 .collect::<Vec<_>>()
298 .join(", ")
299 )),
300 }
301}
302
303#[cfg(test)]
304mod tests {
305 use super::*;
306
307 #[test]
308 fn a_reconnect_replaces_the_device_and_commands_reach_it() {
309 let reg = Registry::default();
310 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
311 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
312 let (tx3, _rx3) = tokio::sync::mpsc::unbounded_channel();
313 reg.register("j", "phone", "ios", "dev-1", tx1);
314 let id = reg.register("j", "phone", "ios", "dev-1", tx2);
315 reg.register("someone", "laptop", "macos", "dev-2", tx3);
316
317 assert_eq!(reg.list(Some("j")).len(), 1);
318 assert_eq!(reg.list(None).len(), 2);
319
320 let sent = reg.send(Some("j"), None, LinkCommand::Pause).unwrap();
321 assert_eq!(sent.id, id);
322 assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
323
324 assert!(
326 reg.send(Some("j"), Some("laptop"), LinkCommand::Pause)
327 .is_err()
328 );
329
330 reg.unregister(&id);
331 assert!(reg.send(Some("j"), None, LinkCommand::Pause).is_err());
332 }
333
334 #[test]
335 fn the_device_playing_is_the_one_meant() {
336 let reg = Registry::default();
337 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
338 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
339 let mac = reg.register("j", "mac", "macos", "dev-1", tx1);
340 let phone = reg.register("j", "phone", "ios", "dev-2", tx2);
341
342 let err = reg.send(Some("j"), None, LinkCommand::Pause).unwrap_err();
344 assert!(err.contains("mac") && err.contains("phone"), "{err}");
345
346 reg.report(
347 &phone,
348 LinkState {
349 playing: true,
350 ..Default::default()
351 },
352 );
353 assert_eq!(
354 reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
355 phone
356 );
357 assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
358
359 reg.report(&phone, LinkState::default());
361 assert_eq!(
362 reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
363 phone
364 );
365 let _ = mac;
366 }
367}