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 pub notified: bool,
36}
37
38impl ClientInfo {
39 pub fn position_ms(&self) -> u64 {
41 let pos = self.state.position_ms;
42 if !self.state.playing {
43 return pos;
44 }
45 let run = (chrono::Utc::now().timestamp_millis() - self.state_at).max(0) as u64;
46 let pos = pos + run;
47 if self.state.duration_ms > 0 {
48 pos.min(self.state.duration_ms)
49 } else {
50 pos
51 }
52 }
53}
54
55struct Entry {
56 info: ClientInfo,
57 device: String,
59 tx: UnboundedSender<LinkCommand>,
60}
61
62#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
65pub struct Order {
66 pub id: String,
67 pub username: Option<String>,
69 pub client: Option<String>,
71 pub artist: String,
72 pub album: String,
73 pub play_next: bool,
75 #[serde(default)]
77 pub playlist: Option<i64>,
78 #[serde(default)]
81 pub titles: Vec<String>,
82 pub created_at: i64,
84}
85
86const ORDER_TTL: i64 = 24 * 60 * 60;
88
89#[derive(Default)]
90pub struct Registry {
91 entries: Mutex<Vec<Entry>>,
92 orders: Mutex<Vec<Order>>,
93}
94
95pub fn registry() -> &'static Registry {
98 static REGISTRY: LazyLock<Registry> = LazyLock::new(|| {
99 let registry = Registry::default();
100 *registry.orders.lock() = outbox::load_orders();
101 registry
102 });
103 ®ISTRY
104}
105
106impl Registry {
107 pub fn register(
110 &self,
111 username: &str,
112 name: &str,
113 platform: &str,
114 device: &str,
115 tx: UnboundedSender<LinkCommand>,
116 ) -> String {
117 let id = uuid::Uuid::now_v7().to_string();
118 for cmd in outbox::take_and_remember(username, device, name, platform) {
121 let _ = tx.send(cmd);
122 }
123 let mut entries = self.entries.lock();
124 entries.retain(|e| !(e.device == device && e.info.username == username));
125 entries.push(Entry {
126 info: ClientInfo {
127 id: id.clone(),
128 name: name.to_string(),
129 platform: platform.to_string(),
130 username: username.to_string(),
131 connected_at: chrono::Utc::now().timestamp(),
132 state: LinkState::default(),
133 last_played_at: None,
134 state_at: chrono::Utc::now().timestamp_millis(),
135 reports: false,
136 notified: false,
137 },
138 device: device.to_string(),
139 tx,
140 });
141 id
142 }
143
144 pub fn report(&self, id: &str, state: LinkState) {
146 let mut entries = self.entries.lock();
147 if let Some(e) = entries.iter_mut().find(|e| e.info.id == id) {
148 if state.playing || e.info.state.playing {
149 e.info.last_played_at = Some(chrono::Utc::now().timestamp());
150 }
151 e.info.state = state;
152 e.info.state_at = chrono::Utc::now().timestamp_millis();
153 e.info.reports = true;
154 }
155 }
156
157 pub fn set_push(&self, username: &str, device: &str, token: &str, sandbox: bool) {
159 outbox::save_push(username, device, token, sandbox);
160 }
161
162 pub fn unregister(&self, id: &str) {
163 self.entries.lock().retain(|e| e.info.id != id);
164 }
165
166 pub fn list(&self, username: Option<&str>) -> Vec<ClientInfo> {
168 let mut out: Vec<ClientInfo> = self
169 .entries
170 .lock()
171 .iter()
172 .filter(|e| username.is_none_or(|u| e.info.username == u))
173 .map(|e| e.info.clone())
174 .collect();
175 out.sort_by_key(|c| std::cmp::Reverse(c.connected_at));
176 out
177 }
178
179 pub fn send(
184 &self,
185 username: Option<&str>,
186 id: Option<&str>,
187 cmd: LinkCommand,
188 ) -> Result<ClientInfo, String> {
189 let clients = self.list(username);
190 let target = match id {
191 Some(id) => clients
192 .iter()
193 .find(|c| c.id == id || c.name.eq_ignore_ascii_case(id)),
194 None if clients.is_empty() => None,
195 None => Some(pick(&clients, chrono::Utc::now().timestamp())?),
196 };
197 let Some(target) = target else {
199 return reach_absent(username, id, &cmd).unwrap_or_else(|| {
200 Err(match id {
201 Some(id) => format!("no linked client {id}; see `clients`"),
202 None => "no koan app is linked to this server; open koan on the device".into(),
203 })
204 });
205 };
206 let entries = self.entries.lock();
207 let entry = entries
208 .iter()
209 .find(|e| e.info.id == target.id)
210 .ok_or("that client has just gone")?;
211 entry
212 .tx
213 .send(cmd)
214 .map_err(|_| "that client has just gone".to_string())?;
215 Ok(target.clone())
216 }
217}
218
219impl Registry {
220 pub fn add_order(&self, order: Order) {
221 outbox::save_order(&order);
222 self.orders.lock().push(order);
223 }
224
225 pub fn orders(&self, username: Option<&str>) -> Vec<Order> {
226 self.orders
227 .lock()
228 .iter()
229 .filter(|o| username.is_none() || o.username.as_deref() == username)
230 .cloned()
231 .collect()
232 }
233
234 pub fn cancel_order(&self, username: Option<&str>, id: &str) -> bool {
235 let mut orders = self.orders.lock();
236 let before = orders.len();
237 orders
238 .retain(|o| !(o.id == id && (username.is_none() || o.username.as_deref() == username)));
239 let gone = orders.len() != before;
240 if gone {
241 outbox::drop_order(id);
242 }
243 gone
244 }
245
246 fn done(&self, id: &str) {
247 self.orders.lock().retain(|o| o.id != id);
248 outbox::drop_order(id);
249 }
250
251 pub fn fulfil_orders(&self, find: impl Fn(&Order) -> Option<Vec<i64>>) {
254 let now = chrono::Utc::now().timestamp();
255 let pending: Vec<Order> = {
256 let mut orders = self.orders.lock();
257 for o in orders.iter().filter(|o| now - o.created_at >= ORDER_TTL) {
258 outbox::drop_order(&o.id);
259 }
260 orders.retain(|o| now - o.created_at < ORDER_TTL);
261 orders.clone()
262 };
263 for order in pending.into_iter().filter(|o| o.playlist.is_none()) {
264 let Some(ids) = find(&order).filter(|ids| !ids.is_empty()) else {
265 continue;
266 };
267 let track_ids = ids.iter().map(i64::to_string).collect();
268 let cmd = if order.play_next {
269 LinkCommand::PlayNext { track_ids }
270 } else {
271 LinkCommand::Enqueue { track_ids }
272 };
273 match self.send(order.username.as_deref(), order.client.as_deref(), cmd) {
274 Ok(c) => {
275 log::info!(
276 "link: {} — {} arrived; queued on {}",
277 order.artist,
278 order.album,
279 c.name
280 );
281 self.done(&order.id);
282 }
283 Err(e) => log::info!("link: {} — {} arrived but {e}", order.artist, order.album),
285 }
286 }
287 }
288}
289
290pub fn fulfil_from(db_path: &std::path::Path) {
292 let registry = registry();
293 if registry.orders.lock().is_empty() {
294 return;
295 }
296 let Ok(db) = koan_core::db::connection::Database::open(db_path) else {
297 return;
298 };
299 let for_playlists: Vec<Order> = registry
302 .orders
303 .lock()
304 .iter()
305 .filter(|o| o.playlist.is_some())
306 .cloned()
307 .collect();
308 let mut edited = false;
309 for order in for_playlists {
310 let Some(playlist) = order.playlist else {
311 continue;
312 };
313 let Some(ids) = order_tracks(&db.conn, &order) else {
314 continue;
315 };
316 match koan_core::db::queries::add_tracks(&db.conn, playlist, &ids) {
317 Ok(_) => {
318 if !cfg!(test) {
320 koan_core::playlists::push_to_remote(playlist);
321 }
322 log::info!(
323 "link: {} — {} arrived; added {} tracks to playlist {playlist}",
324 order.artist,
325 order.album,
326 ids.len()
327 );
328 registry.done(&order.id);
329 edited = true;
330 }
331 Err(e) => log::warn!("link: could not add to playlist {playlist}: {e}"),
332 }
333 }
334 if edited {
335 changed();
336 }
337 registry.fulfil_orders(|order| order_tracks(&db.conn, order));
338}
339
340fn order_tracks(conn: &rusqlite::Connection, order: &Order) -> Option<Vec<i64>> {
343 let tracks = album_tracks(conn, &order.artist, &order.album)?;
344 if order.titles.is_empty() {
345 return Some(tracks.into_iter().map(|(id, _)| id).collect());
346 }
347 let picked: Vec<i64> = order
348 .titles
349 .iter()
350 .filter_map(|want| {
351 let want = want.to_lowercase();
352 tracks
353 .iter()
354 .find(|(_, t)| t.to_lowercase().contains(&want))
355 .map(|(id, _)| *id)
356 })
357 .collect();
358 (!picked.is_empty()).then_some(picked)
359}
360
361pub fn album_tracks(
364 conn: &rusqlite::Connection,
365 artist: &str,
366 album: &str,
367) -> Option<Vec<(i64, String)>> {
368 let like = |s: &str| format!("%{}%", s.replace(['%', '_'], ""));
369 let album_id: i64 = conn
370 .query_row(
371 "SELECT al.id FROM albums al JOIN artists a ON a.id = al.artist_id
372 WHERE a.name LIKE ?1 COLLATE NOCASE AND al.title LIKE ?2 COLLATE NOCASE
373 ORDER BY al.id DESC LIMIT 1",
374 [like(artist), like(album)],
375 |r| r.get(0),
376 )
377 .ok()?;
378 let mut stmt = conn
379 .prepare("SELECT id, title FROM tracks WHERE album_id = ?1 ORDER BY disc, track_number, id")
380 .ok()?;
381 let tracks = stmt
382 .query_map([album_id], |r| Ok((r.get(0)?, r.get(1)?)))
383 .ok()?
384 .filter_map(Result::ok)
385 .collect();
386 Some(tracks)
387}
388
389impl Registry {
390 pub fn broadcast(&self, username: Option<&str>, cmd: LinkCommand) -> Vec<String> {
393 let ids: Vec<String> = self.list(username).into_iter().map(|c| c.id).collect();
394 let entries = self.entries.lock();
395 entries
396 .iter()
397 .filter(|e| ids.contains(&e.info.id) && e.tx.send(cmd.clone()).is_ok())
398 .map(|e| e.info.name.clone())
399 .collect()
400 }
401}
402
403impl Registry {
404 pub fn deliver(&self, username: Option<&str>, cmd: LinkCommand) -> (Vec<String>, Vec<String>) {
409 let sent = self.broadcast(username, cmd.clone());
410 let live: Vec<(String, String)> = self
411 .entries
412 .lock()
413 .iter()
414 .map(|e| (e.device.clone(), e.info.username.clone()))
415 .collect();
416 let queued = outbox::queue_for_absent(username, &live, &cmd);
417 wake(&queued);
418 (sent, queued.into_iter().map(|q| q.name).collect())
419 }
420}
421
422fn wake(devices: &[outbox::Absent]) {
425 let Some(pusher) = crate::push::pusher() else {
426 return;
427 };
428 let targets: Vec<outbox::PushTarget> = outbox::push_targets(None)
429 .into_iter()
430 .filter(|t| {
431 devices
432 .iter()
433 .any(|d| d.device == t.device && d.username == t.username)
434 })
435 .collect();
436 if targets.is_empty() {
437 return;
438 }
439 std::thread::spawn(move || {
440 for t in targets {
441 deliver_push(pusher, &t, &crate::push::Push::Wake);
442 }
443 });
444}
445
446fn reach_absent(
454 username: Option<&str>,
455 id: Option<&str>,
456 cmd: &LinkCommand,
457) -> Option<Result<ClientInfo, String>> {
458 let pusher = crate::push::pusher()?;
459 let targets: Vec<outbox::PushTarget> = outbox::push_targets(username)
460 .into_iter()
461 .filter(|t| id.is_none_or(|id| t.device == id || t.name.eq_ignore_ascii_case(id)))
462 .collect();
463 let target = match targets.as_slice() {
464 [] => return None,
465 [only] => only.clone(),
466 several => {
467 return Some(Err(format!(
468 "no koan app is linked, and several can be reached: {}. Ask which, then pass `client`",
469 several
470 .iter()
471 .map(|t| t.name.as_str())
472 .collect::<Vec<_>>()
473 .join(", ")
474 )));
475 }
476 };
477 let info = ClientInfo {
478 id: target.device.clone(),
479 name: target.name.clone(),
480 platform: target.platform.clone(),
481 username: target.username.clone(),
482 connected_at: 0,
483 state: LinkState::default(),
484 last_played_at: None,
485 state_at: 0,
486 reports: false,
487 notified: true,
488 };
489 let verb = match cmd {
490 LinkCommand::Play { .. } | LinkCommand::JumpTo { .. } | LinkCommand::Resume => Some("Play"),
491 _ => None,
492 };
493 let push = match verb {
494 Some(verb) => crate::push::Push::Notify {
495 title: format!("{verb} on {}", target.name),
496 body: outbox::describe(cmd).unwrap_or_else(|| "From your koan server".into()),
497 command: serde_json::to_value(cmd).ok()?,
498 },
499 None => {
500 outbox::queue_for(&target.device, &target.username, cmd);
501 crate::push::Push::Wake
502 }
503 };
504 std::thread::spawn(move || deliver_push(pusher, &target, &push));
505 Some(Ok(info))
506}
507
508fn deliver_push(
510 pusher: &crate::push::Pusher,
511 target: &outbox::PushTarget,
512 push: &crate::push::Push,
513) {
514 use crate::push::Outcome;
515 match pusher.send(&target.token, target.sandbox, push) {
516 Outcome::Sent => log::info!("push: sent to {}", target.name),
517 Outcome::Gone => {
518 log::info!(
519 "push: {}'s token is no longer valid; forgotten",
520 target.name
521 );
522 outbox::forget_push(&target.username, &target.device);
523 }
524 Outcome::Failed(e) => log::warn!("push: to {} failed: {e}", target.name),
525 }
526}
527
528pub fn changed() {
532 registry().deliver(None, LinkCommand::Sync { full: false });
533}
534
535pub fn changed_if_library_moved(db_path: &std::path::Path) {
539 static LAST: parking_lot::Mutex<Option<(i64, i64, i64)>> = parking_lot::Mutex::new(None);
540 let Ok(db) = koan_core::db::connection::Database::open(db_path) else {
541 return;
542 };
543 let Ok(now) = db.conn.query_row(
544 "SELECT (SELECT COUNT(*) FROM tracks), (SELECT COALESCE(MAX(id), 0) FROM tracks),
545 (SELECT COUNT(*) FROM albums)",
546 [],
547 |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
548 ) else {
549 return;
550 };
551 let before = LAST.lock().replace(now);
552 if before.is_some_and(|b| b != now) {
553 changed();
554 }
555}
556
557mod outbox {
560 use koan_core::remote::link::LinkCommand;
561
562 const KEEP_SECS: i64 = 30 * 24 * 60 * 60;
565
566 fn db() -> Option<koan_core::db::connection::Database> {
568 if cfg!(test) {
569 return None;
570 }
571 koan_core::db::connection::Database::open(&koan_core::config::db_path()).ok()
572 }
573
574 pub fn load_orders() -> Vec<super::Order> {
575 let Some(db) = db() else { return Vec::new() };
576 db.conn
577 .prepare("SELECT body FROM link_orders ORDER BY created_at")
578 .and_then(|mut s| {
579 s.query_map([], |r| r.get::<_, String>(0))?
580 .collect::<Result<Vec<_>, _>>()
581 })
582 .unwrap_or_default()
583 .into_iter()
584 .filter_map(|b| serde_json::from_str(&b).ok())
585 .collect()
586 }
587
588 pub fn save_order(order: &super::Order) {
589 let (Some(db), Ok(body)) = (db(), serde_json::to_string(order)) else {
590 return;
591 };
592 let _ = db.conn.execute(
593 "INSERT OR REPLACE INTO link_orders (id, body, created_at) VALUES (?1, ?2, ?3)",
594 rusqlite::params![order.id, body, order.created_at],
595 );
596 }
597
598 pub fn drop_order(id: &str) {
599 if let Some(db) = db() {
600 let _ = db
601 .conn
602 .execute("DELETE FROM link_orders WHERE id = ?1", [id]);
603 }
604 }
605
606 pub fn take_and_remember(
607 username: &str,
608 device: &str,
609 name: &str,
610 platform: &str,
611 ) -> Vec<LinkCommand> {
612 let Some(db) = db() else { return Vec::new() };
613 let now = chrono::Utc::now().timestamp();
614 let _ = db.conn.execute(
615 "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES (?1, ?2, ?3, ?4, ?5)
616 ON CONFLICT (device, username) DO UPDATE SET name = ?3, platform = ?4, last_seen = ?5",
617 rusqlite::params![device, username, name, platform, now],
618 );
619 let _ = db.conn.execute(
620 "DELETE FROM link_outbox WHERE created_at < ?1",
621 [now - KEEP_SECS],
622 );
623 let waiting: Vec<(i64, String)> = db
624 .conn
625 .prepare("SELECT id, command FROM link_outbox WHERE device = ?1 AND username = ?2 ORDER BY id")
626 .and_then(|mut s| {
627 s.query_map([device, username], |r| Ok((r.get(0)?, r.get(1)?)))?
628 .collect()
629 })
630 .unwrap_or_default();
631 let _ = db.conn.execute(
632 "DELETE FROM link_outbox WHERE device = ?1 AND username = ?2",
633 [device, username],
634 );
635 if !waiting.is_empty() {
636 log::info!("link: {} waiting commands for {name}", waiting.len());
637 }
638 waiting
639 .into_iter()
640 .filter_map(|(_, c)| serde_json::from_str(&c).ok())
641 .collect()
642 }
643
644 pub struct Absent {
646 pub device: String,
647 pub username: String,
648 pub name: String,
649 }
650
651 pub fn queue_for_absent(
653 username: Option<&str>,
654 live: &[(String, String)],
655 cmd: &LinkCommand,
656 ) -> Vec<Absent> {
657 let Some(db) = db() else { return Vec::new() };
658 let known: Vec<(String, String, String)> = db
659 .conn
660 .prepare("SELECT device, username, name FROM link_devices")
661 .and_then(|mut s| {
662 s.query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))?
663 .collect()
664 })
665 .unwrap_or_default();
666 let Ok(text) = serde_json::to_string(cmd) else {
667 return Vec::new();
668 };
669 let is_sync = matches!(cmd, LinkCommand::Sync { .. });
670 let now = chrono::Utc::now().timestamp();
671 let mut queued = Vec::new();
672 for (device, user, name) in known {
673 if username.is_some_and(|u| u != user)
674 || live.iter().any(|(d, u)| *d == device && *u == user)
675 {
676 continue;
677 }
678 if is_sync {
679 let _ = db.conn.execute(
681 "DELETE FROM link_outbox WHERE device = ?1 AND username = ?2 AND command LIKE '{\"type\":\"sync\"%'",
682 [&device, &user],
683 );
684 }
685 if db
686 .conn
687 .execute(
688 "INSERT INTO link_outbox (device, username, command, created_at) VALUES (?1, ?2, ?3, ?4)",
689 rusqlite::params![device, user, text, now],
690 )
691 .is_ok()
692 {
693 queued.push(Absent {
694 device,
695 username: user,
696 name,
697 });
698 }
699 }
700 queued
701 }
702
703 #[derive(Clone)]
705 pub struct PushTarget {
706 pub device: String,
707 pub username: String,
708 pub name: String,
709 pub platform: String,
710 pub token: String,
711 pub sandbox: bool,
712 }
713
714 pub fn queue_for(device: &str, username: &str, cmd: &LinkCommand) {
716 let (Some(db), Ok(text)) = (db(), serde_json::to_string(cmd)) else {
717 return;
718 };
719 let _ = db.conn.execute(
720 "INSERT INTO link_outbox (device, username, command, created_at) VALUES (?1, ?2, ?3, ?4)",
721 rusqlite::params![device, username, text, chrono::Utc::now().timestamp()],
722 );
723 }
724
725 pub fn save_push(username: &str, device: &str, token: &str, sandbox: bool) {
726 let Some(db) = db() else { return };
727 let _ = db.conn.execute(
728 "INSERT INTO link_push (device, username, token, sandbox, updated_at) VALUES (?1, ?2, ?3, ?4, ?5)
729 ON CONFLICT (device, username) DO UPDATE SET token = ?3, sandbox = ?4, updated_at = ?5",
730 rusqlite::params![device, username, token, sandbox, chrono::Utc::now().timestamp()],
731 );
732 }
733
734 pub fn forget_push(username: &str, device: &str) {
735 if let Some(db) = db() {
736 let _ = db.conn.execute(
737 "DELETE FROM link_push WHERE device = ?1 AND username = ?2",
738 [device, username],
739 );
740 }
741 }
742
743 pub fn push_targets(username: Option<&str>) -> Vec<PushTarget> {
745 let Some(db) = db() else { return Vec::new() };
746 db.conn
747 .prepare(
748 "SELECT p.device, p.username, d.name, d.platform, p.token, p.sandbox
749 FROM link_push p JOIN link_devices d ON d.device = p.device AND d.username = p.username
750 WHERE ?1 IS NULL OR p.username = ?1
751 ORDER BY d.last_seen DESC",
752 )
753 .and_then(|mut s| {
754 s.query_map([username], |r| {
755 Ok(PushTarget {
756 device: r.get(0)?,
757 username: r.get(1)?,
758 name: r.get(2)?,
759 platform: r.get(3)?,
760 token: r.get(4)?,
761 sandbox: r.get(5)?,
762 })
763 })?
764 .collect()
765 })
766 .unwrap_or_default()
767 }
768
769 pub fn describe(cmd: &LinkCommand) -> Option<String> {
772 let ids: Vec<i64> = match cmd {
773 LinkCommand::Play { track_ids, .. }
774 | LinkCommand::Enqueue { track_ids }
775 | LinkCommand::PlayNext { track_ids } => {
776 track_ids.iter().filter_map(|t| t.parse().ok()).collect()
777 }
778 LinkCommand::JumpTo { track_id } => vec![track_id.parse().ok()?],
779 _ => return None,
780 };
781 let db = db()?;
782 let row = |id: i64| {
783 db.conn
784 .query_row(
785 "SELECT t.title, COALESCE(a.name, ''), COALESCE(al.title, ''), t.album_id
786 FROM tracks t LEFT JOIN artists a ON a.id = t.artist_id
787 LEFT JOIN albums al ON al.id = t.album_id WHERE t.id = ?1",
788 [id],
789 |r| {
790 Ok((
791 r.get::<_, String>(0)?,
792 r.get::<_, String>(1)?,
793 r.get::<_, String>(2)?,
794 r.get::<_, Option<i64>>(3)?,
795 ))
796 },
797 )
798 .ok()
799 };
800 let (title, artist, album, album_id) = row(*ids.first()?)?;
801 let one_album = ids.len() > 1
802 && album_id.is_some()
803 && ids
804 .iter()
805 .all(|id| row(*id).is_some_and(|r| r.3 == album_id));
806 Some(match (one_album, ids.len()) {
807 (true, _) => format!("{album} — {artist}"),
808 (false, 1) => format!("{title} — {artist}"),
809 (false, n) => format!("{title} — {artist}, and {} more", n - 1),
810 })
811 }
812}
813
814const RECENT: i64 = 6 * 60 * 60;
817
818fn pick(clients: &[ClientInfo], now: i64) -> Result<&ClientInfo, String> {
819 if let Some(c) = clients.iter().find(|c| c.state.playing) {
820 return Ok(c);
821 }
822 if let Some(c) = clients
823 .iter()
824 .filter(|c| c.last_played_at.is_some_and(|t| now - t < RECENT))
825 .max_by_key(|c| c.last_played_at)
826 {
827 return Ok(c);
828 }
829 match clients {
830 [] => Err("no koan app is linked to this server; open koan on the device".into()),
831 [only] => Ok(only),
832 several => Err(format!(
833 "several koan apps are linked and none has played recently: {}. Ask which, then pass `client`",
834 several
835 .iter()
836 .map(|c| c.name.as_str())
837 .collect::<Vec<_>>()
838 .join(", ")
839 )),
840 }
841}
842
843#[cfg(test)]
844mod tests {
845
846 #[test]
847 fn a_playlist_order_adds_the_named_tracks_once_they_arrive() {
848 let dir = tempfile::tempdir().unwrap();
849 let path = dir.path().join("koan.db");
850 let db = koan_core::db::connection::Database::open(&path).unwrap();
851 let playlist =
852 koan_core::db::queries::create_playlist(&db.conn, "cyberpunk", None).unwrap();
853 let order = Order {
854 id: "o1".into(),
855 username: None,
856 client: None,
857 artist: "Perturbator".into(),
858 album: "Dangerous Days".into(),
859 play_next: false,
860 playlist: Some(playlist),
861 titles: vec!["Future Club".into()],
862 created_at: chrono::Utc::now().timestamp(),
863 };
864 registry().add_order(order);
865
866 fulfil_from(&path);
868 assert!(registry().orders(None).iter().any(|o| o.id == "o1"));
869
870 db.conn
871 .execute_batch(
872 "INSERT INTO artists (id, name) VALUES (1, 'Perturbator');
873 INSERT INTO albums (id, title, artist_id) VALUES (1, 'Dangerous Days', 1);
874 INSERT INTO tracks (id, title, album_id, artist_id, track_number, path) VALUES
875 (1, 'Welcome Back', 1, 1, 1, '/1.flac'), (2, 'Future Club', 1, 1, 2, '/2.flac');",
876 )
877 .unwrap();
878 fulfil_from(&path);
879 let held: Vec<i64> = db
880 .conn
881 .prepare("SELECT track_id FROM playlist_tracks WHERE playlist_id = ?1")
882 .unwrap()
883 .query_map([playlist], |r| r.get(0))
884 .unwrap()
885 .collect::<Result<_, _>>()
886 .unwrap();
887 assert_eq!(held, [2]);
888 assert!(!registry().orders(None).iter().any(|o| o.id == "o1"));
889 }
890
891 use super::*;
892
893 #[test]
894 fn a_reconnect_replaces_the_device_and_commands_reach_it() {
895 let reg = Registry::default();
896 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
897 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
898 let (tx3, _rx3) = tokio::sync::mpsc::unbounded_channel();
899 reg.register("j", "phone", "ios", "dev-1", tx1);
900 let id = reg.register("j", "phone", "ios", "dev-1", tx2);
901 reg.register("someone", "laptop", "macos", "dev-2", tx3);
902
903 assert_eq!(reg.list(Some("j")).len(), 1);
904 assert_eq!(reg.list(None).len(), 2);
905
906 let sent = reg.send(Some("j"), None, LinkCommand::Pause).unwrap();
907 assert_eq!(sent.id, id);
908 assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
909
910 assert!(
912 reg.send(Some("j"), Some("laptop"), LinkCommand::Pause)
913 .is_err()
914 );
915
916 reg.unregister(&id);
917 assert!(reg.send(Some("j"), None, LinkCommand::Pause).is_err());
918 }
919
920 #[test]
921 fn the_device_playing_is_the_one_meant() {
922 let reg = Registry::default();
923 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
924 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
925 let mac = reg.register("j", "mac", "macos", "dev-1", tx1);
926 let phone = reg.register("j", "phone", "ios", "dev-2", tx2);
927
928 let err = reg.send(Some("j"), None, LinkCommand::Pause).unwrap_err();
930 assert!(err.contains("mac") && err.contains("phone"), "{err}");
931
932 reg.report(
933 &phone,
934 LinkState {
935 playing: true,
936 ..Default::default()
937 },
938 );
939 assert_eq!(
940 reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
941 phone
942 );
943 assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
944
945 reg.report(&phone, LinkState::default());
947 assert_eq!(
948 reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
949 phone
950 );
951 let _ = mac;
952 }
953}