1use std::sync::LazyLock;
10
11use koan_core::db::queries::{self, UidKind};
12use koan_core::remote::link::{LinkCommand, LinkDevice, LinkState};
13use outbox::Absent;
14use parking_lot::Mutex;
15use tokio::sync::mpsc::UnboundedSender;
16
17#[derive(Debug, Clone)]
19pub struct ClientInfo {
20 pub id: String,
21 pub device: String,
23 pub name: String,
24 pub platform: String,
25 pub username: String,
26 pub connected_at: i64,
28 pub state: LinkState,
30 pub last_played_at: Option<i64>,
32 pub state_at: i64,
34 pub reports: bool,
37 pub notified: bool,
40}
41
42impl ClientInfo {
43 pub fn position_ms(&self) -> u64 {
45 let pos = self.state.position_ms;
46 if !self.state.playing {
47 return pos;
48 }
49 let run = (chrono::Utc::now().timestamp_millis() - self.state_at).max(0) as u64;
50 let pos = pos + run;
51 if self.state.duration_ms > 0 {
52 pos.min(self.state.duration_ms)
53 } else {
54 pos
55 }
56 }
57}
58
59struct Entry {
60 info: ClientInfo,
61 device: String,
63 tx: UnboundedSender<LinkCommand>,
64 wants_devices: bool,
67}
68
69struct Activity {
72 username: String,
73 watcher: String,
75 target: String,
77 token: String,
78 sandbox: bool,
79 sent: Option<crate::push::ActivityState>,
81}
82
83#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
86pub struct Order {
87 pub id: String,
88 pub username: Option<String>,
90 pub client: Option<String>,
92 pub artist: String,
93 pub album: String,
94 pub play_next: bool,
96 #[serde(default)]
98 pub playlist: Option<i64>,
99 #[serde(default)]
102 pub titles: Vec<String>,
103 pub created_at: i64,
105}
106
107const ORDER_TTL: i64 = 24 * 60 * 60;
109
110#[derive(Default)]
111pub struct Registry {
112 entries: Mutex<Vec<Entry>>,
113 orders: Mutex<Vec<Order>>,
114 activities: Mutex<Vec<Activity>>,
115 level_watches: Mutex<Vec<(String, String, String)>>,
119}
120
121pub fn registry() -> &'static Registry {
124 static REGISTRY: LazyLock<Registry> = LazyLock::new(|| {
125 let registry = Registry::default();
126 *registry.orders.lock() = outbox::load_orders();
127 registry
128 });
129 ®ISTRY
130}
131
132impl Registry {
133 pub fn register(
136 &self,
137 username: &str,
138 name: &str,
139 platform: &str,
140 device: &str,
141 tx: UnboundedSender<LinkCommand>,
142 wants_devices: bool,
143 ) -> String {
144 let id = uuid::Uuid::now_v7().to_string();
145 let live = self.live();
148 for cmd in outbox::take_and_remember(username, device, name, platform, &live) {
149 let _ = tx.send(cmd);
150 }
151 let mut entries = self.entries.lock();
152 entries.retain(|e| !(e.device == device && e.info.username == username));
153 entries.push(Entry {
154 info: ClientInfo {
155 id: id.clone(),
156 device: device.to_string(),
157 name: name.to_string(),
158 platform: platform.to_string(),
159 username: username.to_string(),
160 connected_at: chrono::Utc::now().timestamp(),
161 state: LinkState::default(),
162 last_played_at: None,
163 state_at: chrono::Utc::now().timestamp_millis(),
164 reports: false,
165 notified: false,
166 },
167 device: device.to_string(),
168 tx,
169 wants_devices,
170 });
171 drop(entries);
172 self.announce(username);
173 if self
176 .level_watches
177 .lock()
178 .iter()
179 .any(|(u, _, t)| u == username && t == device)
180 {
181 self.send_live(username, device, LinkCommand::WatchLevels { on: true });
182 }
183 id
184 }
185
186 pub fn report(&self, id: &str, state: LinkState) {
188 let mut entries = self.entries.lock();
189 if let Some(e) = entries.iter_mut().find(|e| e.info.id == id) {
190 if state.playing || e.info.state.playing {
191 e.info.last_played_at = Some(chrono::Utc::now().timestamp());
192 }
193 e.info.state = state;
194 e.info.state_at = chrono::Utc::now().timestamp_millis();
195 e.info.reports = true;
196 let (username, device) = (e.info.username.clone(), e.device.clone());
197 drop(entries);
198 self.announce(&username);
199 self.update_activities(&username, &device);
200 }
201 }
202
203 pub fn set_push(&self, username: &str, device: &str, token: &str, sandbox: bool) {
205 outbox::save_push(username, device, token, sandbox);
206 }
207
208 pub fn disconnect(&self, username: &str) {
211 self.entries.lock().retain(|e| e.info.username != username);
212 }
213
214 pub fn unregister(&self, id: &str) {
215 let mut entries = self.entries.lock();
216 let gone = entries
217 .iter()
218 .find(|e| e.info.id == id)
219 .map(|e| (e.device.clone(), e.info.username.clone()));
220 entries.retain(|e| e.info.id != id);
221 drop(entries);
222 if let Some((device, username)) = gone {
223 outbox::touch(&[(device.clone(), username.clone())]);
224 self.announce(&username);
225 let targets: Vec<String> = self
227 .level_watches
228 .lock()
229 .iter()
230 .filter(|(u, w, _)| *u == username && *w == device)
231 .map(|(_, _, t)| t.clone())
232 .collect();
233 for target in targets {
234 self.watch_levels(&username, &device, &target, false);
235 }
236 }
237 }
238
239 pub fn watch_levels(&self, username: &str, watcher: &str, target: &str, on: bool) {
244 let mut watches = self.level_watches.lock();
245 watches.retain(|(u, w, t)| !(u == username && w == watcher && t == target));
246 if on {
247 watches.push((username.into(), watcher.into(), target.into()));
248 }
249 let watched = watches.iter().any(|(u, _, t)| u == username && t == target);
250 drop(watches);
251 if on || !watched {
252 self.send_live(username, target, LinkCommand::WatchLevels { on: watched });
253 }
254 }
255
256 pub fn levels(&self, username: &str, from: &str, f: koan_core::remote::levels::Frame) {
259 let watchers: Vec<String> = self
260 .level_watches
261 .lock()
262 .iter()
263 .filter(|(u, _, t)| u == username && t == from)
264 .map(|(_, w, _)| w.clone())
265 .collect();
266 for watcher in watchers {
267 self.send_live(
268 username,
269 &watcher,
270 LinkCommand::Levels {
271 from: from.to_string(),
272 f,
273 },
274 );
275 }
276 }
277
278 fn send_live(&self, username: &str, device: &str, cmd: LinkCommand) {
280 if let Some(e) = self
281 .entries
282 .lock()
283 .iter()
284 .find(|e| e.info.username == username && e.device == device)
285 {
286 let _ = e.tx.send(cmd);
287 }
288 }
289
290 fn announce(&self, username: &str) {
294 let listening = |e: &Entry| e.info.username == username && e.wants_devices;
295 if !self.entries.lock().iter().any(listening) {
296 return;
297 }
298 let asleep = outbox::push_targets(Some(username));
299 let entries = self.entries.lock();
300 let ours: Vec<&Entry> = entries
301 .iter()
302 .filter(|e| e.info.username == username)
303 .collect();
304 if !ours.iter().any(|e| e.wants_devices) {
305 return;
306 }
307 let mut all: Vec<LinkDevice> = ours
308 .iter()
309 .map(|e| LinkDevice {
310 id: e.device.clone(),
311 name: e.info.name.clone(),
312 platform: e.info.platform.clone(),
313 linked: true,
314 state: e.info.reports.then(|| LinkState {
315 position_ms: e.info.position_ms(),
316 ..e.info.state.clone()
317 }),
318 })
319 .collect();
320 for t in asleep {
321 if !all.iter().any(|d| d.id == t.device) {
322 all.push(LinkDevice {
323 id: t.device,
324 name: t.name,
325 platform: t.platform,
326 linked: false,
327 state: None,
328 });
329 }
330 }
331 for e in ours.iter().filter(|e| e.wants_devices) {
332 let devices = all.iter().filter(|d| d.id != e.device).cloned().collect();
333 let _ = e.tx.send(LinkCommand::Devices { devices });
334 }
335 }
336
337 pub fn relay(
339 &self,
340 username: &str,
341 to: &str,
342 command: LinkCommand,
343 ) -> Result<ClientInfo, String> {
344 if matches!(
347 command,
348 LinkCommand::Devices { .. }
349 | LinkCommand::Levels { .. }
350 | LinkCommand::WatchLevels { .. }
351 ) {
352 return Err("not a command".into());
353 }
354 self.send(Some(username), Some(to), command)
355 }
356
357 pub fn set_activity(
360 &self,
361 username: &str,
362 watcher: &str,
363 activity: Option<(String, String, bool)>,
364 ) {
365 let mut activities = self.activities.lock();
366 activities.retain(|a| !(a.username == username && a.watcher == watcher));
367 let Some((token, target, sandbox)) = activity else {
368 return;
369 };
370 activities.push(Activity {
371 username: username.to_string(),
372 watcher: watcher.to_string(),
373 target: target.clone(),
374 token,
375 sandbox,
376 sent: None,
377 });
378 drop(activities);
379 self.update_activities(username, &target);
380 }
381
382 fn update_activities(&self, username: &str, target: &str) {
385 let Some(pusher) = crate::push::pusher() else {
386 return;
387 };
388 let Some(info) = self
389 .list(Some(username))
390 .into_iter()
391 .find(|c| c.device == target)
392 else {
393 return;
394 };
395 let state = crate::push::ActivityState::of(&info);
396 let mut due = Vec::new();
397 for a in self.activities.lock().iter_mut() {
398 if a.username == username
399 && a.target == target
400 && a.sent.as_ref().is_none_or(|s| s.differs(&state))
401 {
402 a.sent = Some(state.clone());
403 due.push((a.token.clone(), a.sandbox, a.watcher.clone()));
404 }
405 }
406 if due.is_empty() {
407 return;
408 }
409 let username = username.to_string();
410 std::thread::spawn(move || {
411 for (token, sandbox, watcher) in due {
412 let push = crate::push::Push::Activity(state.clone());
413 match pusher.send(&token, sandbox, &push) {
414 crate::push::Outcome::Sent => {}
415 crate::push::Outcome::Gone => {
416 log::info!("push: a Live Activity on {watcher} has ended");
417 registry().set_activity(&username, &watcher, None);
418 }
419 crate::push::Outcome::Failed(e) => {
420 log::warn!("push: Live Activity on {watcher}: {e}");
421 }
422 }
423 }
424 });
425 }
426
427 pub fn list(&self, username: Option<&str>) -> Vec<ClientInfo> {
429 let mut out: Vec<ClientInfo> = self
430 .entries
431 .lock()
432 .iter()
433 .filter(|e| username.is_none_or(|u| e.info.username == u))
434 .map(|e| e.info.clone())
435 .collect();
436 out.sort_by_key(|c| std::cmp::Reverse(c.connected_at));
437 out
438 }
439
440 pub fn send(
445 &self,
446 username: Option<&str>,
447 id: Option<&str>,
448 cmd: LinkCommand,
449 ) -> Result<ClientInfo, String> {
450 let clients = self.list(username);
451 let target = match id {
452 Some(id) => clients
453 .iter()
454 .find(|c| c.id == id || c.device == id || c.name.eq_ignore_ascii_case(id)),
455 None if clients.is_empty() => None,
456 None => Some(pick(&clients, chrono::Utc::now().timestamp())?),
457 };
458 let Some(target) = target else {
460 return reach_absent(username, id, &cmd).unwrap_or_else(|| {
461 Err(match id {
462 Some(id) => format!("no linked client {id}; see `clients`"),
463 None => "no koan app is linked to this server; open koan on the device".into(),
464 })
465 });
466 };
467 let entries = self.entries.lock();
468 let entry = entries
469 .iter()
470 .find(|e| e.info.id == target.id)
471 .ok_or("that client has just gone")?;
472 entry
473 .tx
474 .send(cmd)
475 .map_err(|_| "that client has just gone".to_string())?;
476 Ok(target.clone())
477 }
478}
479
480impl Registry {
481 pub fn add_order(&self, order: Order) {
482 outbox::save_order(&order);
483 self.orders.lock().push(order);
484 }
485
486 pub fn orders(&self, username: Option<&str>) -> Vec<Order> {
487 self.orders
488 .lock()
489 .iter()
490 .filter(|o| username.is_none() || o.username.as_deref() == username)
491 .cloned()
492 .collect()
493 }
494
495 pub fn cancel_order(&self, username: Option<&str>, id: &str) -> bool {
496 let mut orders = self.orders.lock();
497 let before = orders.len();
498 orders
499 .retain(|o| !(o.id == id && (username.is_none() || o.username.as_deref() == username)));
500 let gone = orders.len() != before;
501 if gone {
502 outbox::drop_order(id);
503 }
504 gone
505 }
506
507 fn done(&self, id: &str) {
508 self.orders.lock().retain(|o| o.id != id);
509 outbox::drop_order(id);
510 }
511
512 pub fn fulfil_orders(&self, find: impl Fn(&Order) -> Option<Vec<String>>) {
515 let now = chrono::Utc::now().timestamp();
516 let pending: Vec<Order> = {
517 let mut orders = self.orders.lock();
518 for o in orders.iter().filter(|o| now - o.created_at >= ORDER_TTL) {
519 outbox::drop_order(&o.id);
520 }
521 orders.retain(|o| now - o.created_at < ORDER_TTL);
522 orders.clone()
523 };
524 for order in pending.into_iter().filter(|o| o.playlist.is_none()) {
525 let Some(ids) = find(&order).filter(|ids| !ids.is_empty()) else {
526 continue;
527 };
528 let track_ids = ids;
529 let cmd = if order.play_next {
530 LinkCommand::PlayNext { track_ids }
531 } else {
532 LinkCommand::Enqueue { track_ids }
533 };
534 match self.send(order.username.as_deref(), order.client.as_deref(), cmd) {
535 Ok(c) => {
536 log::info!(
537 "link: {} — {} arrived; queued on {}",
538 order.artist,
539 order.album,
540 c.name
541 );
542 self.done(&order.id);
543 }
544 Err(e) => log::info!("link: {} — {} arrived but {e}", order.artist, order.album),
546 }
547 }
548 }
549}
550
551pub fn fulfil_from(db_path: &std::path::Path) {
553 let registry = registry();
554 if registry.orders.lock().is_empty() {
555 return;
556 }
557 let Ok(db) = koan_core::db::connection::Database::open_existing(db_path) else {
558 return;
559 };
560 let for_playlists: Vec<Order> = registry
563 .orders
564 .lock()
565 .iter()
566 .filter(|o| o.playlist.is_some())
567 .cloned()
568 .collect();
569 let mut edited = false;
570 for order in for_playlists {
571 let Some(playlist) = order.playlist else {
572 continue;
573 };
574 let Some(ids) = order_tracks(&db.conn, &order) else {
575 continue;
576 };
577 match koan_core::db::queries::add_tracks(&db.conn, playlist, &ids) {
578 Ok(_) => {
579 if !cfg!(test) {
581 koan_core::playlists::push_to_remote(playlist);
582 }
583 log::info!(
584 "link: {} — {} arrived; added {} tracks to playlist {playlist}",
585 order.artist,
586 order.album,
587 ids.len()
588 );
589 registry.done(&order.id);
590 edited = true;
591 }
592 Err(e) => log::warn!("link: could not add to playlist {playlist}: {e}"),
593 }
594 }
595 if edited {
596 changed();
597 }
598 registry.fulfil_orders(|order| {
599 let rows = order_tracks(&db.conn, order)?;
600 queries::uids_in_order(&db.conn, UidKind::Track, &rows).ok()
601 });
602}
603
604fn order_tracks(conn: &rusqlite::Connection, order: &Order) -> Option<Vec<i64>> {
607 let tracks = album_tracks(conn, &order.artist, &order.album)?;
608 if order.titles.is_empty() {
609 return Some(tracks.into_iter().map(|(id, _)| id).collect());
610 }
611 let picked: Vec<i64> = order
612 .titles
613 .iter()
614 .filter_map(|want| {
615 let want = want.to_lowercase();
616 tracks
617 .iter()
618 .find(|(_, t)| t.to_lowercase().contains(&want))
619 .map(|(id, _)| *id)
620 })
621 .collect();
622 (!picked.is_empty()).then_some(picked)
623}
624
625pub fn album_tracks(
628 conn: &rusqlite::Connection,
629 artist: &str,
630 album: &str,
631) -> Option<Vec<(i64, String)>> {
632 let like = |s: &str| format!("%{}%", s.replace(['%', '_'], ""));
633 let album_id: i64 = conn
634 .query_row(
635 "SELECT al.id FROM albums al JOIN artists a ON a.id = al.artist_id
636 WHERE a.name LIKE ?1 COLLATE NOCASE AND al.title LIKE ?2 COLLATE NOCASE
637 ORDER BY al.id DESC LIMIT 1",
638 [like(artist), like(album)],
639 |r| r.get(0),
640 )
641 .ok()?;
642 let mut stmt = conn
643 .prepare("SELECT id, title FROM tracks WHERE album_id = ?1 ORDER BY disc, track_number, id")
644 .ok()?;
645 let tracks = stmt
646 .query_map([album_id], |r| Ok((r.get(0)?, r.get(1)?)))
647 .ok()?
648 .filter_map(Result::ok)
649 .collect();
650 Some(tracks)
651}
652
653impl Registry {
654 pub fn broadcast(&self, username: Option<&str>, cmd: LinkCommand) -> Vec<String> {
657 let ids: Vec<String> = self.list(username).into_iter().map(|c| c.id).collect();
658 let entries = self.entries.lock();
659 entries
660 .iter()
661 .filter(|e| ids.contains(&e.info.id) && e.tx.send(cmd.clone()).is_ok())
662 .map(|e| e.info.name.clone())
663 .collect()
664 }
665}
666
667impl Registry {
668 pub fn deliver(&self, username: Option<&str>, cmd: LinkCommand) -> (Vec<String>, Vec<String>) {
673 let (sent, queued) = self.link_or_queue(username, &cmd);
674 wake(&queued.iter().map(Absent::key).collect::<Vec<_>>());
675 (sent, queued.into_iter().map(|q| q.name).collect())
676 }
677
678 fn link_or_queue(
679 &self,
680 username: Option<&str>,
681 cmd: &LinkCommand,
682 ) -> (Vec<String>, Vec<Absent>) {
683 let sent = self.broadcast(username, cmd.clone());
684 let queued = outbox::queue_for_absent(username, &self.live(), cmd);
685 (sent, queued)
686 }
687
688 fn live(&self) -> Vec<(String, String)> {
690 self.entries
691 .lock()
692 .iter()
693 .map(|e| (e.device.clone(), e.info.username.clone()))
694 .collect()
695 }
696}
697
698impl Absent {
699 fn key(&self) -> (String, String) {
700 (self.device.clone(), self.username.clone())
701 }
702}
703
704fn wake(devices: &[(String, String)]) {
707 let Some(pusher) = crate::push::pusher() else {
708 return;
709 };
710 let targets: Vec<outbox::PushTarget> = outbox::push_targets(None)
711 .into_iter()
712 .filter(|t| {
713 devices
714 .iter()
715 .any(|(device, username)| *device == t.device && *username == t.username)
716 })
717 .collect();
718 if targets.is_empty() {
719 return;
720 }
721 std::thread::spawn(move || {
722 for t in targets {
723 deliver_push(pusher, &t, &crate::push::Push::Wake);
724 }
725 });
726}
727
728fn reach_absent(
737 username: Option<&str>,
738 id: Option<&str>,
739 cmd: &LinkCommand,
740) -> Option<Result<ClientInfo, String>> {
741 let pusher = crate::push::pusher()?;
742 let target = outbox::push_targets(username)
745 .into_iter()
746 .find(|t| id.is_none_or(|id| t.device == id || t.name.eq_ignore_ascii_case(id)))?;
747 let info = ClientInfo {
748 id: target.device.clone(),
749 device: target.device.clone(),
750 name: target.name.clone(),
751 platform: target.platform.clone(),
752 username: target.username.clone(),
753 connected_at: 0,
754 state: LinkState::default(),
755 last_played_at: None,
756 state_at: 0,
757 reports: false,
758 notified: true,
759 };
760 let verb = match cmd {
761 LinkCommand::Play { .. } | LinkCommand::JumpTo { .. } | LinkCommand::Resume => Some("Play"),
762 _ => None,
763 };
764 let push = match verb {
765 Some(verb) => crate::push::Push::Notify {
766 title: format!("{verb} on {}", target.name),
767 body: outbox::describe(cmd).unwrap_or_else(|| "From your koan server".into()),
768 command: serde_json::to_value(cmd).ok()?,
769 image: cover_track(cmd)
770 .and_then(outbox::track_row)
771 .and_then(|t| pusher.cover_link(t)),
772 },
773 None => {
774 outbox::queue_for(&target.device, &target.username, cmd);
775 crate::push::Push::Wake
776 }
777 };
778 std::thread::spawn(move || deliver_push(pusher, &target, &push));
779 Some(Ok(info))
780}
781
782fn cover_track(cmd: &LinkCommand) -> Option<&str> {
785 match cmd {
786 LinkCommand::Play {
787 track_ids,
788 start_at,
789 ..
790 } => track_ids
791 .get(*start_at as usize)
792 .or(track_ids.first())
793 .map(String::as_str),
794 LinkCommand::JumpTo { track_id } => Some(track_id),
795 _ => None,
796 }
797}
798
799fn deliver_push(
801 pusher: &crate::push::Pusher,
802 target: &outbox::PushTarget,
803 push: &crate::push::Push,
804) {
805 use crate::push::Outcome;
806 match pusher.send(&target.token, target.sandbox, push) {
807 Outcome::Sent => log::info!("push: sent to {}", target.name),
808 Outcome::Gone => {
809 log::info!(
810 "push: {}'s token is no longer valid; forgotten",
811 target.name
812 );
813 outbox::forget_push(&target.username, &target.device);
814 }
815 Outcome::Failed(e) => log::warn!("push: to {} failed: {e}", target.name),
816 }
817}
818
819pub fn changed() {
824 let (_, queued) = registry().link_or_queue(None, &LinkCommand::Sync { full: false });
825 if queued.is_empty() || crate::push::pusher().is_none() {
826 return;
827 }
828 let mut wakes = WAKES.lock();
829 wakes.add(queued.iter().map(Absent::key), std::time::Instant::now());
830 if !wakes.timer {
831 wakes.timer = true;
832 std::thread::spawn(send_wakes);
833 }
834}
835
836const QUIET: std::time::Duration = std::time::Duration::from_secs(30);
840
841const LONGEST_WAIT: std::time::Duration = std::time::Duration::from_secs(5 * 60);
843
844static WAKES: Mutex<Wakes> = Mutex::new(Wakes {
845 pending: Vec::new(),
846 timer: false,
847});
848
849struct Wakes {
852 pending: Vec<((String, String), std::time::Instant, std::time::Instant)>,
854 timer: bool,
856}
857
858impl Wakes {
859 fn add(
860 &mut self,
861 devices: impl IntoIterator<Item = (String, String)>,
862 now: std::time::Instant,
863 ) {
864 for key in devices {
865 match self.pending.iter_mut().find(|(k, _, _)| *k == key) {
866 Some((_, _, last)) => *last = now,
867 None => self.pending.push((key, now, now)),
868 }
869 }
870 }
871
872 fn due_at(first: std::time::Instant, last: std::time::Instant) -> std::time::Instant {
873 (last + QUIET).min(first + LONGEST_WAIT)
874 }
875
876 fn next(&self) -> Option<std::time::Instant> {
878 self.pending
879 .iter()
880 .map(|(_, first, last)| Self::due_at(*first, *last))
881 .min()
882 }
883
884 fn take_due(&mut self, now: std::time::Instant) -> Vec<(String, String)> {
886 let (due, waiting) = std::mem::take(&mut self.pending)
887 .into_iter()
888 .partition(|(_, first, last)| Self::due_at(*first, *last) <= now);
889 self.pending = waiting;
890 due.into_iter().map(|(key, _, _)| key).collect()
891 }
892}
893
894fn send_wakes() {
897 loop {
898 let (due, next) = {
899 let mut wakes = WAKES.lock();
900 let due = wakes.take_due(std::time::Instant::now());
901 let next = wakes.next();
902 if due.is_empty() && next.is_none() {
903 wakes.timer = false;
904 return;
905 }
906 (due, next)
907 };
908 let live = registry().live();
909 let absent: Vec<(String, String)> = due.into_iter().filter(|d| !live.contains(d)).collect();
910 if !absent.is_empty() {
911 wake(&absent);
912 }
913 if let Some(next) = next {
914 std::thread::sleep(next.saturating_duration_since(std::time::Instant::now()));
915 }
916 }
917}
918
919pub fn changed_if_library_moved(conn: &rusqlite::Connection) {
923 static LAST: parking_lot::Mutex<Option<(i64, i64, i64)>> = parking_lot::Mutex::new(None);
924 let Ok(now) = conn.query_row(
925 "SELECT (SELECT COUNT(*) FROM tracks), (SELECT COALESCE(MAX(id), 0) FROM tracks),
926 (SELECT COUNT(*) FROM albums)",
927 [],
928 |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
929 ) else {
930 return;
931 };
932 let before = LAST.lock().replace(now);
933 if before.is_some_and(|b| b != now) {
934 changed();
935 }
936}
937
938mod outbox {
941 use koan_core::db::queries;
942 use koan_core::remote::link::LinkCommand;
943
944 const KEEP_SECS: i64 = 30 * 24 * 60 * 60;
949
950 fn db() -> Option<koan_core::db::pool::Handle<'static>> {
955 if cfg!(test) {
956 return None;
957 }
958 koan_core::db::pool::shared().get().ok()
959 }
960
961 pub fn load_orders() -> Vec<super::Order> {
962 let Some(db) = db() else { return Vec::new() };
963 db.conn
964 .prepare("SELECT body FROM link_orders ORDER BY created_at")
965 .and_then(|mut s| {
966 s.query_map([], |r| r.get::<_, String>(0))?
967 .collect::<Result<Vec<_>, _>>()
968 })
969 .unwrap_or_default()
970 .into_iter()
971 .filter_map(|b| serde_json::from_str(&b).ok())
972 .collect()
973 }
974
975 pub fn save_order(order: &super::Order) {
976 let (Some(db), Ok(body)) = (db(), serde_json::to_string(order)) else {
977 return;
978 };
979 let _ = db.conn.execute(
980 "INSERT OR REPLACE INTO link_orders (id, body, created_at) VALUES (?1, ?2, ?3)",
981 rusqlite::params![order.id, body, order.created_at],
982 );
983 }
984
985 pub fn drop_order(id: &str) {
986 if let Some(db) = db() {
987 let _ = db
988 .conn
989 .execute("DELETE FROM link_orders WHERE id = ?1", [id]);
990 }
991 }
992
993 pub fn take_and_remember(
994 username: &str,
995 device: &str,
996 name: &str,
997 platform: &str,
998 live: &[(String, String)],
999 ) -> Vec<LinkCommand> {
1000 let Some(db) = db() else { return Vec::new() };
1001 let now = chrono::Utc::now().timestamp();
1002 let _ = db.conn.execute(
1003 "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES (?1, ?2, ?3, ?4, ?5)
1004 ON CONFLICT (device, username) DO UPDATE SET name = ?3, platform = ?4, last_seen = ?5",
1005 rusqlite::params![device, username, name, platform, now],
1006 );
1007 forget_stale(&db.conn, live, now);
1008 let waiting: Vec<(i64, String)> = db
1009 .conn
1010 .prepare("SELECT id, command FROM link_outbox WHERE device = ?1 AND username = ?2 ORDER BY id")
1011 .and_then(|mut s| {
1012 s.query_map([device, username], |r| Ok((r.get(0)?, r.get(1)?)))?
1013 .collect()
1014 })
1015 .unwrap_or_default();
1016 let _ = db.conn.execute(
1017 "DELETE FROM link_outbox WHERE device = ?1 AND username = ?2",
1018 [device, username],
1019 );
1020 if !waiting.is_empty() {
1021 log::info!("link: {} waiting commands for {name}", waiting.len());
1022 }
1023 waiting
1024 .into_iter()
1025 .filter_map(|(_, c)| serde_json::from_str(&c).ok())
1026 .collect()
1027 }
1028
1029 pub(super) fn forget_stale(conn: &rusqlite::Connection, live: &[(String, String)], now: i64) {
1032 touch_with(conn, live, now);
1033 let _ = conn.execute(
1034 "DELETE FROM link_outbox WHERE created_at < ?1",
1035 [now - KEEP_SECS],
1036 );
1037 let _ = conn.execute(
1038 "DELETE FROM link_push WHERE (device, username) IN
1039 (SELECT device, username FROM link_devices WHERE last_seen < ?1)",
1040 [now - KEEP_SECS],
1041 );
1042 let _ = conn.execute(
1043 "DELETE FROM link_devices WHERE last_seen < ?1",
1044 [now - KEEP_SECS],
1045 );
1046 }
1047
1048 pub fn touch(devices: &[(String, String)]) {
1050 if let Some(db) = db() {
1051 touch_with(&db.conn, devices, chrono::Utc::now().timestamp());
1052 }
1053 }
1054
1055 fn touch_with(conn: &rusqlite::Connection, devices: &[(String, String)], now: i64) {
1056 for (device, username) in devices {
1057 let _ = conn.execute(
1058 "UPDATE link_devices SET last_seen = ?1 WHERE device = ?2 AND username = ?3",
1059 rusqlite::params![now, device, username],
1060 );
1061 }
1062 }
1063
1064 pub struct Absent {
1066 pub device: String,
1067 pub username: String,
1068 pub name: String,
1069 }
1070
1071 pub fn queue_for_absent(
1073 username: Option<&str>,
1074 live: &[(String, String)],
1075 cmd: &LinkCommand,
1076 ) -> Vec<Absent> {
1077 let Some(db) = db() else { return Vec::new() };
1078 let known: Vec<(String, String, String)> = db
1079 .conn
1080 .prepare("SELECT device, username, name FROM link_devices")
1081 .and_then(|mut s| {
1082 s.query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))?
1083 .collect()
1084 })
1085 .unwrap_or_default();
1086 let Ok(text) = serde_json::to_string(cmd) else {
1087 return Vec::new();
1088 };
1089 let absent: Vec<_> = known
1090 .into_iter()
1091 .filter(|(device, user, _)| {
1092 !username.is_some_and(|u| u != user)
1093 && !live.iter().any(|(d, u)| d == device && u == user)
1094 })
1095 .collect();
1096 if absent.is_empty() {
1097 return Vec::new();
1098 }
1099 let is_sync = matches!(cmd, LinkCommand::Sync { .. });
1100 let now = chrono::Utc::now().timestamp();
1101 koan_core::db::queries::atomically(&db.conn, || {
1104 let mut queued = Vec::new();
1105 for (device, user, name) in absent {
1106 if is_sync {
1107 let _ = db.conn.execute(
1109 "DELETE FROM link_outbox WHERE device = ?1 AND username = ?2 AND command LIKE '{\"type\":\"sync\"%'",
1110 [&device, &user],
1111 );
1112 }
1113 if db
1114 .conn
1115 .execute(
1116 "INSERT INTO link_outbox (device, username, command, created_at) VALUES (?1, ?2, ?3, ?4)",
1117 rusqlite::params![device, user, text, now],
1118 )
1119 .is_ok()
1120 {
1121 queued.push(Absent {
1122 device,
1123 username: user,
1124 name,
1125 });
1126 }
1127 }
1128 Ok::<_, rusqlite::Error>(queued)
1129 })
1130 .unwrap_or_default()
1131 }
1132
1133 #[derive(Clone)]
1135 pub struct PushTarget {
1136 pub device: String,
1137 pub username: String,
1138 pub name: String,
1139 pub platform: String,
1140 pub token: String,
1141 pub sandbox: bool,
1142 }
1143
1144 pub fn queue_for(device: &str, username: &str, cmd: &LinkCommand) {
1146 let (Some(db), Ok(text)) = (db(), serde_json::to_string(cmd)) else {
1147 return;
1148 };
1149 let _ = db.conn.execute(
1150 "INSERT INTO link_outbox (device, username, command, created_at) VALUES (?1, ?2, ?3, ?4)",
1151 rusqlite::params![device, username, text, chrono::Utc::now().timestamp()],
1152 );
1153 }
1154
1155 pub fn save_push(username: &str, device: &str, token: &str, sandbox: bool) {
1156 let Some(db) = db() else { return };
1157 let _ = db.conn.execute(
1158 "INSERT INTO link_push (device, username, token, sandbox, updated_at) VALUES (?1, ?2, ?3, ?4, ?5)
1159 ON CONFLICT (device, username) DO UPDATE SET token = ?3, sandbox = ?4, updated_at = ?5",
1160 rusqlite::params![device, username, token, sandbox, chrono::Utc::now().timestamp()],
1161 );
1162 }
1163
1164 pub fn forget_push(username: &str, device: &str) {
1165 if let Some(db) = db() {
1166 let _ = db.conn.execute(
1167 "DELETE FROM link_push WHERE device = ?1 AND username = ?2",
1168 [device, username],
1169 );
1170 }
1171 }
1172
1173 pub fn push_targets(username: Option<&str>) -> Vec<PushTarget> {
1175 let Some(db) = db() else { return Vec::new() };
1176 db.conn
1177 .prepare(
1178 "SELECT p.device, p.username, d.name, d.platform, p.token, p.sandbox
1179 FROM link_push p JOIN link_devices d ON d.device = p.device AND d.username = p.username
1180 WHERE ?1 IS NULL OR p.username = ?1
1181 ORDER BY d.last_seen DESC",
1182 )
1183 .and_then(|mut s| {
1184 s.query_map([username], |r| {
1185 Ok(PushTarget {
1186 device: r.get(0)?,
1187 username: r.get(1)?,
1188 name: r.get(2)?,
1189 platform: r.get(3)?,
1190 token: r.get(4)?,
1191 sandbox: r.get(5)?,
1192 })
1193 })?
1194 .collect()
1195 })
1196 .unwrap_or_default()
1197 }
1198
1199 pub fn track_row(id: &str) -> Option<i64> {
1201 let db = db()?;
1202 queries::resolve_id(&db.conn, queries::UidKind::Track, id)
1203 .ok()
1204 .flatten()
1205 }
1206
1207 pub fn describe(cmd: &LinkCommand) -> Option<String> {
1210 let ids: Vec<&String> = match cmd {
1211 LinkCommand::Play { track_ids, .. }
1212 | LinkCommand::Enqueue { track_ids }
1213 | LinkCommand::PlayNext { track_ids } => track_ids.iter().collect(),
1214 LinkCommand::JumpTo { track_id } => vec![track_id],
1215 _ => return None,
1216 };
1217 let db = db()?;
1218 let ids: Vec<i64> = ids
1219 .into_iter()
1220 .filter_map(|t| {
1221 queries::resolve_id(&db.conn, queries::UidKind::Track, t)
1222 .ok()
1223 .flatten()
1224 })
1225 .collect();
1226 let row = |id: i64| {
1227 db.conn
1228 .query_row(
1229 "SELECT t.title, COALESCE(a.name, ''), COALESCE(al.title, ''), t.album_id
1230 FROM tracks t LEFT JOIN artists a ON a.id = t.artist_id
1231 LEFT JOIN albums al ON al.id = t.album_id WHERE t.id = ?1",
1232 [id],
1233 |r| {
1234 Ok((
1235 r.get::<_, String>(0)?,
1236 r.get::<_, String>(1)?,
1237 r.get::<_, String>(2)?,
1238 r.get::<_, Option<i64>>(3)?,
1239 ))
1240 },
1241 )
1242 .ok()
1243 };
1244 let (title, artist, album, album_id) = row(*ids.first()?)?;
1245 let one_album = ids.len() > 1
1246 && album_id.is_some()
1247 && ids
1248 .iter()
1249 .all(|id| row(*id).is_some_and(|r| r.3 == album_id));
1250 Some(match (one_album, ids.len()) {
1251 (true, _) => format!("{album} — {artist}"),
1252 (false, 1) => format!("{title} — {artist}"),
1253 (false, n) => format!("{title} — {artist}, and {} more", n - 1),
1254 })
1255 }
1256}
1257
1258const RECENT: i64 = 6 * 60 * 60;
1261
1262fn pick(clients: &[ClientInfo], now: i64) -> Result<&ClientInfo, String> {
1263 if let Some(c) = clients.iter().find(|c| c.state.playing) {
1264 return Ok(c);
1265 }
1266 if let Some(c) = clients
1267 .iter()
1268 .filter(|c| c.last_played_at.is_some_and(|t| now - t < RECENT))
1269 .max_by_key(|c| c.last_played_at)
1270 {
1271 return Ok(c);
1272 }
1273 match clients {
1274 [] => Err("no koan app is linked to this server; open koan on the device".into()),
1275 [only] => Ok(only),
1276 several => Err(format!(
1277 "several koan apps are linked and none has played recently: {}. Ask which, then pass `client`",
1278 several
1279 .iter()
1280 .map(|c| format!("{} ({}, id {})", c.name, c.platform, c.device))
1281 .collect::<Vec<_>>()
1282 .join(", ")
1283 )),
1284 }
1285}
1286
1287#[cfg(test)]
1288mod tests {
1289
1290 #[test]
1291 fn a_watch_is_never_queued_and_is_renewed_when_the_target_relinks() {
1292 let reg = Registry::default();
1293 let watches = |rx: &mut tokio::sync::mpsc::UnboundedReceiver<LinkCommand>| {
1294 std::iter::from_fn(|| rx.try_recv().ok())
1295 .filter(|c| matches!(c, LinkCommand::WatchLevels { .. }))
1296 .collect::<Vec<_>>()
1297 };
1298 assert!(
1300 reg.relay("rl", "rl-phone", LinkCommand::WatchLevels { on: true })
1301 .is_err()
1302 );
1303 reg.watch_levels("rl", "rl-mac", "rl-phone", true);
1305 let (tx, mut phone) = tokio::sync::mpsc::unbounded_channel();
1306 reg.register("rl", "phone", "ios", "rl-phone", tx, false);
1307 assert_eq!(
1308 watches(&mut phone),
1309 vec![LinkCommand::WatchLevels { on: true }],
1310 "told once, on linking, because it is watched now"
1311 );
1312 let (tx, mut phone) = tokio::sync::mpsc::unbounded_channel();
1314 reg.register("rl", "phone", "ios", "rl-phone", tx, false);
1315 assert_eq!(
1316 watches(&mut phone),
1317 vec![LinkCommand::WatchLevels { on: true }]
1318 );
1319 }
1320
1321 #[test]
1322 fn levels_reach_a_watcher_only_while_it_watches() {
1323 use koan_core::remote::levels::Frame;
1324 let reg = Registry::default();
1325 let (tx_mac, mut mac) = tokio::sync::mpsc::unbounded_channel();
1326 let (tx_phone, mut phone) = tokio::sync::mpsc::unbounded_channel();
1327 let mac_id = reg.register("lv", "mac", "macos", "lv-mac", tx_mac, false);
1328 reg.register("lv", "phone", "ios", "lv-phone", tx_phone, false);
1329 let levels = |rx: &mut tokio::sync::mpsc::UnboundedReceiver<LinkCommand>| {
1330 std::iter::from_fn(|| rx.try_recv().ok())
1331 .filter(|c| {
1332 matches!(
1333 c,
1334 LinkCommand::Levels { .. } | LinkCommand::WatchLevels { .. }
1335 )
1336 })
1337 .collect::<Vec<_>>()
1338 };
1339 let f = Frame(1_000, 1, 2, 3);
1340
1341 reg.levels("lv", "lv-phone", f);
1342 assert!(levels(&mut mac).is_empty(), "nobody watching");
1343
1344 reg.watch_levels("lv", "lv-mac", "lv-phone", true);
1345 assert_eq!(
1346 levels(&mut phone),
1347 vec![LinkCommand::WatchLevels { on: true }]
1348 );
1349 reg.levels("lv", "lv-phone", f);
1350 assert_eq!(
1351 levels(&mut mac),
1352 vec![LinkCommand::Levels {
1353 from: "lv-phone".into(),
1354 f
1355 }]
1356 );
1357
1358 reg.unregister(&mac_id);
1360 assert_eq!(
1361 levels(&mut phone),
1362 vec![LinkCommand::WatchLevels { on: false }]
1363 );
1364 reg.watch_levels("lv", "lv-mac", "lv-phone", true);
1365 reg.watch_levels("lv", "lv-mac", "lv-phone", false);
1366 assert_eq!(
1367 levels(&mut phone),
1368 vec![
1369 LinkCommand::WatchLevels { on: true },
1370 LinkCommand::WatchLevels { on: false }
1371 ]
1372 );
1373 }
1374
1375 #[test]
1376 fn a_playlist_order_adds_the_named_tracks_once_they_arrive() {
1377 let dir = tempfile::tempdir().unwrap();
1378 let path = dir.path().join("koan.db");
1379 let db = koan_core::db::connection::Database::open(&path).unwrap();
1380 let playlist = koan_core::db::queries::create_playlist(
1381 &db.conn,
1382 koan_core::db::queries::LOCAL_USER,
1383 "cyberpunk",
1384 None,
1385 )
1386 .unwrap();
1387 let order = Order {
1388 id: "o1".into(),
1389 username: None,
1390 client: None,
1391 artist: "Perturbator".into(),
1392 album: "Dangerous Days".into(),
1393 play_next: false,
1394 playlist: Some(playlist),
1395 titles: vec!["Future Club".into()],
1396 created_at: chrono::Utc::now().timestamp(),
1397 };
1398 registry().add_order(order);
1399
1400 fulfil_from(&path);
1402 assert!(registry().orders(None).iter().any(|o| o.id == "o1"));
1403
1404 db.conn
1405 .execute_batch(
1406 "INSERT INTO artists (id, name) VALUES (1, 'Perturbator');
1407 INSERT INTO albums (id, title, artist_id) VALUES (1, 'Dangerous Days', 1);
1408 INSERT INTO tracks (id, title, album_id, artist_id, track_number, path) VALUES
1409 (1, 'Welcome Back', 1, 1, 1, '/1.flac'), (2, 'Future Club', 1, 1, 2, '/2.flac');",
1410 )
1411 .unwrap();
1412 fulfil_from(&path);
1413 let held: Vec<i64> = db
1414 .conn
1415 .prepare("SELECT track_id FROM playlist_tracks WHERE playlist_id = ?1")
1416 .unwrap()
1417 .query_map([playlist], |r| r.get(0))
1418 .unwrap()
1419 .collect::<Result<_, _>>()
1420 .unwrap();
1421 assert_eq!(held, [2]);
1422 assert!(!registry().orders(None).iter().any(|o| o.id == "o1"));
1423 }
1424
1425 use super::*;
1426
1427 #[test]
1428 fn a_notification_cover_is_named_by_the_uid_the_command_carries() {
1429 let uid = "0190a5b2-7c3d-7e4f-8a1b-2c3d4e5f6a7b".to_string();
1431 let play = LinkCommand::Play {
1432 track_ids: vec!["x".into(), uid.clone()],
1433 start_at: 1,
1434 position_ms: 0,
1435 paused: false,
1436 handoff: false,
1437 };
1438 assert_eq!(cover_track(&play), Some(uid.as_str()));
1439 }
1440
1441 #[test]
1442 fn a_reconnect_replaces_the_device_and_commands_reach_it() {
1443 let reg = Registry::default();
1444 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
1445 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
1446 let (tx3, _rx3) = tokio::sync::mpsc::unbounded_channel();
1447 reg.register("j", "phone", "ios", "dev-1", tx1, false);
1448 let id = reg.register("j", "phone", "ios", "dev-1", tx2, false);
1449 reg.register("someone", "laptop", "macos", "dev-2", tx3, false);
1450
1451 assert_eq!(reg.list(Some("j")).len(), 1);
1452 assert_eq!(reg.list(None).len(), 2);
1453
1454 let sent = reg.send(Some("j"), None, LinkCommand::Pause).unwrap();
1455 assert_eq!(sent.id, id);
1456 assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
1457
1458 assert!(
1460 reg.send(Some("j"), Some("laptop"), LinkCommand::Pause)
1461 .is_err()
1462 );
1463
1464 reg.unregister(&id);
1465 assert!(reg.send(Some("j"), None, LinkCommand::Pause).is_err());
1466 }
1467
1468 #[test]
1469 fn disconnecting_an_account_closes_only_its_links() {
1470 let reg = Registry::default();
1471 let (tx1, mut rx1) = tokio::sync::mpsc::unbounded_channel();
1472 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
1473 reg.register("j", "phone", "ios", "dev-1", tx1, false);
1474 reg.register("someone", "laptop", "macos", "dev-2", tx2, false);
1475
1476 reg.disconnect("j");
1477 assert!(reg.list(Some("j")).is_empty());
1478 assert!(matches!(
1480 rx1.try_recv(),
1481 Err(tokio::sync::mpsc::error::TryRecvError::Disconnected)
1482 ));
1483 assert_eq!(reg.list(Some("someone")).len(), 1);
1484 assert!(matches!(
1485 rx2.try_recv(),
1486 Err(tokio::sync::mpsc::error::TryRecvError::Empty)
1487 ));
1488 }
1489
1490 #[test]
1491 fn a_burst_of_changes_wakes_each_device_once_when_it_goes_quiet() {
1492 let t0 = std::time::Instant::now();
1493 let s = std::time::Duration::from_secs;
1494 let phone = || ("dev-1".to_string(), "j".to_string());
1495 let ipad = || ("dev-2".to_string(), "j".to_string());
1496 let mut wakes = Wakes {
1497 pending: Vec::new(),
1498 timer: false,
1499 };
1500
1501 wakes.add([phone()], t0);
1502 wakes.add([phone(), ipad()], t0 + s(10));
1503 wakes.add([phone()], t0 + s(20));
1504 assert_eq!(wakes.pending.len(), 2);
1505
1506 assert!(wakes.take_due(t0 + s(39)).is_empty());
1508 assert_eq!(wakes.take_due(t0 + s(40)), [ipad()]);
1509 assert_eq!(wakes.next(), Some(t0 + s(50)));
1510 assert_eq!(wakes.take_due(t0 + s(50)), [phone()]);
1511 assert_eq!(wakes.next(), None);
1512 }
1513
1514 #[test]
1515 fn a_library_that_never_goes_quiet_still_wakes_devices() {
1516 let t0 = std::time::Instant::now();
1517 let phone = || ("dev-1".to_string(), "j".to_string());
1518 let mut wakes = Wakes {
1519 pending: Vec::new(),
1520 timer: false,
1521 };
1522 let mut sent = 0;
1523 for i in 0..40 {
1524 let now = t0 + std::time::Duration::from_secs(i * 10);
1525 sent += wakes.take_due(now).len();
1526 wakes.add([phone()], now);
1527 }
1528 assert_eq!(sent, 1);
1531 assert_eq!(wakes.pending.len(), 1);
1532 }
1533
1534 #[test]
1535 fn the_device_playing_is_the_one_meant() {
1536 let reg = Registry::default();
1537 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
1538 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
1539 let mac = reg.register("j", "mac", "macos", "dev-1", tx1, false);
1540 let phone = reg.register("j", "phone", "ios", "dev-2", tx2, false);
1541
1542 let err = reg.send(Some("j"), None, LinkCommand::Pause).unwrap_err();
1544 assert!(err.contains("mac") && err.contains("phone"), "{err}");
1545
1546 reg.report(
1547 &phone,
1548 LinkState {
1549 playing: true,
1550 ..Default::default()
1551 },
1552 );
1553 assert_eq!(
1554 reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
1555 phone
1556 );
1557 assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
1558
1559 reg.report(&phone, LinkState::default());
1561 assert_eq!(
1562 reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
1563 phone
1564 );
1565 let _ = mac;
1566 }
1567
1568 #[test]
1569 fn a_device_hears_its_peers_and_commands_reach_them_by_device_id() {
1570 let reg = Registry::default();
1571 let (tx1, mut rx1) = tokio::sync::mpsc::unbounded_channel();
1572 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
1573 let (tx3, mut rx3) = tokio::sync::mpsc::unbounded_channel();
1574 reg.register("j", "mac", "macos", "dev-mac", tx1, true);
1575 let phone = reg.register("j", "phone", "ios", "dev-phone", tx2, true);
1576 reg.register("someone", "laptop", "macos", "dev-other", tx3, true);
1577
1578 reg.report(
1579 &phone,
1580 LinkState {
1581 playing: true,
1582 title: Some("Roygbiv".into()),
1583 ..Default::default()
1584 },
1585 );
1586 let mut last = None;
1587 while let Ok(LinkCommand::Devices { devices }) = rx1.try_recv() {
1588 last = Some(devices);
1589 }
1590 let devices = last.expect("the Mac is told of the phone");
1591 assert_eq!(devices.len(), 1, "not itself, not another account's");
1592 assert_eq!(devices[0].id, "dev-phone");
1593 assert_eq!(
1594 devices[0].state.as_ref().and_then(|s| s.title.as_deref()),
1595 Some("Roygbiv")
1596 );
1597 while rx3.try_recv().is_ok() {}
1598 assert!(rx3.try_recv().is_err());
1599
1600 while rx2.try_recv().is_ok() {}
1601 reg.relay("j", "dev-phone", LinkCommand::Pause).unwrap();
1602 assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
1603 assert!(
1604 reg.relay("someone", "dev-phone", LinkCommand::Pause)
1605 .is_err(),
1606 "another account cannot reach it"
1607 );
1608 assert!(
1609 reg.relay("j", "dev-phone", LinkCommand::Devices { devices: vec![] })
1610 .is_err()
1611 );
1612 }
1613
1614 #[test]
1615 fn a_device_linked_for_a_month_is_not_forgotten() {
1616 let conn = rusqlite::Connection::open_in_memory().unwrap();
1617 koan_core::db::schema::create_tables(&conn).unwrap();
1618 let now = 100 * 24 * 60 * 60;
1619 let long_ago = now - 40 * 24 * 60 * 60;
1620 for device in ["mac", "old-phone"] {
1621 conn.execute(
1622 "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES (?1, 'j', ?1, 'ios', ?2)",
1623 rusqlite::params![device, long_ago],
1624 )
1625 .unwrap();
1626 conn.execute(
1627 "INSERT INTO link_push (device, username, token, sandbox, updated_at) VALUES (?1, 'j', 't', 0, ?2)",
1628 rusqlite::params![device, long_ago],
1629 )
1630 .unwrap();
1631 }
1632 outbox::forget_stale(&conn, &[("mac".into(), "j".into())], now);
1633 let left = |table: &str| -> Vec<String> {
1634 conn.prepare(&format!("SELECT device FROM {table} ORDER BY device"))
1635 .unwrap()
1636 .query_map([], |r| r.get(0))
1637 .unwrap()
1638 .collect::<Result<_, _>>()
1639 .unwrap()
1640 };
1641 assert_eq!(left("link_devices"), ["mac"]);
1642 assert_eq!(left("link_push"), ["mac"]);
1643 }
1644}