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