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