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, serde::Serialize, serde::Deserialize)]
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 #[serde(default)]
74 pub playlist: Option<i64>,
75 #[serde(default)]
78 pub titles: Vec<String>,
79 pub created_at: i64,
81}
82
83const ORDER_TTL: i64 = 24 * 60 * 60;
85
86#[derive(Default)]
87pub struct Registry {
88 entries: Mutex<Vec<Entry>>,
89 orders: Mutex<Vec<Order>>,
90}
91
92pub fn registry() -> &'static Registry {
95 static REGISTRY: LazyLock<Registry> = LazyLock::new(|| {
96 let registry = Registry::default();
97 *registry.orders.lock() = outbox::load_orders();
98 registry
99 });
100 ®ISTRY
101}
102
103impl Registry {
104 pub fn register(
107 &self,
108 username: &str,
109 name: &str,
110 platform: &str,
111 device: &str,
112 tx: UnboundedSender<LinkCommand>,
113 ) -> String {
114 let id = uuid::Uuid::now_v7().to_string();
115 for cmd in outbox::take_and_remember(username, device, name, platform) {
118 let _ = tx.send(cmd);
119 }
120 let mut entries = self.entries.lock();
121 entries.retain(|e| !(e.device == device && e.info.username == username));
122 entries.push(Entry {
123 info: ClientInfo {
124 id: id.clone(),
125 name: name.to_string(),
126 platform: platform.to_string(),
127 username: username.to_string(),
128 connected_at: chrono::Utc::now().timestamp(),
129 state: LinkState::default(),
130 last_played_at: None,
131 state_at: chrono::Utc::now().timestamp_millis(),
132 reports: false,
133 },
134 device: device.to_string(),
135 tx,
136 });
137 id
138 }
139
140 pub fn report(&self, id: &str, state: LinkState) {
142 let mut entries = self.entries.lock();
143 if let Some(e) = entries.iter_mut().find(|e| e.info.id == id) {
144 if state.playing || e.info.state.playing {
145 e.info.last_played_at = Some(chrono::Utc::now().timestamp());
146 }
147 e.info.state = state;
148 e.info.state_at = chrono::Utc::now().timestamp_millis();
149 e.info.reports = true;
150 }
151 }
152
153 pub fn unregister(&self, id: &str) {
154 self.entries.lock().retain(|e| e.info.id != id);
155 }
156
157 pub fn list(&self, username: Option<&str>) -> Vec<ClientInfo> {
159 let mut out: Vec<ClientInfo> = self
160 .entries
161 .lock()
162 .iter()
163 .filter(|e| username.is_none_or(|u| e.info.username == u))
164 .map(|e| e.info.clone())
165 .collect();
166 out.sort_by_key(|c| std::cmp::Reverse(c.connected_at));
167 out
168 }
169
170 pub fn send(
175 &self,
176 username: Option<&str>,
177 id: Option<&str>,
178 cmd: LinkCommand,
179 ) -> Result<ClientInfo, String> {
180 let clients = self.list(username);
181 let target = match id {
182 Some(id) => clients
183 .iter()
184 .find(|c| c.id == id || c.name.eq_ignore_ascii_case(id))
185 .ok_or_else(|| format!("no linked client {id}; see `clients`"))?,
186 None => pick(&clients, chrono::Utc::now().timestamp())?,
187 };
188 let entries = self.entries.lock();
189 let entry = entries
190 .iter()
191 .find(|e| e.info.id == target.id)
192 .ok_or("that client has just gone")?;
193 entry
194 .tx
195 .send(cmd)
196 .map_err(|_| "that client has just gone".to_string())?;
197 Ok(target.clone())
198 }
199}
200
201impl Registry {
202 pub fn add_order(&self, order: Order) {
203 outbox::save_order(&order);
204 self.orders.lock().push(order);
205 }
206
207 pub fn orders(&self, username: Option<&str>) -> Vec<Order> {
208 self.orders
209 .lock()
210 .iter()
211 .filter(|o| username.is_none() || o.username.as_deref() == username)
212 .cloned()
213 .collect()
214 }
215
216 pub fn cancel_order(&self, username: Option<&str>, id: &str) -> bool {
217 let mut orders = self.orders.lock();
218 let before = orders.len();
219 orders
220 .retain(|o| !(o.id == id && (username.is_none() || o.username.as_deref() == username)));
221 let gone = orders.len() != before;
222 if gone {
223 outbox::drop_order(id);
224 }
225 gone
226 }
227
228 fn done(&self, id: &str) {
229 self.orders.lock().retain(|o| o.id != id);
230 outbox::drop_order(id);
231 }
232
233 pub fn fulfil_orders(&self, find: impl Fn(&Order) -> Option<Vec<i64>>) {
236 let now = chrono::Utc::now().timestamp();
237 let pending: Vec<Order> = {
238 let mut orders = self.orders.lock();
239 for o in orders.iter().filter(|o| now - o.created_at >= ORDER_TTL) {
240 outbox::drop_order(&o.id);
241 }
242 orders.retain(|o| now - o.created_at < ORDER_TTL);
243 orders.clone()
244 };
245 for order in pending.into_iter().filter(|o| o.playlist.is_none()) {
246 let Some(ids) = find(&order).filter(|ids| !ids.is_empty()) else {
247 continue;
248 };
249 let track_ids = ids.iter().map(i64::to_string).collect();
250 let cmd = if order.play_next {
251 LinkCommand::PlayNext { track_ids }
252 } else {
253 LinkCommand::Enqueue { track_ids }
254 };
255 match self.send(order.username.as_deref(), order.client.as_deref(), cmd) {
256 Ok(c) => {
257 log::info!(
258 "link: {} — {} arrived; queued on {}",
259 order.artist,
260 order.album,
261 c.name
262 );
263 self.done(&order.id);
264 }
265 Err(e) => log::info!("link: {} — {} arrived but {e}", order.artist, order.album),
267 }
268 }
269 }
270}
271
272pub fn fulfil_from(db_path: &std::path::Path) {
274 let registry = registry();
275 if registry.orders.lock().is_empty() {
276 return;
277 }
278 let Ok(db) = koan_core::db::connection::Database::open(db_path) else {
279 return;
280 };
281 let for_playlists: Vec<Order> = registry
284 .orders
285 .lock()
286 .iter()
287 .filter(|o| o.playlist.is_some())
288 .cloned()
289 .collect();
290 let mut edited = false;
291 for order in for_playlists {
292 let Some(playlist) = order.playlist else {
293 continue;
294 };
295 let Some(ids) = order_tracks(&db.conn, &order) else {
296 continue;
297 };
298 match koan_core::db::queries::add_tracks(&db.conn, playlist, &ids) {
299 Ok(_) => {
300 if !cfg!(test) {
302 koan_core::playlists::push_to_remote(playlist);
303 }
304 log::info!(
305 "link: {} — {} arrived; added {} tracks to playlist {playlist}",
306 order.artist,
307 order.album,
308 ids.len()
309 );
310 registry.done(&order.id);
311 edited = true;
312 }
313 Err(e) => log::warn!("link: could not add to playlist {playlist}: {e}"),
314 }
315 }
316 if edited {
317 changed();
318 }
319 registry.fulfil_orders(|order| order_tracks(&db.conn, order));
320}
321
322fn order_tracks(conn: &rusqlite::Connection, order: &Order) -> Option<Vec<i64>> {
325 let tracks = album_tracks(conn, &order.artist, &order.album)?;
326 if order.titles.is_empty() {
327 return Some(tracks.into_iter().map(|(id, _)| id).collect());
328 }
329 let picked: Vec<i64> = order
330 .titles
331 .iter()
332 .filter_map(|want| {
333 let want = want.to_lowercase();
334 tracks
335 .iter()
336 .find(|(_, t)| t.to_lowercase().contains(&want))
337 .map(|(id, _)| *id)
338 })
339 .collect();
340 (!picked.is_empty()).then_some(picked)
341}
342
343pub fn album_tracks(
346 conn: &rusqlite::Connection,
347 artist: &str,
348 album: &str,
349) -> Option<Vec<(i64, String)>> {
350 let like = |s: &str| format!("%{}%", s.replace(['%', '_'], ""));
351 let album_id: i64 = conn
352 .query_row(
353 "SELECT al.id FROM albums al JOIN artists a ON a.id = al.artist_id
354 WHERE a.name LIKE ?1 COLLATE NOCASE AND al.title LIKE ?2 COLLATE NOCASE
355 ORDER BY al.id DESC LIMIT 1",
356 [like(artist), like(album)],
357 |r| r.get(0),
358 )
359 .ok()?;
360 let mut stmt = conn
361 .prepare("SELECT id, title FROM tracks WHERE album_id = ?1 ORDER BY disc, track_number, id")
362 .ok()?;
363 let tracks = stmt
364 .query_map([album_id], |r| Ok((r.get(0)?, r.get(1)?)))
365 .ok()?
366 .filter_map(Result::ok)
367 .collect();
368 Some(tracks)
369}
370
371impl Registry {
372 pub fn broadcast(&self, username: Option<&str>, cmd: LinkCommand) -> Vec<String> {
375 let ids: Vec<String> = self.list(username).into_iter().map(|c| c.id).collect();
376 let entries = self.entries.lock();
377 entries
378 .iter()
379 .filter(|e| ids.contains(&e.info.id) && e.tx.send(cmd.clone()).is_ok())
380 .map(|e| e.info.name.clone())
381 .collect()
382 }
383}
384
385impl Registry {
386 pub fn deliver(&self, username: Option<&str>, cmd: LinkCommand) -> (Vec<String>, Vec<String>) {
391 let sent = self.broadcast(username, cmd.clone());
392 let live: Vec<(String, String)> = self
393 .entries
394 .lock()
395 .iter()
396 .map(|e| (e.device.clone(), e.info.username.clone()))
397 .collect();
398 let queued = outbox::queue_for_absent(username, &live, &cmd);
399 (sent, queued)
400 }
401}
402
403pub fn changed() {
407 registry().deliver(None, LinkCommand::Sync { full: false });
408}
409
410pub fn changed_if_library_moved(db_path: &std::path::Path) {
414 static LAST: parking_lot::Mutex<Option<(i64, i64, i64)>> = parking_lot::Mutex::new(None);
415 let Ok(db) = koan_core::db::connection::Database::open(db_path) else {
416 return;
417 };
418 let Ok(now) = db.conn.query_row(
419 "SELECT (SELECT COUNT(*) FROM tracks), (SELECT COALESCE(MAX(id), 0) FROM tracks),
420 (SELECT COUNT(*) FROM albums)",
421 [],
422 |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
423 ) else {
424 return;
425 };
426 let before = LAST.lock().replace(now);
427 if before.is_some_and(|b| b != now) {
428 changed();
429 }
430}
431
432mod outbox {
435 use koan_core::remote::link::LinkCommand;
436
437 const KEEP_SECS: i64 = 30 * 24 * 60 * 60;
440
441 fn db() -> Option<koan_core::db::connection::Database> {
443 if cfg!(test) {
444 return None;
445 }
446 koan_core::db::connection::Database::open(&koan_core::config::db_path()).ok()
447 }
448
449 pub fn load_orders() -> Vec<super::Order> {
450 let Some(db) = db() else { return Vec::new() };
451 db.conn
452 .prepare("SELECT body FROM link_orders ORDER BY created_at")
453 .and_then(|mut s| {
454 s.query_map([], |r| r.get::<_, String>(0))?
455 .collect::<Result<Vec<_>, _>>()
456 })
457 .unwrap_or_default()
458 .into_iter()
459 .filter_map(|b| serde_json::from_str(&b).ok())
460 .collect()
461 }
462
463 pub fn save_order(order: &super::Order) {
464 let (Some(db), Ok(body)) = (db(), serde_json::to_string(order)) else {
465 return;
466 };
467 let _ = db.conn.execute(
468 "INSERT OR REPLACE INTO link_orders (id, body, created_at) VALUES (?1, ?2, ?3)",
469 rusqlite::params![order.id, body, order.created_at],
470 );
471 }
472
473 pub fn drop_order(id: &str) {
474 if let Some(db) = db() {
475 let _ = db
476 .conn
477 .execute("DELETE FROM link_orders WHERE id = ?1", [id]);
478 }
479 }
480
481 pub fn take_and_remember(
482 username: &str,
483 device: &str,
484 name: &str,
485 platform: &str,
486 ) -> Vec<LinkCommand> {
487 let Some(db) = db() else { return Vec::new() };
488 let now = chrono::Utc::now().timestamp();
489 let _ = db.conn.execute(
490 "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES (?1, ?2, ?3, ?4, ?5)
491 ON CONFLICT (device, username) DO UPDATE SET name = ?3, platform = ?4, last_seen = ?5",
492 rusqlite::params![device, username, name, platform, now],
493 );
494 let _ = db.conn.execute(
495 "DELETE FROM link_outbox WHERE created_at < ?1",
496 [now - KEEP_SECS],
497 );
498 let waiting: Vec<(i64, String)> = db
499 .conn
500 .prepare("SELECT id, command FROM link_outbox WHERE device = ?1 AND username = ?2 ORDER BY id")
501 .and_then(|mut s| {
502 s.query_map([device, username], |r| Ok((r.get(0)?, r.get(1)?)))?
503 .collect()
504 })
505 .unwrap_or_default();
506 let _ = db.conn.execute(
507 "DELETE FROM link_outbox WHERE device = ?1 AND username = ?2",
508 [device, username],
509 );
510 if !waiting.is_empty() {
511 log::info!("link: {} waiting commands for {name}", waiting.len());
512 }
513 waiting
514 .into_iter()
515 .filter_map(|(_, c)| serde_json::from_str(&c).ok())
516 .collect()
517 }
518
519 pub fn queue_for_absent(
521 username: Option<&str>,
522 live: &[(String, String)],
523 cmd: &LinkCommand,
524 ) -> Vec<String> {
525 let Some(db) = db() else { return Vec::new() };
526 let known: Vec<(String, String, String)> = db
527 .conn
528 .prepare("SELECT device, username, name FROM link_devices")
529 .and_then(|mut s| {
530 s.query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))?
531 .collect()
532 })
533 .unwrap_or_default();
534 let Ok(text) = serde_json::to_string(cmd) else {
535 return Vec::new();
536 };
537 let is_sync = matches!(cmd, LinkCommand::Sync { .. });
538 let now = chrono::Utc::now().timestamp();
539 let mut queued = Vec::new();
540 for (device, user, name) in known {
541 if username.is_some_and(|u| u != user)
542 || live.iter().any(|(d, u)| *d == device && *u == user)
543 {
544 continue;
545 }
546 if is_sync {
547 let _ = db.conn.execute(
549 "DELETE FROM link_outbox WHERE device = ?1 AND username = ?2 AND command LIKE '{\"type\":\"sync\"%'",
550 [&device, &user],
551 );
552 }
553 if db
554 .conn
555 .execute(
556 "INSERT INTO link_outbox (device, username, command, created_at) VALUES (?1, ?2, ?3, ?4)",
557 rusqlite::params![device, user, text, now],
558 )
559 .is_ok()
560 {
561 queued.push(name);
562 }
563 }
564 queued
565 }
566}
567
568const RECENT: i64 = 6 * 60 * 60;
571
572fn pick(clients: &[ClientInfo], now: i64) -> Result<&ClientInfo, String> {
573 if let Some(c) = clients.iter().find(|c| c.state.playing) {
574 return Ok(c);
575 }
576 if let Some(c) = clients
577 .iter()
578 .filter(|c| c.last_played_at.is_some_and(|t| now - t < RECENT))
579 .max_by_key(|c| c.last_played_at)
580 {
581 return Ok(c);
582 }
583 match clients {
584 [] => Err("no koan app is linked to this server; open koan on the device".into()),
585 [only] => Ok(only),
586 several => Err(format!(
587 "several koan apps are linked and none has played recently: {}. Ask which, then pass `client`",
588 several
589 .iter()
590 .map(|c| c.name.as_str())
591 .collect::<Vec<_>>()
592 .join(", ")
593 )),
594 }
595}
596
597#[cfg(test)]
598mod tests {
599
600 #[test]
601 fn a_playlist_order_adds_the_named_tracks_once_they_arrive() {
602 let dir = tempfile::tempdir().unwrap();
603 let path = dir.path().join("koan.db");
604 let db = koan_core::db::connection::Database::open(&path).unwrap();
605 let playlist =
606 koan_core::db::queries::create_playlist(&db.conn, "cyberpunk", None).unwrap();
607 let order = Order {
608 id: "o1".into(),
609 username: None,
610 client: None,
611 artist: "Perturbator".into(),
612 album: "Dangerous Days".into(),
613 play_next: false,
614 playlist: Some(playlist),
615 titles: vec!["Future Club".into()],
616 created_at: chrono::Utc::now().timestamp(),
617 };
618 registry().add_order(order);
619
620 fulfil_from(&path);
622 assert!(registry().orders(None).iter().any(|o| o.id == "o1"));
623
624 db.conn
625 .execute_batch(
626 "INSERT INTO artists (id, name) VALUES (1, 'Perturbator');
627 INSERT INTO albums (id, title, artist_id) VALUES (1, 'Dangerous Days', 1);
628 INSERT INTO tracks (id, title, album_id, artist_id, track_number, path) VALUES
629 (1, 'Welcome Back', 1, 1, 1, '/1.flac'), (2, 'Future Club', 1, 1, 2, '/2.flac');",
630 )
631 .unwrap();
632 fulfil_from(&path);
633 let held: Vec<i64> = db
634 .conn
635 .prepare("SELECT track_id FROM playlist_tracks WHERE playlist_id = ?1")
636 .unwrap()
637 .query_map([playlist], |r| r.get(0))
638 .unwrap()
639 .collect::<Result<_, _>>()
640 .unwrap();
641 assert_eq!(held, [2]);
642 assert!(!registry().orders(None).iter().any(|o| o.id == "o1"));
643 }
644
645 use super::*;
646
647 #[test]
648 fn a_reconnect_replaces_the_device_and_commands_reach_it() {
649 let reg = Registry::default();
650 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
651 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
652 let (tx3, _rx3) = tokio::sync::mpsc::unbounded_channel();
653 reg.register("j", "phone", "ios", "dev-1", tx1);
654 let id = reg.register("j", "phone", "ios", "dev-1", tx2);
655 reg.register("someone", "laptop", "macos", "dev-2", tx3);
656
657 assert_eq!(reg.list(Some("j")).len(), 1);
658 assert_eq!(reg.list(None).len(), 2);
659
660 let sent = reg.send(Some("j"), None, LinkCommand::Pause).unwrap();
661 assert_eq!(sent.id, id);
662 assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
663
664 assert!(
666 reg.send(Some("j"), Some("laptop"), LinkCommand::Pause)
667 .is_err()
668 );
669
670 reg.unregister(&id);
671 assert!(reg.send(Some("j"), None, LinkCommand::Pause).is_err());
672 }
673
674 #[test]
675 fn the_device_playing_is_the_one_meant() {
676 let reg = Registry::default();
677 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
678 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
679 let mac = reg.register("j", "mac", "macos", "dev-1", tx1);
680 let phone = reg.register("j", "phone", "ios", "dev-2", tx2);
681
682 let err = reg.send(Some("j"), None, LinkCommand::Pause).unwrap_err();
684 assert!(err.contains("mac") && err.contains("phone"), "{err}");
685
686 reg.report(
687 &phone,
688 LinkState {
689 playing: true,
690 ..Default::default()
691 },
692 );
693 assert_eq!(
694 reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
695 phone
696 );
697 assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
698
699 reg.report(&phone, LinkState::default());
701 assert_eq!(
702 reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
703 phone
704 );
705 let _ = mac;
706 }
707}