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 pub reports: bool,
33}
34
35impl ClientInfo {
36 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 device: String,
56 tx: UnboundedSender<LinkCommand>,
57}
58
59#[derive(Debug, Clone)]
62pub struct Order {
63 pub id: String,
64 pub username: Option<String>,
66 pub client: Option<String>,
68 pub artist: String,
69 pub album: String,
70 pub play_next: bool,
72 pub created_at: i64,
74}
75
76const 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
85pub fn registry() -> &'static Registry {
88 static REGISTRY: LazyLock<Registry> = LazyLock::new(Registry::default);
89 ®ISTRY
90}
91
92impl Registry {
93 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 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 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 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 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 Err(e) => log::info!("link: {} — {} arrived but {e}", order.artist, order.album),
238 }
239 }
240 }
241}
242
243pub 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
255pub 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
279impl Registry {
280 pub fn broadcast(&self, username: Option<&str>, cmd: LinkCommand) -> Vec<String> {
283 let ids: Vec<String> = self.list(username).into_iter().map(|c| c.id).collect();
284 let entries = self.entries.lock();
285 entries
286 .iter()
287 .filter(|e| ids.contains(&e.info.id) && e.tx.send(cmd.clone()).is_ok())
288 .map(|e| e.info.name.clone())
289 .collect()
290 }
291}
292
293const RECENT: i64 = 6 * 60 * 60;
296
297fn pick(clients: &[ClientInfo], now: i64) -> Result<&ClientInfo, String> {
298 if let Some(c) = clients.iter().find(|c| c.state.playing) {
299 return Ok(c);
300 }
301 if let Some(c) = clients
302 .iter()
303 .filter(|c| c.last_played_at.is_some_and(|t| now - t < RECENT))
304 .max_by_key(|c| c.last_played_at)
305 {
306 return Ok(c);
307 }
308 match clients {
309 [] => Err("no koan app is linked to this server; open koan on the device".into()),
310 [only] => Ok(only),
311 several => Err(format!(
312 "several koan apps are linked and none has played recently: {}. Ask which, then pass `client`",
313 several
314 .iter()
315 .map(|c| c.name.as_str())
316 .collect::<Vec<_>>()
317 .join(", ")
318 )),
319 }
320}
321
322#[cfg(test)]
323mod tests {
324 use super::*;
325
326 #[test]
327 fn a_reconnect_replaces_the_device_and_commands_reach_it() {
328 let reg = Registry::default();
329 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
330 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
331 let (tx3, _rx3) = tokio::sync::mpsc::unbounded_channel();
332 reg.register("j", "phone", "ios", "dev-1", tx1);
333 let id = reg.register("j", "phone", "ios", "dev-1", tx2);
334 reg.register("someone", "laptop", "macos", "dev-2", tx3);
335
336 assert_eq!(reg.list(Some("j")).len(), 1);
337 assert_eq!(reg.list(None).len(), 2);
338
339 let sent = reg.send(Some("j"), None, LinkCommand::Pause).unwrap();
340 assert_eq!(sent.id, id);
341 assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
342
343 assert!(
345 reg.send(Some("j"), Some("laptop"), LinkCommand::Pause)
346 .is_err()
347 );
348
349 reg.unregister(&id);
350 assert!(reg.send(Some("j"), None, LinkCommand::Pause).is_err());
351 }
352
353 #[test]
354 fn the_device_playing_is_the_one_meant() {
355 let reg = Registry::default();
356 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
357 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
358 let mac = reg.register("j", "mac", "macos", "dev-1", tx1);
359 let phone = reg.register("j", "phone", "ios", "dev-2", tx2);
360
361 let err = reg.send(Some("j"), None, LinkCommand::Pause).unwrap_err();
363 assert!(err.contains("mac") && err.contains("phone"), "{err}");
364
365 reg.report(
366 &phone,
367 LinkState {
368 playing: true,
369 ..Default::default()
370 },
371 );
372 assert_eq!(
373 reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
374 phone
375 );
376 assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
377
378 reg.report(&phone, LinkState::default());
380 assert_eq!(
381 reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
382 phone
383 );
384 let _ = mac;
385 }
386}