1use std::sync::LazyLock;
10
11use koan_core::db::queries::{self, UidKind};
12use koan_core::remote::acks::{AckOutcome, Envelope};
13use koan_core::remote::link::{LinkCommand, LinkDevice, LinkState};
14use outbox::Absent;
15use parking_lot::Mutex;
16use tokio::sync::mpsc::UnboundedSender;
17
18#[derive(Debug, Clone)]
20pub struct ClientInfo {
21 pub id: String,
22 pub device: String,
24 pub name: String,
25 pub platform: String,
26 pub username: String,
27 pub connected_at: i64,
29 pub state: LinkState,
31 pub last_played_at: Option<i64>,
33 pub state_at: i64,
35 pub reports: bool,
38 pub notified: bool,
41 pub acks: bool,
43}
44
45impl ClientInfo {
46 pub fn position_ms(&self) -> u64 {
48 let pos = self.state.position_ms;
49 if !self.state.playing {
50 return pos;
51 }
52 let run = (chrono::Utc::now().timestamp_millis() - self.state_at).max(0) as u64;
53 let pos = pos + run;
54 if self.state.duration_ms > 0 {
55 pos.min(self.state.duration_ms)
56 } else {
57 pos
58 }
59 }
60}
61
62struct Entry {
63 info: ClientInfo,
64 device: String,
66 tx: UnboundedSender<Envelope>,
67 wants_devices: bool,
70 wants_keys: bool,
73}
74
75pub struct Acking {
78 pub id: u64,
79 pub reply: Box<dyn FnOnce(Option<AckOutcome>) + Send>,
80}
81
82struct Waiting {
84 reply: Box<dyn FnOnce(Option<AckOutcome>) + Send>,
85 received: bool,
87 timer: Option<tokio::task::AbortHandle>,
90}
91
92impl Waiting {
93 fn answer(mut self, outcome: Option<AckOutcome>) {
94 let reply = std::mem::replace(&mut self.reply, Box::new(|_| {}));
95 reply(outcome);
96 }
97}
98
99impl Drop for Waiting {
100 fn drop(&mut self) {
101 if let Some(timer) = self.timer.take() {
102 timer.abort();
103 }
104 }
105}
106
107const FIRST_LOOK: std::time::Duration = std::time::Duration::from_millis(2500);
112const LONGEST_ANSWER: std::time::Duration = std::time::Duration::from_secs(120);
113
114struct Activity {
117 username: String,
118 watcher: String,
120 target: String,
122 token: String,
123 sandbox: bool,
124 sent: Option<crate::push::ActivityState>,
126}
127
128#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
131pub struct Order {
132 pub id: String,
133 pub username: Option<String>,
135 pub client: Option<String>,
137 pub artist: String,
138 pub album: String,
139 pub play_next: bool,
141 #[serde(default)]
143 pub playlist: Option<i64>,
144 #[serde(default)]
147 pub titles: Vec<String>,
148 pub created_at: i64,
150}
151
152const ORDER_TTL: i64 = 24 * 60 * 60;
154
155#[derive(Debug, Clone, PartialEq, Eq)]
159pub struct Grant {
160 pub device: String,
161 pub owner: String,
162 pub grantee: String,
163}
164
165#[derive(Default)]
166pub struct Registry {
167 entries: Mutex<Vec<Entry>>,
168 orders: Mutex<Vec<Order>>,
169 activities: Mutex<Vec<Activity>>,
170 grants: Mutex<Vec<Grant>>,
171 level_watches: Mutex<Vec<LevelWatch>>,
174 addresses: Mutex<std::collections::HashMap<String, (String, std::net::IpAddr)>>,
177 answers: Mutex<std::collections::HashMap<(String, u64), Waiting>>,
179}
180
181struct LevelWatch {
185 owner: String,
186 watcher_user: String,
187 watcher: String,
188 target: String,
189}
190
191pub fn registry() -> &'static Registry {
194 static REGISTRY: LazyLock<Registry> = LazyLock::new(|| {
195 let registry = Registry::default();
196 *registry.orders.lock() = outbox::load_orders();
197 *registry.grants.lock() = outbox::load_grants();
198 *registry.addresses.lock() = outbox::load_addresses();
199 registry
200 });
201 ®ISTRY
202}
203
204impl Registry {
205 #[allow(clippy::too_many_arguments)]
208 pub fn register(
209 &self,
210 username: &str,
211 name: &str,
212 platform: &str,
213 device: &str,
214 tx: UnboundedSender<Envelope>,
215 wants_devices: bool,
216 acks: bool,
217 ) -> String {
218 let id = uuid::Uuid::now_v7().to_string();
219 if let Some((sent, how)) = WOKEN
220 .lock()
221 .remove(&(device.to_string(), username.to_string()))
222 {
223 log::info!(
224 "wake: {name} linked {}ms after its {how}",
225 sent.elapsed().as_millis()
226 );
227 }
228 let live = self.live();
231 for envelope in outbox::take_and_remember(username, device, name, platform, &live) {
232 let _ = tx.send(envelope);
233 }
234 let mut entries = self.entries.lock();
235 entries.retain(|e| !(e.device == device && e.info.username == username));
236 entries.push(Entry {
237 info: ClientInfo {
238 id: id.clone(),
239 device: device.to_string(),
240 name: name.to_string(),
241 platform: platform.to_string(),
242 username: username.to_string(),
243 connected_at: chrono::Utc::now().timestamp(),
244 state: LinkState::default(),
245 last_played_at: None,
246 state_at: chrono::Utc::now().timestamp_millis(),
247 reports: false,
248 notified: false,
249 acks,
250 },
251 device: device.to_string(),
252 tx,
253 wants_devices,
254 wants_keys: false,
255 });
256 drop(entries);
257 self.send_shares(username, device, None);
260 self.announce(username);
261 if self
264 .level_watches
265 .lock()
266 .iter()
267 .any(|w| w.owner == username && w.target == device)
268 {
269 self.send_live(username, device, LinkCommand::WatchLevels { on: true });
270 }
271 id
272 }
273
274 pub fn wants_keys(&self, id: &str) {
276 if let Some(e) = self.entries.lock().iter_mut().find(|e| e.info.id == id) {
277 e.wants_keys = true;
278 }
279 }
280
281 pub fn publish_keys(&self, username: &str, keys: Vec<koan_core::remote::link::LinkDeviceKey>) {
284 for e in self
285 .entries
286 .lock()
287 .iter()
288 .filter(|e| e.info.username == username && e.wants_keys)
289 {
290 let _ =
291 e.tx.send(LinkCommand::DeviceKeys { keys: keys.clone() }.into());
292 }
293 }
294
295 pub fn keyed_accounts(&self) -> Vec<String> {
297 let mut accounts: Vec<String> = self
298 .entries
299 .lock()
300 .iter()
301 .filter(|e| e.wants_keys)
302 .map(|e| e.info.username.clone())
303 .collect();
304 accounts.sort();
305 accounts.dedup();
306 accounts
307 }
308
309 pub fn grantees_of(&self, owner: &str, device: &str) -> Vec<String> {
311 self.grants
312 .lock()
313 .iter()
314 .filter(|g| g.owner == owner && g.device == device)
315 .map(|g| g.grantee.clone())
316 .collect()
317 }
318
319 pub fn report(&self, id: &str, state: LinkState) {
321 let mut entries = self.entries.lock();
322 if let Some(e) = entries.iter_mut().find(|e| e.info.id == id) {
323 if state.playing || e.info.state.playing {
324 e.info.last_played_at = Some(chrono::Utc::now().timestamp());
325 }
326 e.info.state = state;
327 e.info.state_at = chrono::Utc::now().timestamp_millis();
328 e.info.reports = true;
329 let (username, device) = (e.info.username.clone(), e.device.clone());
330 drop(entries);
331 self.announce(&username);
332 self.update_activities(&username, &device);
333 }
334 }
335
336 pub fn set_push(&self, username: &str, device: &str, token: &str, sandbox: bool) {
338 outbox::save_push(username, device, token, sandbox);
339 }
340
341 pub fn disconnect(&self, username: &str) {
344 self.entries.lock().retain(|e| e.info.username != username);
345 }
346
347 pub fn unregister(&self, id: &str) {
348 let mut entries = self.entries.lock();
349 let gone = entries
350 .iter()
351 .find(|e| e.info.id == id)
352 .map(|e| (e.device.clone(), e.info.username.clone()));
353 entries.retain(|e| e.info.id != id);
354 drop(entries);
355 if let Some((device, username)) = gone {
356 outbox::touch(&[(device.clone(), username.clone())]);
357 self.announce(&username);
358 let targets: Vec<String> = self
360 .level_watches
361 .lock()
362 .iter()
363 .filter(|w| w.watcher_user == username && w.watcher == device)
364 .map(|w| w.target.clone())
365 .collect();
366 for target in targets {
367 self.watch_levels(&username, &device, &target, false);
368 }
369 }
370 }
371
372 pub fn watch_levels(&self, username: &str, watcher: &str, target: &str, on: bool) {
377 let owner = self
379 .shared_owner(username, target)
380 .unwrap_or_else(|| username.to_string());
381 let mut watches = self.level_watches.lock();
382 watches.retain(|w| {
383 !(w.watcher_user == username && w.watcher == watcher && w.target == target)
384 });
385 if on {
386 watches.push(LevelWatch {
387 owner: owner.clone(),
388 watcher_user: username.into(),
389 watcher: watcher.into(),
390 target: target.into(),
391 });
392 }
393 let watched = watches
394 .iter()
395 .any(|w| w.owner == owner && w.target == target);
396 drop(watches);
397 if on || !watched {
398 self.send_live(&owner, target, LinkCommand::WatchLevels { on: watched });
399 }
400 }
401
402 pub fn levels(&self, username: &str, from: &str, f: koan_core::remote::levels::Frame) {
405 let watchers: Vec<(String, String)> = self
406 .level_watches
407 .lock()
408 .iter()
409 .filter(|w| w.owner == username && w.target == from)
410 .map(|w| (w.watcher_user.clone(), w.watcher.clone()))
411 .collect();
412 for (watcher_user, watcher) in watchers {
413 self.send_live(
414 &watcher_user,
415 &watcher,
416 LinkCommand::Levels {
417 from: from.to_string(),
418 f,
419 },
420 );
421 }
422 }
423
424 pub(crate) fn send_live(&self, username: &str, device: &str, cmd: LinkCommand) {
426 if let Some(e) = self
427 .entries
428 .lock()
429 .iter()
430 .find(|e| e.info.username == username && e.device == device)
431 {
432 let _ = e.tx.send(cmd.into());
433 }
434 }
435
436 fn announce(&self, username: &str) {
440 self.announce_to(username);
441 let grantees: Vec<String> = {
443 let mut g: Vec<String> = self
444 .grants
445 .lock()
446 .iter()
447 .filter(|g| g.owner == username)
448 .map(|g| g.grantee.clone())
449 .collect();
450 g.sort();
451 g.dedup();
452 g
453 };
454 for grantee in grantees {
455 self.announce_to(&grantee);
456 }
457 }
458
459 fn announce_to(&self, username: &str) {
461 let listening = |e: &Entry| e.info.username == username && e.wants_devices;
462 if !self.entries.lock().iter().any(listening) {
463 return;
464 }
465 let asleep = outbox::push_targets(Some(username));
466 let can_push = crate::push::pusher().is_some();
468 let entries = self.entries.lock();
469 let ours: Vec<&Entry> = entries
470 .iter()
471 .filter(|e| e.info.username == username)
472 .collect();
473 if !ours.iter().any(|e| e.wants_devices) {
474 return;
475 }
476 let mut all: Vec<LinkDevice> = ours
477 .iter()
478 .map(|e| LinkDevice {
479 id: e.device.clone(),
480 name: e.info.name.clone(),
481 platform: e.info.platform.clone(),
482 linked: true,
483 state: e.info.reports.then(|| LinkState {
484 position_ms: e.info.position_ms(),
485 ..e.info.state.clone()
486 }),
487 last_seen: None,
488 wakeable: Some(can_push && asleep.iter().any(|t| t.device == e.device)),
489 owner: None,
490 acks: e.info.acks,
491 })
492 .collect();
493 for t in asleep {
494 if !all.iter().any(|d| d.id == t.device) {
495 all.push(LinkDevice {
496 id: t.device,
497 name: t.name,
498 platform: t.platform,
499 linked: false,
500 state: None,
501 last_seen: Some(t.last_seen),
502 wakeable: Some(can_push),
503 owner: None,
504 acks: false,
505 });
506 }
507 }
508 all.extend(self.shared_devices(username, &entries, can_push));
509 for e in ours.iter().filter(|e| e.wants_devices) {
510 let devices = all.iter().filter(|d| d.id != e.device).cloned().collect();
511 let _ = e.tx.send(LinkCommand::Devices { devices }.into());
512 }
513 }
514
515 fn shared_devices(&self, grantee: &str, entries: &[Entry], can_push: bool) -> Vec<LinkDevice> {
519 let grants: Vec<Grant> = self
520 .grants
521 .lock()
522 .iter()
523 .filter(|g| g.grantee == grantee)
524 .cloned()
525 .collect();
526 let mut out = Vec::new();
527 for g in grants {
528 let linked = entries
529 .iter()
530 .find(|e| e.device == g.device && e.info.username == g.owner);
531 match linked {
532 Some(e) => out.push(LinkDevice {
533 id: e.device.clone(),
534 name: e.info.name.clone(),
535 platform: e.info.platform.clone(),
536 linked: true,
537 state: e.info.reports.then(|| LinkState {
538 position_ms: e.info.position_ms(),
539 ..e.info.state.clone()
540 }),
541 last_seen: None,
542 wakeable: Some(can_push),
543 owner: Some(g.owner.clone()),
544 acks: e.info.acks,
545 }),
546 None => {
547 if let Some(t) = outbox::push_targets(Some(&g.owner))
548 .into_iter()
549 .find(|t| t.device == g.device)
550 {
551 out.push(LinkDevice {
552 id: t.device,
553 name: t.name,
554 platform: t.platform,
555 linked: false,
556 state: None,
557 last_seen: Some(t.last_seen),
558 wakeable: Some(can_push),
559 owner: Some(g.owner.clone()),
560 acks: false,
561 });
562 }
563 }
564 }
565 }
566 out
567 }
568
569 pub fn share(
573 &self,
574 owner: &str,
575 device: &str,
576 grantee: &str,
577 allow: bool,
578 ) -> Result<(), String> {
579 if grantee == owner {
580 return Err("a device is already its own account's".into());
581 }
582 let grant = Grant {
583 device: device.to_string(),
584 owner: owner.to_string(),
585 grantee: grantee.to_string(),
586 };
587 {
588 let mut grants = self.grants.lock();
589 grants.retain(|g| *g != grant);
590 if allow {
591 grants.push(grant.clone());
592 }
593 }
594 if allow {
595 outbox::save_grant(&grant);
596 } else {
597 outbox::drop_grant(&grant);
598 }
599 log::info!(
600 "share: {owner}'s {device} {} {grantee}",
601 if allow {
602 "shared with"
603 } else {
604 "no longer shared with"
605 }
606 );
607 self.send_shares(owner, device, None);
608 self.announce_to(grantee);
609 Ok(())
610 }
611
612 pub fn send_shares(&self, owner: &str, device: &str, error: Option<String>) {
615 let mut grantees: Vec<String> = self
616 .grants
617 .lock()
618 .iter()
619 .filter(|g| g.owner == owner && g.device == device)
620 .map(|g| g.grantee.clone())
621 .collect();
622 grantees.sort();
623 if let Some(e) = self
624 .entries
625 .lock()
626 .iter()
627 .find(|e| e.info.username == owner && e.device == device && e.wants_devices)
628 {
629 let accounts = outbox::accounts()
630 .into_iter()
631 .filter(|a| a != owner)
632 .collect();
633 let _ = e.tx.send(Envelope::from(LinkCommand::Shares {
634 grantees,
635 error,
636 accounts,
637 }));
638 }
639 }
640
641 pub fn seen_at(&self, device: &str, username: &str, addr: std::net::IpAddr) {
643 self.addresses
644 .lock()
645 .insert(device.to_string(), (username.to_string(), addr));
646 outbox::save_address(device, username, addr);
647 }
648
649 fn wake_owner(&self, username: &str, from: &str, to: &str) -> String {
656 if let Some(owner) = self.shared_owner(username, to) {
657 return owner;
658 }
659 let addresses = self.addresses.lock();
660 let here = addresses
661 .get(from)
662 .filter(|(user, _)| user == username)
663 .map(|(_, addr)| *addr);
664 match (here, addresses.get(to)) {
665 (Some(here), Some((owner, there))) if *there == here => owner.clone(),
666 _ => username.to_string(),
667 }
668 }
669
670 fn grantee_of(&self, owner: &str, from: &str, to: &str) -> Option<String> {
673 let to_user = self
674 .entries
675 .lock()
676 .iter()
677 .find(|e| e.device == to && e.info.username != owner)
678 .map(|e| e.info.username.clone())?;
679 self.grants
680 .lock()
681 .iter()
682 .any(|g| g.owner == owner && g.device == from && g.grantee == to_user)
683 .then_some(to_user)
684 }
685
686 fn shared_owner(&self, username: &str, to: &str) -> Option<String> {
689 self.grants
690 .lock()
691 .iter()
692 .find(|g| g.grantee == username && g.device == to)
693 .map(|g| g.owner.clone())
694 }
695
696 pub fn wake(&self, username: &str, from: &str, to: &str, notify: bool) {
702 let from_name = self
703 .list(Some(username))
704 .into_iter()
705 .find(|c| c.device == from)
706 .map_or_else(|| "Another device".to_string(), |c| c.name);
707 let owner = self.wake_owner(username, from, to);
708 if owner != username {
709 log::info!("wake: {to} is {owner}'s, woken for {username}");
710 }
711 let username = &owner;
712 if self.list(Some(username)).iter().any(|c| c.device == to) {
713 log::info!("wake: {to} is already linked");
714 return;
715 }
716 let Some(pusher) = crate::push::pusher() else {
717 log::info!("wake: no push key, so {to} cannot be woken");
718 return;
719 };
720 let Some(target) = outbox::push_targets(Some(username))
721 .into_iter()
722 .find(|t| t.device == to)
723 else {
724 log::info!("wake: {to} has given no push token");
725 return;
726 };
727 let (push, how) = if notify {
728 (
729 crate::push::Push::Summon {
730 title: format!("{from_name} wants to play here"),
731 body: "Tap to open kōan and let it.".into(),
732 },
733 "notification",
734 )
735 } else {
736 (crate::push::Push::WakeNow, "wake push")
737 };
738 WOKEN.lock().insert(
739 (target.device.clone(), target.username.clone()),
740 (std::time::Instant::now(), how),
741 );
742 std::thread::spawn(move || {
743 let started = std::time::Instant::now();
744 deliver_push(pusher, &target, &push);
745 log::info!(
746 "wake: {how} for {} answered by APNs in {}ms",
747 target.name,
748 started.elapsed().as_millis()
749 );
750 });
751 }
752
753 pub fn forget(&self, username: &str, device: &str) -> Result<(), String> {
760 if let Some(owner) = self.shared_owner(username, device) {
763 return self.share(&owner, device, username, false);
764 }
765 if self.list(Some(username)).iter().any(|c| c.device == device) {
766 return Err(format!("{device} is linked; it would be back at once"));
767 }
768 outbox::forget_device(device, username);
769 log::info!("devices: {username} forgot {device}");
770 self.broadcast(
771 Some(username),
772 LinkCommand::Forgotten {
773 device: device.to_string(),
774 },
775 );
776 self.announce(username);
777 Ok(())
778 }
779
780 pub fn relay(
782 &self,
783 username: &str,
784 to: &str,
785 command: LinkCommand,
786 ) -> Result<ClientInfo, String> {
787 self.relay_from(username, None, to, command)
788 }
789
790 pub fn relay_from(
794 &self,
795 username: &str,
796 from: Option<&str>,
797 to: &str,
798 command: LinkCommand,
799 ) -> Result<ClientInfo, String> {
800 self.relay_acked(username, from, to, command, None)
801 }
802
803 pub fn relay_acked(
805 &self,
806 username: &str,
807 from: Option<&str>,
808 to: &str,
809 command: LinkCommand,
810 acking: Option<Acking>,
811 ) -> Result<ClientInfo, String> {
812 if !command.relayable() {
816 return Err("not a command".into());
817 }
818 if let Some(owner) = self.shared_owner(username, to) {
822 if !command.allowed_playback() {
823 return Err(format!("{to} is shared for playback only"));
824 }
825 let command = LinkCommand::Shared {
826 command: Box::new(command),
827 };
828 return self.send_with(Some(&owner), Some(to), command, acking);
829 }
830 if let Some(from) = from
834 && let Some(grantee) = self.grantee_of(username, from, to)
835 {
836 if !matches!(command, LinkCommand::Play { handoff: true, .. }) {
837 return Err(format!("{to} is another account's"));
838 }
839 let command = LinkCommand::Shared {
840 command: Box::new(command),
841 };
842 return self.send_with(Some(&grantee), Some(to), command, acking);
843 }
844 self.send_with(Some(username), Some(to), command, acking)
845 }
846
847 pub fn set_activity(
850 &self,
851 username: &str,
852 watcher: &str,
853 activity: Option<(String, String, bool)>,
854 ) {
855 let mut activities = self.activities.lock();
856 activities.retain(|a| !(a.username == username && a.watcher == watcher));
857 let Some((token, target, sandbox)) = activity else {
858 return;
859 };
860 activities.push(Activity {
861 username: username.to_string(),
862 watcher: watcher.to_string(),
863 target: target.clone(),
864 token,
865 sandbox,
866 sent: None,
867 });
868 drop(activities);
869 self.update_activities(username, &target);
870 }
871
872 fn update_activities(&self, username: &str, target: &str) {
875 let Some(pusher) = crate::push::pusher() else {
876 return;
877 };
878 let owner = self
881 .shared_owner(username, target)
882 .unwrap_or_else(|| username.to_string());
883 let Some(info) = self
884 .list(Some(&owner))
885 .into_iter()
886 .find(|c| c.device == target)
887 else {
888 return;
889 };
890 let watchers: Vec<String> = self
891 .grants
892 .lock()
893 .iter()
894 .filter(|g| g.owner == owner && g.device == target)
895 .map(|g| g.grantee.clone())
896 .chain(std::iter::once(owner.clone()))
897 .collect();
898 let state = crate::push::ActivityState::of(&info);
899 let mut due = Vec::new();
900 for a in self.activities.lock().iter_mut() {
901 if watchers.contains(&a.username)
902 && a.target == target
903 && a.sent.as_ref().is_none_or(|s| s.differs(&state))
904 {
905 a.sent = Some(state.clone());
906 due.push((
907 a.token.clone(),
908 a.sandbox,
909 a.watcher.clone(),
910 a.username.clone(),
911 ));
912 }
913 }
914 if due.is_empty() {
915 return;
916 }
917 std::thread::spawn(move || {
918 for (token, sandbox, watcher, username) in due {
919 let push = crate::push::Push::Activity(state.clone());
920 match pusher.send(&token, sandbox, &push) {
921 crate::push::Outcome::Sent => {}
922 crate::push::Outcome::Gone => {
923 log::info!("push: a Live Activity on {watcher} has ended");
924 registry().set_activity(&username, &watcher, None);
925 }
926 crate::push::Outcome::Failed(e) => {
927 log::warn!("push: Live Activity on {watcher}: {e}");
928 }
929 }
930 }
931 });
932 }
933
934 pub fn list(&self, username: Option<&str>) -> Vec<ClientInfo> {
936 let mut out: Vec<ClientInfo> = self
937 .entries
938 .lock()
939 .iter()
940 .filter(|e| username.is_none_or(|u| e.info.username == u))
941 .map(|e| e.info.clone())
942 .collect();
943 out.sort_by_key(|c| std::cmp::Reverse(c.connected_at));
944 out
945 }
946
947 pub fn send(
952 &self,
953 username: Option<&str>,
954 id: Option<&str>,
955 cmd: LinkCommand,
956 ) -> Result<ClientInfo, String> {
957 self.send_with(username, id, cmd, None)
958 }
959
960 pub fn send_with(
966 &self,
967 username: Option<&str>,
968 id: Option<&str>,
969 cmd: LinkCommand,
970 acking: Option<Acking>,
971 ) -> Result<ClientInfo, String> {
972 let clients = self.list(username);
973 let target = match id {
974 Some(id) => clients
975 .iter()
976 .find(|c| c.id == id || c.device == id || c.name.eq_ignore_ascii_case(id)),
977 None if clients.is_empty() => None,
978 None => Some(pick(&clients, chrono::Utc::now().timestamp())?),
979 };
980 let Some(target) = target else {
982 let reached = reach_absent(username, id, &cmd, acking.as_ref().map(|a| a.id))
983 .unwrap_or_else(|| {
984 Err(match id {
985 Some(id) => format!("no linked client {id}; see `clients`"),
986 None => {
987 "no koan app is linked to this server; open koan on the device".into()
988 }
989 })
990 });
991 if let (Ok(_), Some(acking)) = (&reached, acking) {
992 (acking.reply)(Some(AckOutcome::Queued));
993 }
994 return reached;
995 };
996 let target = target.clone();
997 let entries = self.entries.lock();
998 let entry = entries
999 .iter()
1000 .find(|e| e.info.id == target.id)
1001 .ok_or("that client has just gone")?;
1002 let device = entry.device.clone();
1003 let (ack, unanswered) = match acking {
1007 Some(acking) if entry.info.acks => {
1008 self.answers.lock().insert(
1009 (device.clone(), acking.id),
1010 Waiting {
1011 reply: acking.reply,
1012 received: false,
1013 timer: None,
1014 },
1015 );
1016 (Some(acking.id), None)
1017 }
1018 other => (None, other),
1019 };
1020 let sent = entry.tx.send(Envelope {
1021 command: cmd.clone(),
1022 ack,
1023 });
1024 drop(entries);
1025 if sent.is_err() {
1026 if let Some(w) = ack.and_then(|ack| self.answers.lock().remove(&(device, ack))) {
1027 w.answer(Some(AckOutcome::Failed {
1028 error: "its link has just gone".into(),
1029 }));
1030 }
1031 return Err("that client has just gone".to_string());
1032 }
1033 if let Some(acking) = unanswered {
1034 (acking.reply)(None);
1035 }
1036 if let Some(ack) = ack {
1037 self.watch(target.username.clone(), device, ack, cmd);
1038 }
1039 Ok(target)
1040 }
1041
1042 fn watch(&self, username: String, device: String, ack: u64, cmd: LinkCommand) {
1047 let Ok(runtime) = tokio::runtime::Handle::try_current() else {
1048 return;
1049 };
1050 let key = (device.clone(), ack);
1051 let timer = runtime.spawn(async move {
1052 tokio::time::sleep(FIRST_LOOK).await;
1053 let looking = device.clone();
1054 let taken = tokio::task::spawn_blocking(move || {
1055 registry().look(&username, &looking, ack, &cmd)
1056 })
1057 .await
1058 .unwrap_or(false);
1059 if taken {
1060 tokio::time::sleep(LONGEST_ANSWER.saturating_sub(FIRST_LOOK)).await;
1061 registry().give_up(&device, ack);
1062 }
1063 });
1064 match self.answers.lock().get_mut(&key) {
1065 Some(w) => w.timer = Some(timer.abort_handle()),
1066 None => timer.abort(),
1068 }
1069 }
1070
1071 pub fn answered(&self, device: &str, ack: u64, outcome: AckOutcome) {
1073 let waiting = self.answers.lock().remove(&(device.to_string(), ack));
1074 if let Some(w) = waiting {
1075 w.answer(Some(outcome));
1076 }
1077 }
1078
1079 pub fn received(&self, device: &str, ack: u64) {
1081 if let Some(w) = self.answers.lock().get_mut(&(device.to_string(), ack)) {
1082 w.received = true;
1083 }
1084 }
1085
1086 fn look(&self, username: &str, device: &str, ack: u64, cmd: &LinkCommand) -> bool {
1094 let key = (device.to_string(), ack);
1095 let mut answers = self.answers.lock();
1096 match answers.get(&key) {
1097 None => return false,
1098 Some(w) if w.received => return true,
1099 Some(_) => {}
1100 }
1101 let w = answers.remove(&key).expect("just seen");
1102 drop(answers);
1103 log::info!("link: {device} has not taken a command; reaching it as away");
1104 let outcome = match reach_absent(Some(username), Some(device), cmd, Some(ack)) {
1105 Some(Ok(_)) => AckOutcome::Queued,
1106 _ => AckOutcome::Failed {
1107 error: format!("{device} did not take it"),
1108 },
1109 };
1110 w.answer(Some(outcome));
1111 false
1112 }
1113
1114 fn give_up(&self, device: &str, ack: u64) {
1117 let waiting = self.answers.lock().remove(&(device.to_string(), ack));
1118 if let Some(w) = waiting {
1119 w.answer(None);
1120 }
1121 }
1122}
1123
1124impl Registry {
1125 pub fn add_order(&self, order: Order) {
1126 outbox::save_order(&order);
1127 self.orders.lock().push(order);
1128 }
1129
1130 pub fn orders(&self, username: Option<&str>) -> Vec<Order> {
1131 self.orders
1132 .lock()
1133 .iter()
1134 .filter(|o| username.is_none() || o.username.as_deref() == username)
1135 .cloned()
1136 .collect()
1137 }
1138
1139 pub fn cancel_order(&self, username: Option<&str>, id: &str) -> bool {
1140 let mut orders = self.orders.lock();
1141 let before = orders.len();
1142 orders
1143 .retain(|o| !(o.id == id && (username.is_none() || o.username.as_deref() == username)));
1144 let gone = orders.len() != before;
1145 if gone {
1146 outbox::drop_order(id);
1147 }
1148 gone
1149 }
1150
1151 fn done(&self, id: &str) {
1152 self.orders.lock().retain(|o| o.id != id);
1153 outbox::drop_order(id);
1154 }
1155
1156 pub fn fulfil_orders(&self, find: impl Fn(&Order) -> Option<Vec<String>>) {
1159 let now = chrono::Utc::now().timestamp();
1160 let pending: Vec<Order> = {
1161 let mut orders = self.orders.lock();
1162 for o in orders.iter().filter(|o| now - o.created_at >= ORDER_TTL) {
1163 outbox::drop_order(&o.id);
1164 }
1165 orders.retain(|o| now - o.created_at < ORDER_TTL);
1166 orders.clone()
1167 };
1168 for order in pending.into_iter().filter(|o| o.playlist.is_none()) {
1169 let Some(ids) = find(&order).filter(|ids| !ids.is_empty()) else {
1170 continue;
1171 };
1172 let track_ids = ids;
1173 let cmd = if order.play_next {
1174 LinkCommand::PlayNext { track_ids }
1175 } else {
1176 LinkCommand::Enqueue { track_ids }
1177 };
1178 match self.send(order.username.as_deref(), order.client.as_deref(), cmd) {
1179 Ok(c) => {
1180 log::info!(
1181 "link: {} — {} arrived; queued on {}",
1182 order.artist,
1183 order.album,
1184 c.name
1185 );
1186 self.done(&order.id);
1187 }
1188 Err(e) => log::info!("link: {} — {} arrived but {e}", order.artist, order.album),
1190 }
1191 }
1192 }
1193}
1194
1195pub fn fulfil_from(db_path: &std::path::Path) {
1197 let registry = registry();
1198 if registry.orders.lock().is_empty() {
1199 return;
1200 }
1201 let Ok(db) = koan_core::db::connection::Database::open_existing(db_path) else {
1202 return;
1203 };
1204 let for_playlists: Vec<Order> = registry
1207 .orders
1208 .lock()
1209 .iter()
1210 .filter(|o| o.playlist.is_some())
1211 .cloned()
1212 .collect();
1213 let mut edited = false;
1214 for order in for_playlists {
1215 let Some(playlist) = order.playlist else {
1216 continue;
1217 };
1218 let Some(ids) = order_tracks(&db.conn, &order) else {
1219 continue;
1220 };
1221 match koan_core::db::queries::add_tracks(&db.conn, playlist, &ids) {
1222 Ok(_) => {
1223 if !cfg!(test) {
1225 koan_core::playlists::push_to_remote(playlist);
1226 }
1227 log::info!(
1228 "link: {} — {} arrived; added {} tracks to playlist {playlist}",
1229 order.artist,
1230 order.album,
1231 ids.len()
1232 );
1233 registry.done(&order.id);
1234 edited = true;
1235 }
1236 Err(e) => log::warn!("link: could not add to playlist {playlist}: {e}"),
1237 }
1238 }
1239 if edited {
1240 changed();
1241 }
1242 registry.fulfil_orders(|order| {
1243 let rows = order_tracks(&db.conn, order)?;
1244 queries::uids_in_order(&db.conn, UidKind::Track, &rows).ok()
1245 });
1246}
1247
1248fn order_tracks(conn: &rusqlite::Connection, order: &Order) -> Option<Vec<i64>> {
1251 let tracks = album_tracks(conn, &order.artist, &order.album)?;
1252 if order.titles.is_empty() {
1253 return Some(tracks.into_iter().map(|(id, _)| id).collect());
1254 }
1255 let picked: Vec<i64> = order
1256 .titles
1257 .iter()
1258 .filter_map(|want| {
1259 let want = want.to_lowercase();
1260 tracks
1261 .iter()
1262 .find(|(_, t)| t.to_lowercase().contains(&want))
1263 .map(|(id, _)| *id)
1264 })
1265 .collect();
1266 (!picked.is_empty()).then_some(picked)
1267}
1268
1269pub fn album_tracks(
1272 conn: &rusqlite::Connection,
1273 artist: &str,
1274 album: &str,
1275) -> Option<Vec<(i64, String)>> {
1276 let like = |s: &str| format!("%{}%", s.replace(['%', '_'], ""));
1277 let album_id: i64 = conn
1278 .query_row(
1279 "SELECT al.id FROM albums al JOIN artists a ON a.id = al.artist_id
1280 WHERE a.name LIKE ?1 COLLATE NOCASE AND al.title LIKE ?2 COLLATE NOCASE
1281 ORDER BY al.id DESC LIMIT 1",
1282 [like(artist), like(album)],
1283 |r| r.get(0),
1284 )
1285 .ok()?;
1286 let mut stmt = conn
1287 .prepare("SELECT id, title FROM tracks WHERE album_id = ?1 ORDER BY disc, track_number, id")
1288 .ok()?;
1289 let tracks = stmt
1290 .query_map([album_id], |r| Ok((r.get(0)?, r.get(1)?)))
1291 .ok()?
1292 .filter_map(Result::ok)
1293 .collect();
1294 Some(tracks)
1295}
1296
1297impl Registry {
1298 pub fn broadcast(&self, username: Option<&str>, cmd: LinkCommand) -> Vec<String> {
1301 let ids: Vec<String> = self.list(username).into_iter().map(|c| c.id).collect();
1302 let entries = self.entries.lock();
1303 entries
1304 .iter()
1305 .filter(|e| ids.contains(&e.info.id) && e.tx.send(cmd.clone().into()).is_ok())
1306 .map(|e| e.info.name.clone())
1307 .collect()
1308 }
1309}
1310
1311impl Registry {
1312 pub fn deliver(&self, username: Option<&str>, cmd: LinkCommand) -> (Vec<String>, Vec<String>) {
1317 let (sent, queued) = self.link_or_queue(username, &cmd);
1318 wake(&queued.iter().map(Absent::key).collect::<Vec<_>>());
1319 (sent, queued.into_iter().map(|q| q.name).collect())
1320 }
1321
1322 fn link_or_queue(
1323 &self,
1324 username: Option<&str>,
1325 cmd: &LinkCommand,
1326 ) -> (Vec<String>, Vec<Absent>) {
1327 let sent = self.broadcast(username, cmd.clone());
1328 let queued = outbox::queue_for_absent(username, &self.live(), cmd);
1329 (sent, queued)
1330 }
1331
1332 fn live(&self) -> Vec<(String, String)> {
1334 self.entries
1335 .lock()
1336 .iter()
1337 .map(|e| (e.device.clone(), e.info.username.clone()))
1338 .collect()
1339 }
1340}
1341
1342impl Absent {
1343 fn key(&self) -> (String, String) {
1344 (self.device.clone(), self.username.clone())
1345 }
1346}
1347
1348fn wake(devices: &[(String, String)]) {
1351 let Some(pusher) = crate::push::pusher() else {
1352 return;
1353 };
1354 let targets: Vec<outbox::PushTarget> = outbox::push_targets(None)
1355 .into_iter()
1356 .filter(|t| {
1357 devices
1358 .iter()
1359 .any(|(device, username)| *device == t.device && *username == t.username)
1360 })
1361 .collect();
1362 if targets.is_empty() {
1363 return;
1364 }
1365 std::thread::spawn(move || {
1366 for t in targets {
1367 deliver_push(pusher, &t, &crate::push::Push::Wake);
1368 }
1369 });
1370}
1371
1372fn reach_absent(
1381 username: Option<&str>,
1382 id: Option<&str>,
1383 cmd: &LinkCommand,
1384 ack: Option<u64>,
1385) -> Option<Result<ClientInfo, String>> {
1386 let envelope = Envelope {
1387 command: cmd.clone(),
1388 ack,
1389 };
1390 let pusher = crate::push::pusher()?;
1391 let target = outbox::push_targets(username)
1394 .into_iter()
1395 .find(|t| id.is_none_or(|id| t.device == id || t.name.eq_ignore_ascii_case(id)))?;
1396 let info = ClientInfo {
1397 id: target.device.clone(),
1398 device: target.device.clone(),
1399 name: target.name.clone(),
1400 platform: target.platform.clone(),
1401 username: target.username.clone(),
1402 connected_at: 0,
1403 state: LinkState::default(),
1404 last_played_at: None,
1405 state_at: 0,
1406 reports: false,
1407 notified: true,
1408 acks: false,
1409 };
1410 let asked = match cmd {
1412 LinkCommand::Shared { command } => command.as_ref(),
1413 cmd => cmd,
1414 };
1415 let verb = match asked {
1416 LinkCommand::Play { .. } | LinkCommand::JumpTo { .. } | LinkCommand::Resume => Some("Play"),
1417 _ => None,
1418 };
1419 let push = match verb {
1420 Some(verb) => crate::push::Push::Notify {
1421 title: format!("{verb} on {}", target.name),
1422 body: outbox::describe(asked).unwrap_or_else(|| "From your koan server".into()),
1423 command: serde_json::to_value(&envelope).ok()?,
1424 image: cover_track(asked)
1425 .and_then(outbox::track_row)
1426 .and_then(|t| pusher.cover_link(t)),
1427 },
1428 None => {
1429 outbox::queue_for(&target.device, &target.username, &envelope);
1430 crate::push::Push::Wake
1431 }
1432 };
1433 std::thread::spawn(move || deliver_push(pusher, &target, &push));
1434 Some(Ok(info))
1435}
1436
1437fn cover_track(cmd: &LinkCommand) -> Option<&str> {
1440 match cmd {
1441 LinkCommand::Play {
1442 track_ids,
1443 start_at,
1444 ..
1445 } => track_ids
1446 .get(*start_at as usize)
1447 .or(track_ids.first())
1448 .map(String::as_str),
1449 LinkCommand::JumpTo { track_id } => Some(track_id),
1450 _ => None,
1451 }
1452}
1453
1454fn deliver_push(
1456 pusher: &crate::push::Pusher,
1457 target: &outbox::PushTarget,
1458 push: &crate::push::Push,
1459) {
1460 use crate::push::Outcome;
1461 match pusher.send(&target.token, target.sandbox, push) {
1462 Outcome::Sent => log::info!("push: sent to {}", target.name),
1463 Outcome::Gone => {
1464 log::info!(
1465 "push: {}'s token is no longer valid; forgotten",
1466 target.name
1467 );
1468 outbox::forget_push(&target.username, &target.device);
1469 }
1470 Outcome::Failed(e) => log::warn!("push: to {} failed: {e}", target.name),
1471 }
1472}
1473
1474type Woken = std::collections::HashMap<(String, String), (std::time::Instant, &'static str)>;
1478
1479static WOKEN: LazyLock<Mutex<Woken>> = LazyLock::new(Default::default);
1480
1481pub fn smart_activity(
1485 db: &koan_core::db::connection::Database,
1486 user: i64,
1487 fields: &[koan_core::smart::Field],
1488) {
1489 match koan_core::db::queries::smart::refresh_after_activity(&db.conn, user, fields) {
1490 Ok(moved) if !moved.is_empty() => changed(),
1491 Ok(_) => {}
1492 Err(e) => log::warn!("smart playlists not refreshed after activity: {e}"),
1493 }
1494}
1495
1496pub fn changed() {
1501 let (_, queued) = registry().link_or_queue(None, &LinkCommand::Sync { full: false });
1502 if queued.is_empty() || crate::push::pusher().is_none() {
1503 return;
1504 }
1505 let mut wakes = WAKES.lock();
1506 wakes.add(queued.iter().map(Absent::key), std::time::Instant::now());
1507 if !wakes.timer {
1508 wakes.timer = true;
1509 std::thread::spawn(send_wakes);
1510 }
1511}
1512
1513const QUIET: std::time::Duration = std::time::Duration::from_secs(30);
1517
1518const LONGEST_WAIT: std::time::Duration = std::time::Duration::from_secs(5 * 60);
1520
1521static WAKES: Mutex<Wakes> = Mutex::new(Wakes {
1522 pending: Vec::new(),
1523 timer: false,
1524});
1525
1526struct Wakes {
1529 pending: Vec<((String, String), std::time::Instant, std::time::Instant)>,
1531 timer: bool,
1533}
1534
1535impl Wakes {
1536 fn add(
1537 &mut self,
1538 devices: impl IntoIterator<Item = (String, String)>,
1539 now: std::time::Instant,
1540 ) {
1541 for key in devices {
1542 match self.pending.iter_mut().find(|(k, _, _)| *k == key) {
1543 Some((_, _, last)) => *last = now,
1544 None => self.pending.push((key, now, now)),
1545 }
1546 }
1547 }
1548
1549 fn due_at(first: std::time::Instant, last: std::time::Instant) -> std::time::Instant {
1550 (last + QUIET).min(first + LONGEST_WAIT)
1551 }
1552
1553 fn next(&self) -> Option<std::time::Instant> {
1555 self.pending
1556 .iter()
1557 .map(|(_, first, last)| Self::due_at(*first, *last))
1558 .min()
1559 }
1560
1561 fn take_due(&mut self, now: std::time::Instant) -> Vec<(String, String)> {
1563 let (due, waiting) = std::mem::take(&mut self.pending)
1564 .into_iter()
1565 .partition(|(_, first, last)| Self::due_at(*first, *last) <= now);
1566 self.pending = waiting;
1567 due.into_iter().map(|(key, _, _)| key).collect()
1568 }
1569}
1570
1571fn send_wakes() {
1574 loop {
1575 let (due, next) = {
1576 let mut wakes = WAKES.lock();
1577 let due = wakes.take_due(std::time::Instant::now());
1578 let next = wakes.next();
1579 if due.is_empty() && next.is_none() {
1580 wakes.timer = false;
1581 return;
1582 }
1583 (due, next)
1584 };
1585 let live = registry().live();
1586 let absent: Vec<(String, String)> = due.into_iter().filter(|d| !live.contains(d)).collect();
1587 if !absent.is_empty() {
1588 wake(&absent);
1589 }
1590 if let Some(next) = next {
1591 std::thread::sleep(next.saturating_duration_since(std::time::Instant::now()));
1592 }
1593 }
1594}
1595
1596pub fn changed_if_library_moved(conn: &rusqlite::Connection) {
1600 static LAST: parking_lot::Mutex<Option<(i64, i64, i64)>> = parking_lot::Mutex::new(None);
1601 let Ok(now) = conn.query_row(
1602 "SELECT (SELECT COUNT(*) FROM tracks), (SELECT COALESCE(MAX(id), 0) FROM tracks),
1603 (SELECT COUNT(*) FROM albums)",
1604 [],
1605 |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
1606 ) else {
1607 return;
1608 };
1609 let before = LAST.lock().replace(now);
1610 if before.is_some_and(|b| b != now) {
1611 changed();
1612 }
1613}
1614
1615mod outbox {
1618 use koan_core::db::queries;
1619 use koan_core::remote::link::LinkCommand;
1620
1621 const KEEP_SECS: i64 = 30 * 24 * 60 * 60;
1626
1627 fn db() -> Option<koan_core::db::pool::Handle<'static>> {
1632 if cfg!(test) {
1633 return None;
1634 }
1635 koan_core::db::pool::shared().get().ok()
1636 }
1637
1638 pub fn load_orders() -> Vec<super::Order> {
1639 let Some(db) = db() else { return Vec::new() };
1640 db.conn
1641 .prepare("SELECT body FROM link_orders ORDER BY created_at")
1642 .and_then(|mut s| {
1643 s.query_map([], |r| r.get::<_, String>(0))?
1644 .collect::<Result<Vec<_>, _>>()
1645 })
1646 .unwrap_or_default()
1647 .into_iter()
1648 .filter_map(|b| serde_json::from_str(&b).ok())
1649 .collect()
1650 }
1651
1652 pub fn save_order(order: &super::Order) {
1653 let (Some(db), Ok(body)) = (db(), serde_json::to_string(order)) else {
1654 return;
1655 };
1656 let _ = db.conn.execute(
1657 "INSERT OR REPLACE INTO link_orders (id, body, created_at) VALUES (?1, ?2, ?3)",
1658 rusqlite::params![order.id, body, order.created_at],
1659 );
1660 }
1661
1662 pub fn drop_order(id: &str) {
1663 if let Some(db) = db() {
1664 let _ = db
1665 .conn
1666 .execute("DELETE FROM link_orders WHERE id = ?1", [id]);
1667 }
1668 }
1669
1670 pub fn take_and_remember(
1671 username: &str,
1672 device: &str,
1673 name: &str,
1674 platform: &str,
1675 live: &[(String, String)],
1676 ) -> Vec<koan_core::remote::acks::Envelope> {
1677 let Some(db) = db() else { return Vec::new() };
1678 let now = chrono::Utc::now().timestamp();
1679 let _ = db.conn.execute(
1680 "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES (?1, ?2, ?3, ?4, ?5)
1681 ON CONFLICT (device, username) DO UPDATE SET name = ?3, platform = ?4, last_seen = ?5",
1682 rusqlite::params![device, username, name, platform, now],
1683 );
1684 forget_stale(&db.conn, live, now);
1685 let waiting: Vec<(i64, String)> = db
1686 .conn
1687 .prepare("SELECT id, command FROM link_outbox WHERE device = ?1 AND username = ?2 ORDER BY id")
1688 .and_then(|mut s| {
1689 s.query_map([device, username], |r| Ok((r.get(0)?, r.get(1)?)))?
1690 .collect()
1691 })
1692 .unwrap_or_default();
1693 let _ = db.conn.execute(
1694 "DELETE FROM link_outbox WHERE device = ?1 AND username = ?2",
1695 [device, username],
1696 );
1697 if !waiting.is_empty() {
1698 log::info!("link: {} waiting commands for {name}", waiting.len());
1699 }
1700 waiting
1701 .into_iter()
1702 .filter_map(|(_, c)| serde_json::from_str(&c).ok())
1703 .collect()
1704 }
1705
1706 pub(super) fn forget_stale(conn: &rusqlite::Connection, live: &[(String, String)], now: i64) {
1709 touch_with(conn, live, now);
1710 let _ = conn.execute(
1711 "DELETE FROM link_outbox WHERE created_at < ?1",
1712 [now - KEEP_SECS],
1713 );
1714 let _ = conn.execute(
1715 "DELETE FROM link_push WHERE (device, username) IN
1716 (SELECT device, username FROM link_devices WHERE last_seen < ?1)",
1717 [now - KEEP_SECS],
1718 );
1719 let _ = conn.execute(
1720 "DELETE FROM link_devices WHERE last_seen < ?1",
1721 [now - KEEP_SECS],
1722 );
1723 }
1724
1725 pub fn touch(devices: &[(String, String)]) {
1727 if let Some(db) = db() {
1728 touch_with(&db.conn, devices, chrono::Utc::now().timestamp());
1729 }
1730 }
1731
1732 fn touch_with(conn: &rusqlite::Connection, devices: &[(String, String)], now: i64) {
1733 for (device, username) in devices {
1734 let _ = conn.execute(
1735 "UPDATE link_devices SET last_seen = ?1 WHERE device = ?2 AND username = ?3",
1736 rusqlite::params![now, device, username],
1737 );
1738 }
1739 }
1740
1741 pub struct Absent {
1743 pub device: String,
1744 pub username: String,
1745 pub name: String,
1746 }
1747
1748 pub fn queue_for_absent(
1750 username: Option<&str>,
1751 live: &[(String, String)],
1752 cmd: &LinkCommand,
1753 ) -> Vec<Absent> {
1754 let Some(db) = db() else { return Vec::new() };
1755 let known: Vec<(String, String, String)> = db
1756 .conn
1757 .prepare("SELECT device, username, name FROM link_devices")
1758 .and_then(|mut s| {
1759 s.query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))?
1760 .collect()
1761 })
1762 .unwrap_or_default();
1763 let Ok(text) = serde_json::to_string(cmd) else {
1764 return Vec::new();
1765 };
1766 let absent: Vec<_> = known
1767 .into_iter()
1768 .filter(|(device, user, _)| {
1769 !username.is_some_and(|u| u != user)
1770 && !live.iter().any(|(d, u)| d == device && u == user)
1771 })
1772 .collect();
1773 if absent.is_empty() {
1774 return Vec::new();
1775 }
1776 let is_sync = matches!(cmd, LinkCommand::Sync { .. });
1777 let now = chrono::Utc::now().timestamp();
1778 koan_core::db::queries::atomically(&db.conn, || {
1781 let mut queued = Vec::new();
1782 for (device, user, name) in absent {
1783 if is_sync {
1784 let _ = db.conn.execute(
1786 "DELETE FROM link_outbox WHERE device = ?1 AND username = ?2 AND command LIKE '{\"type\":\"sync\"%'",
1787 [&device, &user],
1788 );
1789 }
1790 if db
1791 .conn
1792 .execute(
1793 "INSERT INTO link_outbox (device, username, command, created_at) VALUES (?1, ?2, ?3, ?4)",
1794 rusqlite::params![device, user, text, now],
1795 )
1796 .is_ok()
1797 {
1798 queued.push(Absent {
1799 device,
1800 username: user,
1801 name,
1802 });
1803 }
1804 }
1805 Ok::<_, rusqlite::Error>(queued)
1806 })
1807 .unwrap_or_default()
1808 }
1809
1810 #[derive(Clone)]
1812 pub struct PushTarget {
1813 pub device: String,
1814 pub username: String,
1815 pub name: String,
1816 pub platform: String,
1817 pub token: String,
1818 pub sandbox: bool,
1819 pub last_seen: i64,
1821 }
1822
1823 pub fn queue_for(device: &str, username: &str, envelope: &koan_core::remote::acks::Envelope) {
1827 let (Some(db), Ok(text)) = (db(), serde_json::to_string(envelope)) else {
1828 return;
1829 };
1830 let _ = db.conn.execute(
1831 "INSERT INTO link_outbox (device, username, command, created_at) VALUES (?1, ?2, ?3, ?4)",
1832 rusqlite::params![device, username, text, chrono::Utc::now().timestamp()],
1833 );
1834 }
1835
1836 pub fn accounts() -> Vec<String> {
1840 let Some(db) = db() else { return Vec::new() };
1841 koan_core::db::queries::auth::list_users(&db.conn)
1842 .map(|users| users.into_iter().map(|u| u.username).collect())
1843 .unwrap_or_default()
1844 }
1845
1846 pub fn load_addresses() -> std::collections::HashMap<String, (String, std::net::IpAddr)> {
1849 let Some(db) = db() else {
1850 return Default::default();
1851 };
1852 load_addresses_in(&db.conn)
1853 }
1854
1855 pub(super) fn load_addresses_in(
1856 conn: &rusqlite::Connection,
1857 ) -> std::collections::HashMap<String, (String, std::net::IpAddr)> {
1858 conn.prepare(
1859 "SELECT device, username, addr FROM link_devices WHERE addr IS NOT NULL ORDER BY last_seen",
1860 )
1861 .and_then(|mut s| {
1862 s.query_map([], |r| {
1863 Ok((
1864 r.get::<_, String>(0)?,
1865 r.get::<_, String>(1)?,
1866 r.get::<_, String>(2)?,
1867 ))
1868 })?
1869 .collect::<Result<Vec<_>, _>>()
1870 })
1871 .unwrap_or_default()
1872 .into_iter()
1873 .filter_map(|(device, username, addr)| Some((device, (username, addr.parse().ok()?))))
1874 .collect()
1875 }
1876
1877 pub fn save_address(device: &str, username: &str, addr: std::net::IpAddr) {
1878 if let Some(db) = db() {
1879 save_address_in(&db.conn, device, username, addr);
1880 }
1881 }
1882
1883 pub(super) fn save_address_in(
1884 conn: &rusqlite::Connection,
1885 device: &str,
1886 username: &str,
1887 addr: std::net::IpAddr,
1888 ) {
1889 let _ = conn.execute(
1890 "UPDATE link_devices SET addr = ?1 WHERE device = ?2 AND username = ?3",
1891 rusqlite::params![addr.to_string(), device, username],
1892 );
1893 }
1894
1895 pub fn load_grants() -> Vec<super::Grant> {
1896 let Some(db) = db() else { return Vec::new() };
1897 db.conn
1898 .prepare("SELECT device, owner, grantee FROM link_grants ORDER BY created_at")
1899 .and_then(|mut s| {
1900 s.query_map([], |r| {
1901 Ok(super::Grant {
1902 device: r.get(0)?,
1903 owner: r.get(1)?,
1904 grantee: r.get(2)?,
1905 })
1906 })?
1907 .collect()
1908 })
1909 .unwrap_or_default()
1910 }
1911
1912 pub fn save_grant(g: &super::Grant) {
1913 if let Some(db) = db() {
1914 let _ = db.conn.execute(
1915 "INSERT OR IGNORE INTO link_grants (device, owner, grantee, created_at) VALUES (?1, ?2, ?3, ?4)",
1916 rusqlite::params![g.device, g.owner, g.grantee, chrono::Utc::now().timestamp()],
1917 );
1918 }
1919 }
1920
1921 pub fn drop_grant(g: &super::Grant) {
1922 if let Some(db) = db() {
1923 let _ = db.conn.execute(
1924 "DELETE FROM link_grants WHERE device = ?1 AND owner = ?2 AND grantee = ?3",
1925 [&g.device, &g.owner, &g.grantee],
1926 );
1927 }
1928 }
1929
1930 pub fn save_push(username: &str, device: &str, token: &str, sandbox: bool) {
1931 let Some(db) = db() else { return };
1932 let _ = db.conn.execute(
1933 "INSERT INTO link_push (device, username, token, sandbox, updated_at) VALUES (?1, ?2, ?3, ?4, ?5)
1934 ON CONFLICT (device, username) DO UPDATE SET token = ?3, sandbox = ?4, updated_at = ?5",
1935 rusqlite::params![device, username, token, sandbox, chrono::Utc::now().timestamp()],
1936 );
1937 }
1938
1939 pub fn forget_push(username: &str, device: &str) {
1940 if let Some(db) = db() {
1941 let _ = db.conn.execute(
1942 "DELETE FROM link_push WHERE device = ?1 AND username = ?2",
1943 [device, username],
1944 );
1945 }
1946 }
1947
1948 pub fn push_targets(username: Option<&str>) -> Vec<PushTarget> {
1950 let Some(db) = db() else { return Vec::new() };
1951 push_targets_in(&db.conn, username)
1952 }
1953
1954 pub(super) fn push_targets_in(
1955 conn: &rusqlite::Connection,
1956 username: Option<&str>,
1957 ) -> Vec<PushTarget> {
1958 conn
1959 .prepare(
1960 "SELECT p.device, p.username, d.name, d.platform, p.token, p.sandbox, d.last_seen
1961 FROM link_push p JOIN link_devices d ON d.device = p.device AND d.username = p.username
1962 WHERE ?1 IS NULL OR p.username = ?1
1963 ORDER BY d.last_seen DESC",
1964 )
1965 .and_then(|mut s| {
1966 s.query_map([username], |r| {
1967 Ok(PushTarget {
1968 device: r.get(0)?,
1969 username: r.get(1)?,
1970 name: r.get(2)?,
1971 platform: r.get(3)?,
1972 token: r.get(4)?,
1973 sandbox: r.get(5)?,
1974 last_seen: r.get(6)?,
1975 })
1976 })?
1977 .collect()
1978 })
1979 .unwrap_or_default()
1980 }
1981
1982 pub fn forget_device(device: &str, username: &str) {
1985 if let Some(db) = db() {
1986 forget_device_in(&db.conn, device, username);
1987 }
1988 }
1989
1990 pub(super) fn forget_device_in(conn: &rusqlite::Connection, device: &str, username: &str) {
1991 for table in ["link_push", "link_outbox", "link_devices"] {
1992 let _ = conn.execute(
1993 &format!("DELETE FROM {table} WHERE device = ?1 AND username = ?2"),
1994 [device, username],
1995 );
1996 }
1997 }
1998
1999 pub fn track_row(id: &str) -> Option<i64> {
2001 let db = db()?;
2002 queries::resolve_id(&db.conn, queries::UidKind::Track, id)
2003 .ok()
2004 .flatten()
2005 }
2006
2007 pub fn describe(cmd: &LinkCommand) -> Option<String> {
2010 let ids: Vec<&String> = match cmd {
2011 LinkCommand::Play { track_ids, .. }
2012 | LinkCommand::Enqueue { track_ids }
2013 | LinkCommand::PlayNext { track_ids } => track_ids.iter().collect(),
2014 LinkCommand::JumpTo { track_id } => vec![track_id],
2015 _ => return None,
2016 };
2017 let db = db()?;
2018 let ids: Vec<i64> = ids
2019 .into_iter()
2020 .filter_map(|t| {
2021 queries::resolve_id(&db.conn, queries::UidKind::Track, t)
2022 .ok()
2023 .flatten()
2024 })
2025 .collect();
2026 let row = |id: i64| {
2027 db.conn
2028 .query_row(
2029 "SELECT t.title, COALESCE(a.name, ''), COALESCE(al.title, ''), t.album_id
2030 FROM tracks t LEFT JOIN artists a ON a.id = t.artist_id
2031 LEFT JOIN albums al ON al.id = t.album_id WHERE t.id = ?1",
2032 [id],
2033 |r| {
2034 Ok((
2035 r.get::<_, String>(0)?,
2036 r.get::<_, String>(1)?,
2037 r.get::<_, String>(2)?,
2038 r.get::<_, Option<i64>>(3)?,
2039 ))
2040 },
2041 )
2042 .ok()
2043 };
2044 let (title, artist, album, album_id) = row(*ids.first()?)?;
2045 let one_album = ids.len() > 1
2046 && album_id.is_some()
2047 && ids
2048 .iter()
2049 .all(|id| row(*id).is_some_and(|r| r.3 == album_id));
2050 Some(match (one_album, ids.len()) {
2051 (true, _) => format!("{album} — {artist}"),
2052 (false, 1) => format!("{title} — {artist}"),
2053 (false, n) => format!("{title} — {artist}, and {} more", n - 1),
2054 })
2055 }
2056}
2057
2058const RECENT: i64 = 6 * 60 * 60;
2061
2062fn pick(clients: &[ClientInfo], now: i64) -> Result<&ClientInfo, String> {
2063 if let Some(c) = clients.iter().find(|c| c.state.playing) {
2064 return Ok(c);
2065 }
2066 if let Some(c) = clients
2067 .iter()
2068 .filter(|c| c.last_played_at.is_some_and(|t| now - t < RECENT))
2069 .max_by_key(|c| c.last_played_at)
2070 {
2071 return Ok(c);
2072 }
2073 match clients {
2074 [] => Err("no koan app is linked to this server; open koan on the device".into()),
2075 [only] => Ok(only),
2076 several => Err(format!(
2077 "several koan apps are linked and none has played recently: {}. Ask which, then pass `client`",
2078 several
2079 .iter()
2080 .map(|c| format!("{} ({}, id {})", c.name, c.platform, c.device))
2081 .collect::<Vec<_>>()
2082 .join(", ")
2083 )),
2084 }
2085}
2086
2087#[cfg(test)]
2088mod tests {
2089
2090 #[test]
2095 fn a_quiet_link_that_took_the_command_is_kept_and_waited_for() {
2096 use koan_core::remote::acks::AckOutcome;
2097 let reg = Registry::default();
2098 let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
2099 reg.register("q", "mac", "macos", "dev-quiet", tx, false, true);
2100 let answers = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
2101 let acking = |id: u64| {
2102 let answers = answers.clone();
2103 Acking {
2104 id,
2105 reply: Box::new(move |outcome| answers.lock().unwrap().push((id, outcome))),
2106 }
2107 };
2108
2109 reg.relay_acked("q", None, "dev-quiet", LinkCommand::Pause, Some(acking(1)))
2111 .unwrap();
2112 reg.received("dev-quiet", 1);
2113 assert!(reg.look("q", "dev-quiet", 1, &LinkCommand::Pause));
2114 assert!(answers.lock().unwrap().is_empty(), "still waiting");
2115 reg.answered("dev-quiet", 1, AckOutcome::Done);
2116 assert_eq!(
2117 answers.lock().unwrap().last(),
2118 Some(&(1, Some(AckOutcome::Done)))
2119 );
2120
2121 reg.relay_acked("q", None, "dev-quiet", LinkCommand::Pause, Some(acking(2)))
2124 .unwrap();
2125 assert!(!reg.look("q", "dev-quiet", 2, &LinkCommand::Pause));
2126 assert!(matches!(
2127 answers.lock().unwrap().last(),
2128 Some((2, Some(AckOutcome::Failed { .. })))
2129 ));
2130 assert_eq!(reg.list(Some("q")).len(), 1, "the link is not dropped");
2131 }
2132
2133 #[tokio::test]
2136 async fn an_answer_stops_its_timer() {
2137 let timer = tokio::spawn(tokio::time::sleep(std::time::Duration::from_secs(3600)));
2138 let handle = timer.abort_handle();
2139 let waiting = Waiting {
2140 reply: Box::new(|_| {}),
2141 received: true,
2142 timer: Some(timer.abort_handle()),
2143 };
2144 waiting.answer(Some(koan_core::remote::acks::AckOutcome::Done));
2145 assert!(timer.await.unwrap_err().is_cancelled());
2146 assert!(handle.is_finished());
2147 }
2148
2149 #[test]
2150 fn a_command_relayed_with_an_id_is_answered_once() {
2151 use koan_core::remote::acks::AckOutcome;
2152 let reg = Registry::default();
2153 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
2154 reg.register("a", "phone", "ios", "dev-answers", tx, false, true);
2155 let (old_tx, mut old_rx) = tokio::sync::mpsc::unbounded_channel();
2156 reg.register("a", "mac", "macos", "dev-old", old_tx, false, false);
2157 let answers = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
2158 let acking = |id: u64| {
2159 let answers = answers.clone();
2160 Acking {
2161 id,
2162 reply: Box::new(move |outcome| answers.lock().unwrap().push((id, outcome))),
2163 }
2164 };
2165
2166 reg.relay_acked(
2168 "a",
2169 None,
2170 "dev-answers",
2171 LinkCommand::Pause,
2172 Some(acking(7)),
2173 )
2174 .unwrap();
2175 assert_eq!(rx.try_recv().unwrap().ack, Some(7));
2176 reg.answered("dev-answers", 7, AckOutcome::Done);
2177 reg.answered("dev-answers", 7, AckOutcome::Done);
2178
2179 reg.relay_acked("a", None, "dev-old", LinkCommand::Pause, Some(acking(8)))
2182 .unwrap();
2183 assert_eq!(old_rx.try_recv().unwrap().ack, None);
2184
2185 assert_eq!(
2186 *answers.lock().unwrap(),
2187 [(7, Some(AckOutcome::Done)), (8, None)]
2188 );
2189 }
2190
2191 #[test]
2192 fn a_watch_is_never_queued_and_is_renewed_when_the_target_relinks() {
2193 let reg = Registry::default();
2194 let watches = |rx: &mut tokio::sync::mpsc::UnboundedReceiver<Envelope>| {
2195 std::iter::from_fn(|| rx.try_recv().ok().map(|e| e.command))
2196 .filter(|c| matches!(c, LinkCommand::WatchLevels { .. }))
2197 .collect::<Vec<_>>()
2198 };
2199 assert!(
2201 reg.relay("rl", "rl-phone", LinkCommand::WatchLevels { on: true })
2202 .is_err()
2203 );
2204 reg.watch_levels("rl", "rl-mac", "rl-phone", true);
2206 let (tx, mut phone) = tokio::sync::mpsc::unbounded_channel();
2207 reg.register("rl", "phone", "ios", "rl-phone", tx, false, false);
2208 assert_eq!(
2209 watches(&mut phone),
2210 vec![LinkCommand::WatchLevels { on: true }],
2211 "told once, on linking, because it is watched now"
2212 );
2213 let (tx, mut phone) = tokio::sync::mpsc::unbounded_channel();
2215 reg.register("rl", "phone", "ios", "rl-phone", tx, false, false);
2216 assert_eq!(
2217 watches(&mut phone),
2218 vec![LinkCommand::WatchLevels { on: true }]
2219 );
2220 }
2221
2222 #[test]
2223 fn levels_reach_a_watcher_only_while_it_watches() {
2224 use koan_core::remote::levels::Frame;
2225 let reg = Registry::default();
2226 let (tx_mac, mut mac) = tokio::sync::mpsc::unbounded_channel();
2227 let (tx_phone, mut phone) = tokio::sync::mpsc::unbounded_channel();
2228 let mac_id = reg.register("lv", "mac", "macos", "lv-mac", tx_mac, false, false);
2229 reg.register("lv", "phone", "ios", "lv-phone", tx_phone, false, false);
2230 let levels = |rx: &mut tokio::sync::mpsc::UnboundedReceiver<Envelope>| {
2231 std::iter::from_fn(|| rx.try_recv().ok().map(|e| e.command))
2232 .filter(|c| {
2233 matches!(
2234 c,
2235 LinkCommand::Levels { .. } | LinkCommand::WatchLevels { .. }
2236 )
2237 })
2238 .collect::<Vec<_>>()
2239 };
2240 let f = Frame(1_000, 1, 2, 3);
2241
2242 reg.levels("lv", "lv-phone", f);
2243 assert!(levels(&mut mac).is_empty(), "nobody watching");
2244
2245 reg.watch_levels("lv", "lv-mac", "lv-phone", true);
2246 assert_eq!(
2247 levels(&mut phone),
2248 vec![LinkCommand::WatchLevels { on: true }]
2249 );
2250 reg.levels("lv", "lv-phone", f);
2251 assert_eq!(
2252 levels(&mut mac),
2253 vec![LinkCommand::Levels {
2254 from: "lv-phone".into(),
2255 f
2256 }]
2257 );
2258
2259 reg.unregister(&mac_id);
2261 assert_eq!(
2262 levels(&mut phone),
2263 vec![LinkCommand::WatchLevels { on: false }]
2264 );
2265 reg.watch_levels("lv", "lv-mac", "lv-phone", true);
2266 reg.watch_levels("lv", "lv-mac", "lv-phone", false);
2267 assert_eq!(
2268 levels(&mut phone),
2269 vec![
2270 LinkCommand::WatchLevels { on: true },
2271 LinkCommand::WatchLevels { on: false }
2272 ]
2273 );
2274 }
2275
2276 #[test]
2277 fn a_playlist_order_adds_the_named_tracks_once_they_arrive() {
2278 let dir = tempfile::tempdir().unwrap();
2279 let path = dir.path().join("koan.db");
2280 let db = koan_core::db::connection::Database::open(&path).unwrap();
2281 let playlist = koan_core::db::queries::create_playlist(
2282 &db.conn,
2283 koan_core::db::queries::LOCAL_USER,
2284 "cyberpunk",
2285 None,
2286 )
2287 .unwrap();
2288 let order = Order {
2289 id: "o1".into(),
2290 username: None,
2291 client: None,
2292 artist: "Perturbator".into(),
2293 album: "Dangerous Days".into(),
2294 play_next: false,
2295 playlist: Some(playlist),
2296 titles: vec!["Future Club".into()],
2297 created_at: chrono::Utc::now().timestamp(),
2298 };
2299 registry().add_order(order);
2300
2301 fulfil_from(&path);
2303 assert!(registry().orders(None).iter().any(|o| o.id == "o1"));
2304
2305 db.conn
2306 .execute_batch(
2307 "INSERT INTO artists (id, name) VALUES (1, 'Perturbator');
2308 INSERT INTO albums (id, title, artist_id) VALUES (1, 'Dangerous Days', 1);
2309 INSERT INTO tracks (id, title, album_id, artist_id, track_number, path) VALUES
2310 (1, 'Welcome Back', 1, 1, 1, '/1.flac'), (2, 'Future Club', 1, 1, 2, '/2.flac');",
2311 )
2312 .unwrap();
2313 fulfil_from(&path);
2314 let held: Vec<i64> = db
2315 .conn
2316 .prepare("SELECT track_id FROM playlist_tracks WHERE playlist_id = ?1")
2317 .unwrap()
2318 .query_map([playlist], |r| r.get(0))
2319 .unwrap()
2320 .collect::<Result<_, _>>()
2321 .unwrap();
2322 assert_eq!(held, [2]);
2323 assert!(!registry().orders(None).iter().any(|o| o.id == "o1"));
2324 }
2325
2326 use super::*;
2327
2328 #[test]
2329 fn a_notification_cover_is_named_by_the_uid_the_command_carries() {
2330 let uid = "0190a5b2-7c3d-7e4f-8a1b-2c3d4e5f6a7b".to_string();
2332 let play = LinkCommand::Play {
2333 track_ids: vec!["x".into(), uid.clone()],
2334 start_at: 1,
2335 position_ms: 0,
2336 paused: false,
2337 handoff: false,
2338 };
2339 assert_eq!(cover_track(&play), Some(uid.as_str()));
2340 }
2341
2342 #[test]
2343 fn a_reconnect_replaces_the_device_and_commands_reach_it() {
2344 let reg = Registry::default();
2345 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
2346 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
2347 let (tx3, _rx3) = tokio::sync::mpsc::unbounded_channel();
2348 reg.register("j", "phone", "ios", "dev-1", tx1, false, false);
2349 let id = reg.register("j", "phone", "ios", "dev-1", tx2, false, false);
2350 reg.register("someone", "laptop", "macos", "dev-2", tx3, false, false);
2351
2352 assert_eq!(reg.list(Some("j")).len(), 1);
2353 assert_eq!(reg.list(None).len(), 2);
2354
2355 let sent = reg.send(Some("j"), None, LinkCommand::Pause).unwrap();
2356 assert_eq!(sent.id, id);
2357 assert_eq!(rx2.try_recv().unwrap().command, LinkCommand::Pause);
2358
2359 assert!(
2361 reg.send(Some("j"), Some("laptop"), LinkCommand::Pause)
2362 .is_err()
2363 );
2364
2365 reg.unregister(&id);
2366 assert!(reg.send(Some("j"), None, LinkCommand::Pause).is_err());
2367 }
2368
2369 #[test]
2370 fn disconnecting_an_account_closes_only_its_links() {
2371 let reg = Registry::default();
2372 let (tx1, mut rx1) = tokio::sync::mpsc::unbounded_channel();
2373 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
2374 reg.register("j", "phone", "ios", "dev-1", tx1, false, false);
2375 reg.register("someone", "laptop", "macos", "dev-2", tx2, false, false);
2376
2377 reg.disconnect("j");
2378 assert!(reg.list(Some("j")).is_empty());
2379 assert!(matches!(
2381 rx1.try_recv(),
2382 Err(tokio::sync::mpsc::error::TryRecvError::Disconnected)
2383 ));
2384 assert_eq!(reg.list(Some("someone")).len(), 1);
2385 assert!(matches!(
2386 rx2.try_recv(),
2387 Err(tokio::sync::mpsc::error::TryRecvError::Empty)
2388 ));
2389 }
2390
2391 #[test]
2392 fn a_burst_of_changes_wakes_each_device_once_when_it_goes_quiet() {
2393 let t0 = std::time::Instant::now();
2394 let s = std::time::Duration::from_secs;
2395 let phone = || ("dev-1".to_string(), "j".to_string());
2396 let ipad = || ("dev-2".to_string(), "j".to_string());
2397 let mut wakes = Wakes {
2398 pending: Vec::new(),
2399 timer: false,
2400 };
2401
2402 wakes.add([phone()], t0);
2403 wakes.add([phone(), ipad()], t0 + s(10));
2404 wakes.add([phone()], t0 + s(20));
2405 assert_eq!(wakes.pending.len(), 2);
2406
2407 assert!(wakes.take_due(t0 + s(39)).is_empty());
2409 assert_eq!(wakes.take_due(t0 + s(40)), [ipad()]);
2410 assert_eq!(wakes.next(), Some(t0 + s(50)));
2411 assert_eq!(wakes.take_due(t0 + s(50)), [phone()]);
2412 assert_eq!(wakes.next(), None);
2413 }
2414
2415 #[test]
2416 fn a_library_that_never_goes_quiet_still_wakes_devices() {
2417 let t0 = std::time::Instant::now();
2418 let phone = || ("dev-1".to_string(), "j".to_string());
2419 let mut wakes = Wakes {
2420 pending: Vec::new(),
2421 timer: false,
2422 };
2423 let mut sent = 0;
2424 for i in 0..40 {
2425 let now = t0 + std::time::Duration::from_secs(i * 10);
2426 sent += wakes.take_due(now).len();
2427 wakes.add([phone()], now);
2428 }
2429 assert_eq!(sent, 1);
2432 assert_eq!(wakes.pending.len(), 1);
2433 }
2434
2435 #[test]
2436 fn the_device_playing_is_the_one_meant() {
2437 let reg = Registry::default();
2438 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
2439 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
2440 let mac = reg.register("j", "mac", "macos", "dev-1", tx1, false, false);
2441 let phone = reg.register("j", "phone", "ios", "dev-2", tx2, false, false);
2442
2443 let err = reg.send(Some("j"), None, LinkCommand::Pause).unwrap_err();
2445 assert!(err.contains("mac") && err.contains("phone"), "{err}");
2446
2447 reg.report(
2448 &phone,
2449 LinkState {
2450 playing: true,
2451 ..Default::default()
2452 },
2453 );
2454 assert_eq!(
2455 reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
2456 phone
2457 );
2458 assert_eq!(rx2.try_recv().unwrap().command, LinkCommand::Pause);
2459
2460 reg.report(&phone, LinkState::default());
2462 assert_eq!(
2463 reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
2464 phone
2465 );
2466 let _ = mac;
2467 }
2468
2469 fn drain(rx: &mut tokio::sync::mpsc::UnboundedReceiver<Envelope>) -> Vec<LinkCommand> {
2470 std::iter::from_fn(|| rx.try_recv().ok().map(|e| e.command)).collect()
2471 }
2472
2473 fn shared() -> (
2475 Registry,
2476 tokio::sync::mpsc::UnboundedReceiver<Envelope>,
2477 tokio::sync::mpsc::UnboundedReceiver<Envelope>,
2478 tokio::sync::mpsc::UnboundedReceiver<Envelope>,
2479 ) {
2480 let reg = Registry::default();
2481 let (phone_tx, mut phone) = tokio::sync::mpsc::unbounded_channel();
2482 let (mac_tx, mac) = tokio::sync::mpsc::unbounded_channel();
2483 let (k_tx, k) = tokio::sync::mpsc::unbounded_channel();
2484 let id = reg.register("j", "phone", "ios", "dev-phone", phone_tx, true, false);
2485 reg.register("j", "mac", "macos", "dev-mac", mac_tx, true, false);
2486 reg.register("k", "laptop", "macos", "dev-k", k_tx, true, false);
2487 reg.report(
2488 &id,
2489 LinkState {
2490 playing: true,
2491 title: Some("Roygbiv".into()),
2492 outputs: Some(Default::default()),
2493 ..Default::default()
2494 },
2495 );
2496 reg.share("j", "dev-phone", "k", true).unwrap();
2497 drain(&mut phone);
2498 (reg, phone, mac, k)
2499 }
2500
2501 #[test]
2502 fn a_grantee_sees_the_shared_device_and_nothing_else_of_the_owner() {
2503 let (_reg, _phone, _mac, mut k) = shared();
2504 let listed = drain(&mut k)
2505 .into_iter()
2506 .rev()
2507 .find_map(|c| match c {
2508 LinkCommand::Devices { devices } => Some(devices),
2509 _ => None,
2510 })
2511 .expect("told of it");
2512 assert_eq!(listed.len(), 1, "the phone, not the Mac");
2513 let phone = &listed[0];
2514 assert_eq!(phone.id, "dev-phone");
2515 assert_eq!(phone.owner.as_deref(), Some("j"));
2516 let state = phone.state.as_ref().expect("its state");
2517 assert_eq!(state.title.as_deref(), Some("Roygbiv"));
2518 assert!(
2519 state.outputs.is_some(),
2520 "outputs, to choose from: the output is in the playback set"
2521 );
2522 }
2523
2524 #[test]
2530 fn a_grantee_controls_playback_as_itself_and_nothing_of_the_owners_account() {
2531 let (reg, mut phone, mut mac, _k) = shared();
2532 let wrapped = |c: LinkCommand| LinkCommand::Shared {
2533 command: Box::new(c),
2534 };
2535 for cmd in [
2536 LinkCommand::Pause,
2537 LinkCommand::SetRendererVolume { volume: 40 },
2538 LinkCommand::SleepTimer {
2539 timer: Some(koan_core::player::state::SleepTimer::After { minutes: 30 }),
2540 },
2541 LinkCommand::HandOff { to: "dev-k".into() },
2542 ] {
2543 reg.relay("k", "dev-phone", cmd.clone()).unwrap();
2544 assert_eq!(drain(&mut phone), [wrapped(cmd)]);
2545 }
2546 for cmd in [
2547 LinkCommand::Sync { full: false },
2548 LinkCommand::Evict { track_ids: vec![] },
2549 LinkCommand::Shared {
2550 command: Box::new(LinkCommand::Pause),
2551 },
2552 ] {
2553 assert!(!cmd.allowed_playback());
2554 assert!(reg.relay("k", "dev-phone", cmd).is_err());
2555 }
2556 assert!(drain(&mut phone).is_empty());
2557 assert!(
2558 reg.relay("k", "dev-mac", LinkCommand::Pause).is_err(),
2559 "not shared, not reachable"
2560 );
2561 assert!(
2562 drain(&mut mac)
2563 .iter()
2564 .all(|c| matches!(c, LinkCommand::Devices { .. } | LinkCommand::Shares { .. })),
2565 "news, and no command"
2566 );
2567 assert!(reg.shared_owner("k", "dev-mac").is_none(), "nor wakeable");
2568 }
2569
2570 #[test]
2574 fn a_hand_off_runs_both_ways_across_a_grant_and_nothing_else_does() {
2575 let (reg, _phone, _mac, mut k) = shared();
2576 drain(&mut k);
2577 let play = LinkCommand::Play {
2578 track_ids: vec!["t".into()],
2579 start_at: 0,
2580 position_ms: 1000,
2581 paused: false,
2582 handoff: true,
2583 };
2584 reg.relay_from("j", Some("dev-phone"), "dev-k", play.clone())
2585 .unwrap();
2586 assert!(drain(&mut k).contains(&LinkCommand::Shared {
2587 command: Box::new(play.clone())
2588 }));
2589 assert!(
2590 reg.relay_from("j", Some("dev-phone"), "dev-k", LinkCommand::Pause)
2591 .is_err(),
2592 "the owner's phone does not command the grantee's devices"
2593 );
2594 assert!(
2595 reg.relay_from("j", Some("dev-mac"), "dev-k", play.clone())
2596 .is_err(),
2597 "only the shared device"
2598 );
2599 let not_a_hand_off = LinkCommand::Play {
2600 track_ids: vec!["t".into()],
2601 start_at: 0,
2602 position_ms: 0,
2603 paused: false,
2604 handoff: false,
2605 };
2606 assert!(
2607 reg.relay_from("j", Some("dev-phone"), "dev-k", not_a_hand_off)
2608 .is_err(),
2609 "a hand-off, not any play"
2610 );
2611 }
2612
2613 #[test]
2614 fn revoking_ends_control_at_once() {
2615 let (reg, mut phone, _mac, mut k) = shared();
2616 reg.share("j", "dev-phone", "k", false).unwrap();
2617 assert!(reg.relay("k", "dev-phone", LinkCommand::Pause).is_err());
2618 assert!(
2619 drain(&mut phone)
2620 .iter()
2621 .all(|c| !matches!(c, LinkCommand::Pause))
2622 );
2623 assert!(reg.shared_owner("k", "dev-phone").is_none());
2624 let listed = drain(&mut k)
2625 .into_iter()
2626 .rev()
2627 .find_map(|c| match c {
2628 LinkCommand::Devices { devices } => Some(devices),
2629 _ => None,
2630 })
2631 .expect("told it is gone");
2632 assert!(listed.is_empty());
2633 }
2634
2635 #[test]
2636 fn a_device_hears_whom_it_is_shared_with_every_time_it_links() {
2637 let reg = Registry::default();
2638 let (tx, mut phone) = tokio::sync::mpsc::unbounded_channel();
2639 reg.register("j", "phone", "ios", "dev-phone", tx, true, false);
2640 let first = drain(&mut phone);
2641 assert!(
2642 first.contains(&LinkCommand::Shares {
2643 grantees: vec![],
2644 error: None,
2645 accounts: vec![],
2646 }),
2647 "an empty list replaces one from another server"
2648 );
2649 reg.send_shares(
2650 "j",
2651 "dev-phone",
2652 Some("There is no account called x".into()),
2653 );
2654 assert!(drain(&mut phone).iter().any(
2655 |c| matches!(c, LinkCommand::Shares { error: Some(e), .. } if e.contains("no account"))
2656 ));
2657 }
2658
2659 #[test]
2660 fn the_owner_is_told_who_it_shares_with() {
2661 let reg = Registry::default();
2662 let (tx, mut phone) = tokio::sync::mpsc::unbounded_channel();
2663 reg.register("j", "phone", "ios", "dev-phone", tx, true, false);
2664 reg.share("j", "dev-phone", "k", true).unwrap();
2665 reg.share("j", "dev-phone", "m", true).unwrap();
2666 let last = drain(&mut phone).into_iter().rev().find_map(|c| match c {
2667 LinkCommand::Shares { grantees, .. } => Some(grantees),
2668 _ => None,
2669 });
2670 assert_eq!(last, Some(vec!["k".to_string(), "m".to_string()]));
2671 assert!(
2672 reg.share("j", "dev-phone", "j", true).is_err(),
2673 "not with itself"
2674 );
2675 }
2676
2677 #[test]
2681 fn a_device_is_woken_for_another_account_on_its_network_or_by_grant() {
2682 let reg = Registry::default();
2683 let home: std::net::IpAddr = "203.0.113.7".parse().unwrap();
2684 let away: std::net::IpAddr = "198.51.100.2".parse().unwrap();
2685 reg.seen_at("dev-ipad", "sarita", home);
2686 reg.seen_at("dev-phone", "admin", home);
2687 assert_eq!(
2688 reg.wake_owner("admin", "dev-phone", "dev-ipad"),
2689 "sarita",
2690 "same address: the iPad's own token"
2691 );
2692 reg.seen_at("dev-phone", "admin", away);
2693 assert_eq!(
2694 reg.wake_owner("admin", "dev-phone", "dev-ipad"),
2695 "admin",
2696 "elsewhere, with no grant: only admin's own, which it is not"
2697 );
2698 reg.share("sarita", "dev-ipad", "admin", true).unwrap();
2699 assert_eq!(
2700 reg.wake_owner("admin", "dev-phone", "dev-ipad"),
2701 "sarita",
2702 "shared: from anywhere"
2703 );
2704 reg.seen_at("dev-other", "mallory", home);
2706 assert_eq!(reg.wake_owner("admin", "dev-other", "dev-tv"), "admin");
2707 }
2708
2709 #[test]
2712 fn where_a_device_last_linked_from_survives_a_restart() {
2713 let conn = rusqlite::Connection::open_in_memory().unwrap();
2714 koan_core::db::schema::create_tables(&conn).unwrap();
2715 for (device, user, seen) in [("dev-ipad", "sarita", 1), ("dev-phone", "admin", 2)] {
2716 conn.execute(
2717 "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES (?1, ?2, ?1, 'ios', ?3)",
2718 rusqlite::params![device, user, seen],
2719 )
2720 .unwrap();
2721 }
2722 let home: std::net::IpAddr = "203.0.113.7".parse().unwrap();
2723 outbox::save_address_in(&conn, "dev-ipad", "sarita", home);
2724 outbox::save_address_in(&conn, "dev-phone", "admin", home);
2725
2726 let reg = Registry::default();
2728 *reg.addresses.lock() = outbox::load_addresses_in(&conn);
2729 assert_eq!(
2730 reg.wake_owner("admin", "dev-phone", "dev-ipad"),
2731 "sarita",
2732 "still woken through its own account after the restart"
2733 );
2734 }
2735
2736 #[test]
2740 fn a_forgotten_device_is_dropped_by_the_accounts_devices() {
2741 let reg = Registry::default();
2742 let (mac_tx, mut mac) = tokio::sync::mpsc::unbounded_channel();
2743 let (other_tx, mut other) = tokio::sync::mpsc::unbounded_channel();
2744 reg.register("j", "mac", "macos", "dev-mac", mac_tx, true, false);
2745 reg.register("k", "laptop", "macos", "dev-k", other_tx, true, false);
2746 while mac.try_recv().is_ok() {}
2747 while other.try_recv().is_ok() {}
2748
2749 reg.forget("j", "dev-phone").unwrap();
2750 let heard: Vec<LinkCommand> =
2751 std::iter::from_fn(|| mac.try_recv().ok().map(|e| e.command)).collect();
2752 assert!(heard.contains(&LinkCommand::Forgotten {
2753 device: "dev-phone".into()
2754 }));
2755 assert!(
2756 std::iter::from_fn(|| other.try_recv().ok().map(|e| e.command))
2757 .all(|c| !matches!(c, LinkCommand::Forgotten { .. })),
2758 "another account hears nothing of it"
2759 );
2760 assert!(
2761 reg.forget("j", "dev-mac").is_err(),
2762 "linked: it would be back"
2763 );
2764 assert!(
2765 reg.relay(
2766 "j",
2767 "dev-mac",
2768 LinkCommand::Forgotten { device: "x".into() }
2769 )
2770 .is_err(),
2771 "news from the server, not a command a device may send"
2772 );
2773 }
2774
2775 #[test]
2778 fn forgetting_a_shared_device_declines_the_share() {
2779 let (reg, mut phone, _mac, mut k) = shared();
2780 drain(&mut k);
2781 reg.forget("k", "dev-phone").unwrap();
2782 assert!(reg.shared_owner("k", "dev-phone").is_none());
2783 assert!(
2784 drain(&mut phone)
2785 .iter()
2786 .any(|c| matches!(c, LinkCommand::Shares { grantees, .. } if grantees.is_empty()))
2787 );
2788 let listed = drain(&mut k).into_iter().rev().find_map(|c| match c {
2789 LinkCommand::Devices { devices } => Some(devices),
2790 _ => None,
2791 });
2792 assert_eq!(listed.map(|d| d.len()), Some(0));
2793 }
2794
2795 #[test]
2798 fn forgetting_a_device_stops_pushes_to_it_and_only_for_its_account() {
2799 let conn = rusqlite::Connection::open_in_memory().unwrap();
2800 koan_core::db::schema::create_tables(&conn).unwrap();
2801 for user in ["j", "k"] {
2802 conn.execute(
2803 "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES ('dev-phone', ?1, 'phone', 'ios', 1)",
2804 [user],
2805 )
2806 .unwrap();
2807 conn.execute(
2808 "INSERT INTO link_push (device, username, token, sandbox, updated_at) VALUES ('dev-phone', ?1, 't', 0, 1)",
2809 [user],
2810 )
2811 .unwrap();
2812 }
2813 outbox::forget_device_in(&conn, "dev-phone", "k");
2814 assert_eq!(
2815 outbox::push_targets_in(&conn, Some("j")).len(),
2816 1,
2817 "k forgetting its own leaves j's alone"
2818 );
2819 outbox::forget_device_in(&conn, "dev-phone", "j");
2820 assert!(outbox::push_targets_in(&conn, Some("j")).is_empty());
2821 assert!(outbox::push_targets_in(&conn, None).is_empty());
2822 }
2823
2824 #[test]
2825 fn a_device_hears_its_peers_and_commands_reach_them_by_device_id() {
2826 let reg = Registry::default();
2827 let (tx1, mut rx1) = tokio::sync::mpsc::unbounded_channel();
2828 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
2829 let (tx3, mut rx3) = tokio::sync::mpsc::unbounded_channel();
2830 reg.register("j", "mac", "macos", "dev-mac", tx1, true, false);
2831 let phone = reg.register("j", "phone", "ios", "dev-phone", tx2, true, false);
2832 reg.register("someone", "laptop", "macos", "dev-other", tx3, true, false);
2833
2834 reg.report(
2835 &phone,
2836 LinkState {
2837 playing: true,
2838 title: Some("Roygbiv".into()),
2839 ..Default::default()
2840 },
2841 );
2842 let devices = drain(&mut rx1)
2843 .into_iter()
2844 .rev()
2845 .find_map(|c| match c {
2846 LinkCommand::Devices { devices } => Some(devices),
2847 _ => None,
2848 })
2849 .expect("the Mac is told of the phone");
2850 assert_eq!(devices.len(), 1, "not itself, not another account's");
2851 assert_eq!(devices[0].id, "dev-phone");
2852 assert_eq!(
2853 devices[0].state.as_ref().and_then(|s| s.title.as_deref()),
2854 Some("Roygbiv")
2855 );
2856 while rx3.try_recv().is_ok() {}
2857 assert!(rx3.try_recv().is_err());
2858
2859 while rx2.try_recv().is_ok() {}
2860 reg.relay("j", "dev-phone", LinkCommand::Pause).unwrap();
2861 assert_eq!(rx2.try_recv().unwrap().command, LinkCommand::Pause);
2862 assert!(
2863 reg.relay("someone", "dev-phone", LinkCommand::Pause)
2864 .is_err(),
2865 "another account cannot reach it"
2866 );
2867 assert!(
2868 reg.relay("j", "dev-phone", LinkCommand::Devices { devices: vec![] })
2869 .is_err()
2870 );
2871 }
2872
2873 #[test]
2874 fn a_device_linked_for_a_month_is_not_forgotten() {
2875 let conn = rusqlite::Connection::open_in_memory().unwrap();
2876 koan_core::db::schema::create_tables(&conn).unwrap();
2877 let now = 100 * 24 * 60 * 60;
2878 let long_ago = now - 40 * 24 * 60 * 60;
2879 for device in ["mac", "old-phone"] {
2880 conn.execute(
2881 "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES (?1, 'j', ?1, 'ios', ?2)",
2882 rusqlite::params![device, long_ago],
2883 )
2884 .unwrap();
2885 conn.execute(
2886 "INSERT INTO link_push (device, username, token, sandbox, updated_at) VALUES (?1, 'j', 't', 0, ?2)",
2887 rusqlite::params![device, long_ago],
2888 )
2889 .unwrap();
2890 }
2891 outbox::forget_stale(&conn, &[("mac".into(), "j".into())], now);
2892 let left = |table: &str| -> Vec<String> {
2893 conn.prepare(&format!("SELECT device FROM {table} ORDER BY device"))
2894 .unwrap()
2895 .query_map([], |r| r.get(0))
2896 .unwrap()
2897 .collect::<Result<_, _>>()
2898 .unwrap()
2899 };
2900 assert_eq!(left("link_devices"), ["mac"]);
2901 assert_eq!(left("link_push"), ["mac"]);
2902 }
2903}