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 if cmd.live_only() {
1387 return None;
1388 }
1389 let envelope = Envelope {
1390 command: cmd.clone(),
1391 ack,
1392 };
1393 let pusher = crate::push::pusher()?;
1394 let target = outbox::push_targets(username)
1397 .into_iter()
1398 .find(|t| id.is_none_or(|id| t.device == id || t.name.eq_ignore_ascii_case(id)))?;
1399 let info = ClientInfo {
1400 id: target.device.clone(),
1401 device: target.device.clone(),
1402 name: target.name.clone(),
1403 platform: target.platform.clone(),
1404 username: target.username.clone(),
1405 connected_at: 0,
1406 state: LinkState::default(),
1407 last_played_at: None,
1408 state_at: 0,
1409 reports: false,
1410 notified: true,
1411 acks: false,
1412 };
1413 let asked = match cmd {
1415 LinkCommand::Shared { command } => command.as_ref(),
1416 cmd => cmd,
1417 };
1418 let verb = match asked {
1419 LinkCommand::Play { .. } | LinkCommand::JumpTo { .. } | LinkCommand::Resume => Some("Play"),
1420 _ => None,
1421 };
1422 let push = match verb {
1423 Some(verb) => crate::push::Push::Notify {
1424 title: format!("{verb} on {}", target.name),
1425 body: outbox::describe(asked).unwrap_or_else(|| "From your koan server".into()),
1426 command: serde_json::to_value(&envelope).ok()?,
1427 image: cover_track(asked)
1428 .and_then(outbox::track_row)
1429 .and_then(|t| pusher.cover_link(t)),
1430 },
1431 None => {
1432 outbox::queue_for(&target.device, &target.username, &envelope);
1433 crate::push::Push::Wake
1434 }
1435 };
1436 std::thread::spawn(move || deliver_push(pusher, &target, &push));
1437 Some(Ok(info))
1438}
1439
1440fn cover_track(cmd: &LinkCommand) -> Option<&str> {
1443 match cmd {
1444 LinkCommand::Play {
1445 track_ids,
1446 start_at,
1447 ..
1448 } => track_ids
1449 .get(*start_at as usize)
1450 .or(track_ids.first())
1451 .map(String::as_str),
1452 LinkCommand::JumpTo { track_id } => Some(track_id),
1453 _ => None,
1454 }
1455}
1456
1457fn deliver_push(
1459 pusher: &crate::push::Pusher,
1460 target: &outbox::PushTarget,
1461 push: &crate::push::Push,
1462) {
1463 use crate::push::Outcome;
1464 match pusher.send(&target.token, target.sandbox, push) {
1465 Outcome::Sent => log::info!("push: sent to {}", target.name),
1466 Outcome::Gone => {
1467 log::info!(
1468 "push: {}'s token is no longer valid; forgotten",
1469 target.name
1470 );
1471 outbox::forget_push(&target.username, &target.device);
1472 }
1473 Outcome::Failed(e) => log::warn!("push: to {} failed: {e}", target.name),
1474 }
1475}
1476
1477type Woken = std::collections::HashMap<(String, String), (std::time::Instant, &'static str)>;
1481
1482static WOKEN: LazyLock<Mutex<Woken>> = LazyLock::new(Default::default);
1483
1484pub fn smart_activity(
1488 db: &koan_core::db::connection::Database,
1489 user: i64,
1490 fields: &[koan_core::smart::Field],
1491) {
1492 match koan_core::db::queries::smart::refresh_after_activity(&db.conn, user, fields) {
1493 Ok(moved) if !moved.is_empty() => changed(),
1494 Ok(_) => {}
1495 Err(e) => log::warn!("smart playlists not refreshed after activity: {e}"),
1496 }
1497}
1498
1499pub fn changed() {
1504 let (_, queued) = registry().link_or_queue(None, &LinkCommand::Sync { full: false });
1505 if queued.is_empty() || crate::push::pusher().is_none() {
1506 return;
1507 }
1508 let mut wakes = WAKES.lock();
1509 wakes.add(queued.iter().map(Absent::key), std::time::Instant::now());
1510 if !wakes.timer {
1511 wakes.timer = true;
1512 std::thread::spawn(send_wakes);
1513 }
1514}
1515
1516const QUIET: std::time::Duration = std::time::Duration::from_secs(30);
1520
1521const LONGEST_WAIT: std::time::Duration = std::time::Duration::from_secs(5 * 60);
1523
1524static WAKES: Mutex<Wakes> = Mutex::new(Wakes {
1525 pending: Vec::new(),
1526 timer: false,
1527});
1528
1529struct Wakes {
1532 pending: Vec<((String, String), std::time::Instant, std::time::Instant)>,
1534 timer: bool,
1536}
1537
1538impl Wakes {
1539 fn add(
1540 &mut self,
1541 devices: impl IntoIterator<Item = (String, String)>,
1542 now: std::time::Instant,
1543 ) {
1544 for key in devices {
1545 match self.pending.iter_mut().find(|(k, _, _)| *k == key) {
1546 Some((_, _, last)) => *last = now,
1547 None => self.pending.push((key, now, now)),
1548 }
1549 }
1550 }
1551
1552 fn due_at(first: std::time::Instant, last: std::time::Instant) -> std::time::Instant {
1553 (last + QUIET).min(first + LONGEST_WAIT)
1554 }
1555
1556 fn next(&self) -> Option<std::time::Instant> {
1558 self.pending
1559 .iter()
1560 .map(|(_, first, last)| Self::due_at(*first, *last))
1561 .min()
1562 }
1563
1564 fn take_due(&mut self, now: std::time::Instant) -> Vec<(String, String)> {
1566 let (due, waiting) = std::mem::take(&mut self.pending)
1567 .into_iter()
1568 .partition(|(_, first, last)| Self::due_at(*first, *last) <= now);
1569 self.pending = waiting;
1570 due.into_iter().map(|(key, _, _)| key).collect()
1571 }
1572}
1573
1574fn send_wakes() {
1577 loop {
1578 let (due, next) = {
1579 let mut wakes = WAKES.lock();
1580 let due = wakes.take_due(std::time::Instant::now());
1581 let next = wakes.next();
1582 if due.is_empty() && next.is_none() {
1583 wakes.timer = false;
1584 return;
1585 }
1586 (due, next)
1587 };
1588 let live = registry().live();
1589 let absent: Vec<(String, String)> = due.into_iter().filter(|d| !live.contains(d)).collect();
1590 if !absent.is_empty() {
1591 wake(&absent);
1592 }
1593 if let Some(next) = next {
1594 std::thread::sleep(next.saturating_duration_since(std::time::Instant::now()));
1595 }
1596 }
1597}
1598
1599pub fn changed_if_library_moved(conn: &rusqlite::Connection) {
1603 static LAST: parking_lot::Mutex<Option<(i64, i64, i64)>> = parking_lot::Mutex::new(None);
1604 let Ok(now) = conn.query_row(
1605 "SELECT (SELECT COUNT(*) FROM tracks), (SELECT COALESCE(MAX(id), 0) FROM tracks),
1606 (SELECT COUNT(*) FROM albums)",
1607 [],
1608 |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
1609 ) else {
1610 return;
1611 };
1612 let before = LAST.lock().replace(now);
1613 if before.is_some_and(|b| b != now) {
1614 changed();
1615 }
1616}
1617
1618mod outbox {
1621 use koan_core::db::queries;
1622 use koan_core::remote::link::LinkCommand;
1623
1624 const KEEP_SECS: i64 = 30 * 24 * 60 * 60;
1629
1630 fn db() -> Option<koan_core::db::pool::Handle<'static>> {
1635 if cfg!(test) {
1636 return None;
1637 }
1638 koan_core::db::pool::shared().get().ok()
1639 }
1640
1641 pub fn load_orders() -> Vec<super::Order> {
1642 let Some(db) = db() else { return Vec::new() };
1643 db.conn
1644 .prepare("SELECT body FROM link_orders ORDER BY created_at")
1645 .and_then(|mut s| {
1646 s.query_map([], |r| r.get::<_, String>(0))?
1647 .collect::<Result<Vec<_>, _>>()
1648 })
1649 .unwrap_or_default()
1650 .into_iter()
1651 .filter_map(|b| serde_json::from_str(&b).ok())
1652 .collect()
1653 }
1654
1655 pub fn save_order(order: &super::Order) {
1656 let (Some(db), Ok(body)) = (db(), serde_json::to_string(order)) else {
1657 return;
1658 };
1659 let _ = db.conn.execute(
1660 "INSERT OR REPLACE INTO link_orders (id, body, created_at) VALUES (?1, ?2, ?3)",
1661 rusqlite::params![order.id, body, order.created_at],
1662 );
1663 }
1664
1665 pub fn drop_order(id: &str) {
1666 if let Some(db) = db() {
1667 let _ = db
1668 .conn
1669 .execute("DELETE FROM link_orders WHERE id = ?1", [id]);
1670 }
1671 }
1672
1673 pub fn take_and_remember(
1674 username: &str,
1675 device: &str,
1676 name: &str,
1677 platform: &str,
1678 live: &[(String, String)],
1679 ) -> Vec<koan_core::remote::acks::Envelope> {
1680 let Some(db) = db() else { return Vec::new() };
1681 let now = chrono::Utc::now().timestamp();
1682 let _ = db.conn.execute(
1683 "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES (?1, ?2, ?3, ?4, ?5)
1684 ON CONFLICT (device, username) DO UPDATE SET name = ?3, platform = ?4, last_seen = ?5",
1685 rusqlite::params![device, username, name, platform, now],
1686 );
1687 forget_stale(&db.conn, live, now);
1688 let waiting: Vec<(i64, String)> = db
1689 .conn
1690 .prepare("SELECT id, command FROM link_outbox WHERE device = ?1 AND username = ?2 ORDER BY id")
1691 .and_then(|mut s| {
1692 s.query_map([device, username], |r| Ok((r.get(0)?, r.get(1)?)))?
1693 .collect()
1694 })
1695 .unwrap_or_default();
1696 let _ = db.conn.execute(
1697 "DELETE FROM link_outbox WHERE device = ?1 AND username = ?2",
1698 [device, username],
1699 );
1700 if !waiting.is_empty() {
1701 log::info!("link: {} waiting commands for {name}", waiting.len());
1702 }
1703 waiting
1704 .into_iter()
1705 .filter_map(|(_, c)| serde_json::from_str(&c).ok())
1706 .collect()
1707 }
1708
1709 pub(super) fn forget_stale(conn: &rusqlite::Connection, live: &[(String, String)], now: i64) {
1712 touch_with(conn, live, now);
1713 let _ = conn.execute(
1714 "DELETE FROM link_outbox WHERE created_at < ?1",
1715 [now - KEEP_SECS],
1716 );
1717 let _ = conn.execute(
1718 "DELETE FROM link_push WHERE (device, username) IN
1719 (SELECT device, username FROM link_devices WHERE last_seen < ?1)",
1720 [now - KEEP_SECS],
1721 );
1722 let _ = conn.execute(
1723 "DELETE FROM link_devices WHERE last_seen < ?1",
1724 [now - KEEP_SECS],
1725 );
1726 }
1727
1728 pub fn touch(devices: &[(String, String)]) {
1730 if let Some(db) = db() {
1731 touch_with(&db.conn, devices, chrono::Utc::now().timestamp());
1732 }
1733 }
1734
1735 fn touch_with(conn: &rusqlite::Connection, devices: &[(String, String)], now: i64) {
1736 for (device, username) in devices {
1737 let _ = conn.execute(
1738 "UPDATE link_devices SET last_seen = ?1 WHERE device = ?2 AND username = ?3",
1739 rusqlite::params![now, device, username],
1740 );
1741 }
1742 }
1743
1744 pub struct Absent {
1746 pub device: String,
1747 pub username: String,
1748 pub name: String,
1749 }
1750
1751 pub fn queue_for_absent(
1753 username: Option<&str>,
1754 live: &[(String, String)],
1755 cmd: &LinkCommand,
1756 ) -> Vec<Absent> {
1757 let Some(db) = db() else { return Vec::new() };
1758 let known: Vec<(String, String, String)> = db
1759 .conn
1760 .prepare("SELECT device, username, name FROM link_devices")
1761 .and_then(|mut s| {
1762 s.query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))?
1763 .collect()
1764 })
1765 .unwrap_or_default();
1766 let Ok(text) = serde_json::to_string(cmd) else {
1767 return Vec::new();
1768 };
1769 let absent: Vec<_> = known
1770 .into_iter()
1771 .filter(|(device, user, _)| {
1772 !username.is_some_and(|u| u != user)
1773 && !live.iter().any(|(d, u)| d == device && u == user)
1774 })
1775 .collect();
1776 if absent.is_empty() {
1777 return Vec::new();
1778 }
1779 let is_sync = matches!(cmd, LinkCommand::Sync { .. });
1780 let now = chrono::Utc::now().timestamp();
1781 koan_core::db::queries::atomically(&db.conn, || {
1784 let mut queued = Vec::new();
1785 for (device, user, name) in absent {
1786 if is_sync {
1787 let _ = db.conn.execute(
1789 "DELETE FROM link_outbox WHERE device = ?1 AND username = ?2 AND command LIKE '{\"type\":\"sync\"%'",
1790 [&device, &user],
1791 );
1792 }
1793 if db
1794 .conn
1795 .execute(
1796 "INSERT INTO link_outbox (device, username, command, created_at) VALUES (?1, ?2, ?3, ?4)",
1797 rusqlite::params![device, user, text, now],
1798 )
1799 .is_ok()
1800 {
1801 queued.push(Absent {
1802 device,
1803 username: user,
1804 name,
1805 });
1806 }
1807 }
1808 Ok::<_, rusqlite::Error>(queued)
1809 })
1810 .unwrap_or_default()
1811 }
1812
1813 #[derive(Clone)]
1815 pub struct PushTarget {
1816 pub device: String,
1817 pub username: String,
1818 pub name: String,
1819 pub platform: String,
1820 pub token: String,
1821 pub sandbox: bool,
1822 pub last_seen: i64,
1824 }
1825
1826 pub fn queue_for(device: &str, username: &str, envelope: &koan_core::remote::acks::Envelope) {
1830 let (Some(db), Ok(text)) = (db(), serde_json::to_string(envelope)) else {
1831 return;
1832 };
1833 let _ = db.conn.execute(
1834 "INSERT INTO link_outbox (device, username, command, created_at) VALUES (?1, ?2, ?3, ?4)",
1835 rusqlite::params![device, username, text, chrono::Utc::now().timestamp()],
1836 );
1837 }
1838
1839 pub fn accounts() -> Vec<String> {
1843 let Some(db) = db() else { return Vec::new() };
1844 koan_core::db::queries::auth::list_users(&db.conn)
1845 .map(|users| users.into_iter().map(|u| u.username).collect())
1846 .unwrap_or_default()
1847 }
1848
1849 pub fn load_addresses() -> std::collections::HashMap<String, (String, std::net::IpAddr)> {
1852 let Some(db) = db() else {
1853 return Default::default();
1854 };
1855 load_addresses_in(&db.conn)
1856 }
1857
1858 pub(super) fn load_addresses_in(
1859 conn: &rusqlite::Connection,
1860 ) -> std::collections::HashMap<String, (String, std::net::IpAddr)> {
1861 conn.prepare(
1862 "SELECT device, username, addr FROM link_devices WHERE addr IS NOT NULL ORDER BY last_seen",
1863 )
1864 .and_then(|mut s| {
1865 s.query_map([], |r| {
1866 Ok((
1867 r.get::<_, String>(0)?,
1868 r.get::<_, String>(1)?,
1869 r.get::<_, String>(2)?,
1870 ))
1871 })?
1872 .collect::<Result<Vec<_>, _>>()
1873 })
1874 .unwrap_or_default()
1875 .into_iter()
1876 .filter_map(|(device, username, addr)| Some((device, (username, addr.parse().ok()?))))
1877 .collect()
1878 }
1879
1880 pub fn save_address(device: &str, username: &str, addr: std::net::IpAddr) {
1881 if let Some(db) = db() {
1882 save_address_in(&db.conn, device, username, addr);
1883 }
1884 }
1885
1886 pub(super) fn save_address_in(
1887 conn: &rusqlite::Connection,
1888 device: &str,
1889 username: &str,
1890 addr: std::net::IpAddr,
1891 ) {
1892 let _ = conn.execute(
1893 "UPDATE link_devices SET addr = ?1 WHERE device = ?2 AND username = ?3",
1894 rusqlite::params![addr.to_string(), device, username],
1895 );
1896 }
1897
1898 pub fn load_grants() -> Vec<super::Grant> {
1899 let Some(db) = db() else { return Vec::new() };
1900 db.conn
1901 .prepare("SELECT device, owner, grantee FROM link_grants ORDER BY created_at")
1902 .and_then(|mut s| {
1903 s.query_map([], |r| {
1904 Ok(super::Grant {
1905 device: r.get(0)?,
1906 owner: r.get(1)?,
1907 grantee: r.get(2)?,
1908 })
1909 })?
1910 .collect()
1911 })
1912 .unwrap_or_default()
1913 }
1914
1915 pub fn save_grant(g: &super::Grant) {
1916 if let Some(db) = db() {
1917 let _ = db.conn.execute(
1918 "INSERT OR IGNORE INTO link_grants (device, owner, grantee, created_at) VALUES (?1, ?2, ?3, ?4)",
1919 rusqlite::params![g.device, g.owner, g.grantee, chrono::Utc::now().timestamp()],
1920 );
1921 }
1922 }
1923
1924 pub fn drop_grant(g: &super::Grant) {
1925 if let Some(db) = db() {
1926 let _ = db.conn.execute(
1927 "DELETE FROM link_grants WHERE device = ?1 AND owner = ?2 AND grantee = ?3",
1928 [&g.device, &g.owner, &g.grantee],
1929 );
1930 }
1931 }
1932
1933 pub fn save_push(username: &str, device: &str, token: &str, sandbox: bool) {
1934 let Some(db) = db() else { return };
1935 let _ = db.conn.execute(
1936 "INSERT INTO link_push (device, username, token, sandbox, updated_at) VALUES (?1, ?2, ?3, ?4, ?5)
1937 ON CONFLICT (device, username) DO UPDATE SET token = ?3, sandbox = ?4, updated_at = ?5",
1938 rusqlite::params![device, username, token, sandbox, chrono::Utc::now().timestamp()],
1939 );
1940 }
1941
1942 pub fn forget_push(username: &str, device: &str) {
1943 if let Some(db) = db() {
1944 let _ = db.conn.execute(
1945 "DELETE FROM link_push WHERE device = ?1 AND username = ?2",
1946 [device, username],
1947 );
1948 }
1949 }
1950
1951 pub fn push_targets(username: Option<&str>) -> Vec<PushTarget> {
1953 let Some(db) = db() else { return Vec::new() };
1954 push_targets_in(&db.conn, username)
1955 }
1956
1957 pub(super) fn push_targets_in(
1958 conn: &rusqlite::Connection,
1959 username: Option<&str>,
1960 ) -> Vec<PushTarget> {
1961 conn
1962 .prepare(
1963 "SELECT p.device, p.username, d.name, d.platform, p.token, p.sandbox, d.last_seen
1964 FROM link_push p JOIN link_devices d ON d.device = p.device AND d.username = p.username
1965 WHERE ?1 IS NULL OR p.username = ?1
1966 ORDER BY d.last_seen DESC",
1967 )
1968 .and_then(|mut s| {
1969 s.query_map([username], |r| {
1970 Ok(PushTarget {
1971 device: r.get(0)?,
1972 username: r.get(1)?,
1973 name: r.get(2)?,
1974 platform: r.get(3)?,
1975 token: r.get(4)?,
1976 sandbox: r.get(5)?,
1977 last_seen: r.get(6)?,
1978 })
1979 })?
1980 .collect()
1981 })
1982 .unwrap_or_default()
1983 }
1984
1985 pub fn forget_device(device: &str, username: &str) {
1988 if let Some(db) = db() {
1989 forget_device_in(&db.conn, device, username);
1990 }
1991 }
1992
1993 pub(super) fn forget_device_in(conn: &rusqlite::Connection, device: &str, username: &str) {
1994 for table in ["link_push", "link_outbox", "link_devices"] {
1995 let _ = conn.execute(
1996 &format!("DELETE FROM {table} WHERE device = ?1 AND username = ?2"),
1997 [device, username],
1998 );
1999 }
2000 }
2001
2002 pub fn track_row(id: &str) -> Option<i64> {
2004 let db = db()?;
2005 queries::resolve_id(&db.conn, queries::UidKind::Track, id)
2006 .ok()
2007 .flatten()
2008 }
2009
2010 pub fn describe(cmd: &LinkCommand) -> Option<String> {
2013 let ids: Vec<&String> = match cmd {
2014 LinkCommand::Play { track_ids, .. }
2015 | LinkCommand::Enqueue { track_ids }
2016 | LinkCommand::PlayNext { track_ids } => track_ids.iter().collect(),
2017 LinkCommand::JumpTo { track_id } => vec![track_id],
2018 _ => return None,
2019 };
2020 let db = db()?;
2021 let ids: Vec<i64> = ids
2022 .into_iter()
2023 .filter_map(|t| {
2024 queries::resolve_id(&db.conn, queries::UidKind::Track, t)
2025 .ok()
2026 .flatten()
2027 })
2028 .collect();
2029 let row = |id: i64| {
2030 db.conn
2031 .query_row(
2032 "SELECT t.title, COALESCE(a.name, ''), COALESCE(al.title, ''), t.album_id
2033 FROM tracks t LEFT JOIN artists a ON a.id = t.artist_id
2034 LEFT JOIN albums al ON al.id = t.album_id WHERE t.id = ?1",
2035 [id],
2036 |r| {
2037 Ok((
2038 r.get::<_, String>(0)?,
2039 r.get::<_, String>(1)?,
2040 r.get::<_, String>(2)?,
2041 r.get::<_, Option<i64>>(3)?,
2042 ))
2043 },
2044 )
2045 .ok()
2046 };
2047 let (title, artist, album, album_id) = row(*ids.first()?)?;
2048 let one_album = ids.len() > 1
2049 && album_id.is_some()
2050 && ids
2051 .iter()
2052 .all(|id| row(*id).is_some_and(|r| r.3 == album_id));
2053 Some(match (one_album, ids.len()) {
2054 (true, _) => format!("{album} — {artist}"),
2055 (false, 1) => format!("{title} — {artist}"),
2056 (false, n) => format!("{title} — {artist}, and {} more", n - 1),
2057 })
2058 }
2059}
2060
2061const RECENT: i64 = 6 * 60 * 60;
2064
2065fn pick(clients: &[ClientInfo], now: i64) -> Result<&ClientInfo, String> {
2066 if let Some(c) = clients.iter().find(|c| c.state.playing) {
2067 return Ok(c);
2068 }
2069 if let Some(c) = clients
2070 .iter()
2071 .filter(|c| c.last_played_at.is_some_and(|t| now - t < RECENT))
2072 .max_by_key(|c| c.last_played_at)
2073 {
2074 return Ok(c);
2075 }
2076 match clients {
2077 [] => Err("no koan app is linked to this server; open koan on the device".into()),
2078 [only] => Ok(only),
2079 several => Err(format!(
2080 "several koan apps are linked and none has played recently: {}. Ask which, then pass `client`",
2081 several
2082 .iter()
2083 .map(|c| format!("{} ({}, id {})", c.name, c.platform, c.device))
2084 .collect::<Vec<_>>()
2085 .join(", ")
2086 )),
2087 }
2088}
2089
2090#[cfg(test)]
2091mod tests {
2092
2093 #[test]
2098 fn a_quiet_link_that_took_the_command_is_kept_and_waited_for() {
2099 use koan_core::remote::acks::AckOutcome;
2100 let reg = Registry::default();
2101 let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
2102 reg.register("q", "mac", "macos", "dev-quiet", tx, false, true);
2103 let answers = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
2104 let acking = |id: u64| {
2105 let answers = answers.clone();
2106 Acking {
2107 id,
2108 reply: Box::new(move |outcome| answers.lock().unwrap().push((id, outcome))),
2109 }
2110 };
2111
2112 reg.relay_acked("q", None, "dev-quiet", LinkCommand::Pause, Some(acking(1)))
2114 .unwrap();
2115 reg.received("dev-quiet", 1);
2116 assert!(reg.look("q", "dev-quiet", 1, &LinkCommand::Pause));
2117 assert!(answers.lock().unwrap().is_empty(), "still waiting");
2118 reg.answered("dev-quiet", 1, AckOutcome::Done);
2119 assert_eq!(
2120 answers.lock().unwrap().last(),
2121 Some(&(1, Some(AckOutcome::Done)))
2122 );
2123
2124 reg.relay_acked("q", None, "dev-quiet", LinkCommand::Pause, Some(acking(2)))
2127 .unwrap();
2128 assert!(!reg.look("q", "dev-quiet", 2, &LinkCommand::Pause));
2129 assert!(matches!(
2130 answers.lock().unwrap().last(),
2131 Some((2, Some(AckOutcome::Failed { .. })))
2132 ));
2133 assert_eq!(reg.list(Some("q")).len(), 1, "the link is not dropped");
2134 }
2135
2136 #[tokio::test]
2139 async fn an_answer_stops_its_timer() {
2140 let timer = tokio::spawn(tokio::time::sleep(std::time::Duration::from_secs(3600)));
2141 let handle = timer.abort_handle();
2142 let waiting = Waiting {
2143 reply: Box::new(|_| {}),
2144 received: true,
2145 timer: Some(timer.abort_handle()),
2146 };
2147 waiting.answer(Some(koan_core::remote::acks::AckOutcome::Done));
2148 assert!(timer.await.unwrap_err().is_cancelled());
2149 assert!(handle.is_finished());
2150 }
2151
2152 #[test]
2153 fn a_command_relayed_with_an_id_is_answered_once() {
2154 use koan_core::remote::acks::AckOutcome;
2155 let reg = Registry::default();
2156 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
2157 reg.register("a", "phone", "ios", "dev-answers", tx, false, true);
2158 let (old_tx, mut old_rx) = tokio::sync::mpsc::unbounded_channel();
2159 reg.register("a", "mac", "macos", "dev-old", old_tx, false, false);
2160 let answers = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
2161 let acking = |id: u64| {
2162 let answers = answers.clone();
2163 Acking {
2164 id,
2165 reply: Box::new(move |outcome| answers.lock().unwrap().push((id, outcome))),
2166 }
2167 };
2168
2169 reg.relay_acked(
2171 "a",
2172 None,
2173 "dev-answers",
2174 LinkCommand::Pause,
2175 Some(acking(7)),
2176 )
2177 .unwrap();
2178 assert_eq!(rx.try_recv().unwrap().ack, Some(7));
2179 reg.answered("dev-answers", 7, AckOutcome::Done);
2180 reg.answered("dev-answers", 7, AckOutcome::Done);
2181
2182 reg.relay_acked("a", None, "dev-old", LinkCommand::Pause, Some(acking(8)))
2185 .unwrap();
2186 assert_eq!(old_rx.try_recv().unwrap().ack, None);
2187
2188 assert_eq!(
2189 *answers.lock().unwrap(),
2190 [(7, Some(AckOutcome::Done)), (8, None)]
2191 );
2192 }
2193
2194 #[test]
2195 fn a_watch_is_never_queued_and_is_renewed_when_the_target_relinks() {
2196 let reg = Registry::default();
2197 let watches = |rx: &mut tokio::sync::mpsc::UnboundedReceiver<Envelope>| {
2198 std::iter::from_fn(|| rx.try_recv().ok().map(|e| e.command))
2199 .filter(|c| matches!(c, LinkCommand::WatchLevels { .. }))
2200 .collect::<Vec<_>>()
2201 };
2202 assert!(
2204 reg.relay("rl", "rl-phone", LinkCommand::WatchLevels { on: true })
2205 .is_err()
2206 );
2207 reg.watch_levels("rl", "rl-mac", "rl-phone", true);
2209 let (tx, mut phone) = tokio::sync::mpsc::unbounded_channel();
2210 reg.register("rl", "phone", "ios", "rl-phone", tx, false, false);
2211 assert_eq!(
2212 watches(&mut phone),
2213 vec![LinkCommand::WatchLevels { on: true }],
2214 "told once, on linking, because it is watched now"
2215 );
2216 let (tx, mut phone) = tokio::sync::mpsc::unbounded_channel();
2218 reg.register("rl", "phone", "ios", "rl-phone", tx, false, false);
2219 assert_eq!(
2220 watches(&mut phone),
2221 vec![LinkCommand::WatchLevels { on: true }]
2222 );
2223 }
2224
2225 #[test]
2226 fn levels_reach_a_watcher_only_while_it_watches() {
2227 use koan_core::remote::levels::Frame;
2228 let reg = Registry::default();
2229 let (tx_mac, mut mac) = tokio::sync::mpsc::unbounded_channel();
2230 let (tx_phone, mut phone) = tokio::sync::mpsc::unbounded_channel();
2231 let mac_id = reg.register("lv", "mac", "macos", "lv-mac", tx_mac, false, false);
2232 reg.register("lv", "phone", "ios", "lv-phone", tx_phone, false, false);
2233 let levels = |rx: &mut tokio::sync::mpsc::UnboundedReceiver<Envelope>| {
2234 std::iter::from_fn(|| rx.try_recv().ok().map(|e| e.command))
2235 .filter(|c| {
2236 matches!(
2237 c,
2238 LinkCommand::Levels { .. } | LinkCommand::WatchLevels { .. }
2239 )
2240 })
2241 .collect::<Vec<_>>()
2242 };
2243 let f = Frame(1_000, 1, 2, 3);
2244
2245 reg.levels("lv", "lv-phone", f);
2246 assert!(levels(&mut mac).is_empty(), "nobody watching");
2247
2248 reg.watch_levels("lv", "lv-mac", "lv-phone", true);
2249 assert_eq!(
2250 levels(&mut phone),
2251 vec![LinkCommand::WatchLevels { on: true }]
2252 );
2253 reg.levels("lv", "lv-phone", f);
2254 assert_eq!(
2255 levels(&mut mac),
2256 vec![LinkCommand::Levels {
2257 from: "lv-phone".into(),
2258 f
2259 }]
2260 );
2261
2262 reg.unregister(&mac_id);
2264 assert_eq!(
2265 levels(&mut phone),
2266 vec![LinkCommand::WatchLevels { on: false }]
2267 );
2268 reg.watch_levels("lv", "lv-mac", "lv-phone", true);
2269 reg.watch_levels("lv", "lv-mac", "lv-phone", false);
2270 assert_eq!(
2271 levels(&mut phone),
2272 vec![
2273 LinkCommand::WatchLevels { on: true },
2274 LinkCommand::WatchLevels { on: false }
2275 ]
2276 );
2277 }
2278
2279 #[test]
2280 fn a_playlist_order_adds_the_named_tracks_once_they_arrive() {
2281 let dir = tempfile::tempdir().unwrap();
2282 let path = dir.path().join("koan.db");
2283 let db = koan_core::db::connection::Database::open(&path).unwrap();
2284 let playlist = koan_core::db::queries::create_playlist(
2285 &db.conn,
2286 koan_core::db::queries::LOCAL_USER,
2287 "cyberpunk",
2288 None,
2289 )
2290 .unwrap();
2291 let order = Order {
2292 id: "o1".into(),
2293 username: None,
2294 client: None,
2295 artist: "Perturbator".into(),
2296 album: "Dangerous Days".into(),
2297 play_next: false,
2298 playlist: Some(playlist),
2299 titles: vec!["Future Club".into()],
2300 created_at: chrono::Utc::now().timestamp(),
2301 };
2302 registry().add_order(order);
2303
2304 fulfil_from(&path);
2306 assert!(registry().orders(None).iter().any(|o| o.id == "o1"));
2307
2308 db.conn
2309 .execute_batch(
2310 "INSERT INTO artists (id, name) VALUES (1, 'Perturbator');
2311 INSERT INTO albums (id, title, artist_id) VALUES (1, 'Dangerous Days', 1);
2312 INSERT INTO tracks (id, title, album_id, artist_id, track_number, path) VALUES
2313 (1, 'Welcome Back', 1, 1, 1, '/1.flac'), (2, 'Future Club', 1, 1, 2, '/2.flac');",
2314 )
2315 .unwrap();
2316 fulfil_from(&path);
2317 let held: Vec<i64> = db
2318 .conn
2319 .prepare("SELECT track_id FROM playlist_tracks WHERE playlist_id = ?1")
2320 .unwrap()
2321 .query_map([playlist], |r| r.get(0))
2322 .unwrap()
2323 .collect::<Result<_, _>>()
2324 .unwrap();
2325 assert_eq!(held, [2]);
2326 assert!(!registry().orders(None).iter().any(|o| o.id == "o1"));
2327 }
2328
2329 use super::*;
2330
2331 #[test]
2332 fn a_notification_cover_is_named_by_the_uid_the_command_carries() {
2333 let uid = "0190a5b2-7c3d-7e4f-8a1b-2c3d4e5f6a7b".to_string();
2335 let play = LinkCommand::Play {
2336 track_ids: vec!["x".into(), uid.clone()],
2337 start_at: 1,
2338 position_ms: 0,
2339 paused: false,
2340 handoff: false,
2341 };
2342 assert_eq!(cover_track(&play), Some(uid.as_str()));
2343 }
2344
2345 #[test]
2346 fn a_reconnect_replaces_the_device_and_commands_reach_it() {
2347 let reg = Registry::default();
2348 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
2349 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
2350 let (tx3, _rx3) = tokio::sync::mpsc::unbounded_channel();
2351 reg.register("j", "phone", "ios", "dev-1", tx1, false, false);
2352 let id = reg.register("j", "phone", "ios", "dev-1", tx2, false, false);
2353 reg.register("someone", "laptop", "macos", "dev-2", tx3, false, false);
2354
2355 assert_eq!(reg.list(Some("j")).len(), 1);
2356 assert_eq!(reg.list(None).len(), 2);
2357
2358 let sent = reg.send(Some("j"), None, LinkCommand::Pause).unwrap();
2359 assert_eq!(sent.id, id);
2360 assert_eq!(rx2.try_recv().unwrap().command, LinkCommand::Pause);
2361
2362 assert!(
2364 reg.send(Some("j"), Some("laptop"), LinkCommand::Pause)
2365 .is_err()
2366 );
2367
2368 reg.unregister(&id);
2369 assert!(reg.send(Some("j"), None, LinkCommand::Pause).is_err());
2370 }
2371
2372 #[test]
2373 fn disconnecting_an_account_closes_only_its_links() {
2374 let reg = Registry::default();
2375 let (tx1, mut rx1) = tokio::sync::mpsc::unbounded_channel();
2376 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
2377 reg.register("j", "phone", "ios", "dev-1", tx1, false, false);
2378 reg.register("someone", "laptop", "macos", "dev-2", tx2, false, false);
2379
2380 reg.disconnect("j");
2381 assert!(reg.list(Some("j")).is_empty());
2382 assert!(matches!(
2384 rx1.try_recv(),
2385 Err(tokio::sync::mpsc::error::TryRecvError::Disconnected)
2386 ));
2387 assert_eq!(reg.list(Some("someone")).len(), 1);
2388 assert!(matches!(
2389 rx2.try_recv(),
2390 Err(tokio::sync::mpsc::error::TryRecvError::Empty)
2391 ));
2392 }
2393
2394 #[test]
2395 fn a_burst_of_changes_wakes_each_device_once_when_it_goes_quiet() {
2396 let t0 = std::time::Instant::now();
2397 let s = std::time::Duration::from_secs;
2398 let phone = || ("dev-1".to_string(), "j".to_string());
2399 let ipad = || ("dev-2".to_string(), "j".to_string());
2400 let mut wakes = Wakes {
2401 pending: Vec::new(),
2402 timer: false,
2403 };
2404
2405 wakes.add([phone()], t0);
2406 wakes.add([phone(), ipad()], t0 + s(10));
2407 wakes.add([phone()], t0 + s(20));
2408 assert_eq!(wakes.pending.len(), 2);
2409
2410 assert!(wakes.take_due(t0 + s(39)).is_empty());
2412 assert_eq!(wakes.take_due(t0 + s(40)), [ipad()]);
2413 assert_eq!(wakes.next(), Some(t0 + s(50)));
2414 assert_eq!(wakes.take_due(t0 + s(50)), [phone()]);
2415 assert_eq!(wakes.next(), None);
2416 }
2417
2418 #[test]
2419 fn a_library_that_never_goes_quiet_still_wakes_devices() {
2420 let t0 = std::time::Instant::now();
2421 let phone = || ("dev-1".to_string(), "j".to_string());
2422 let mut wakes = Wakes {
2423 pending: Vec::new(),
2424 timer: false,
2425 };
2426 let mut sent = 0;
2427 for i in 0..40 {
2428 let now = t0 + std::time::Duration::from_secs(i * 10);
2429 sent += wakes.take_due(now).len();
2430 wakes.add([phone()], now);
2431 }
2432 assert_eq!(sent, 1);
2435 assert_eq!(wakes.pending.len(), 1);
2436 }
2437
2438 #[test]
2439 fn the_device_playing_is_the_one_meant() {
2440 let reg = Registry::default();
2441 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
2442 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
2443 let mac = reg.register("j", "mac", "macos", "dev-1", tx1, false, false);
2444 let phone = reg.register("j", "phone", "ios", "dev-2", tx2, false, false);
2445
2446 let err = reg.send(Some("j"), None, LinkCommand::Pause).unwrap_err();
2448 assert!(err.contains("mac") && err.contains("phone"), "{err}");
2449
2450 reg.report(
2451 &phone,
2452 LinkState {
2453 playing: true,
2454 ..Default::default()
2455 },
2456 );
2457 assert_eq!(
2458 reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
2459 phone
2460 );
2461 assert_eq!(rx2.try_recv().unwrap().command, LinkCommand::Pause);
2462
2463 reg.report(&phone, LinkState::default());
2465 assert_eq!(
2466 reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
2467 phone
2468 );
2469 let _ = mac;
2470 }
2471
2472 fn drain(rx: &mut tokio::sync::mpsc::UnboundedReceiver<Envelope>) -> Vec<LinkCommand> {
2473 std::iter::from_fn(|| rx.try_recv().ok().map(|e| e.command)).collect()
2474 }
2475
2476 fn shared() -> (
2478 Registry,
2479 tokio::sync::mpsc::UnboundedReceiver<Envelope>,
2480 tokio::sync::mpsc::UnboundedReceiver<Envelope>,
2481 tokio::sync::mpsc::UnboundedReceiver<Envelope>,
2482 ) {
2483 let reg = Registry::default();
2484 let (phone_tx, mut phone) = tokio::sync::mpsc::unbounded_channel();
2485 let (mac_tx, mac) = tokio::sync::mpsc::unbounded_channel();
2486 let (k_tx, k) = tokio::sync::mpsc::unbounded_channel();
2487 let id = reg.register("j", "phone", "ios", "dev-phone", phone_tx, true, false);
2488 reg.register("j", "mac", "macos", "dev-mac", mac_tx, true, false);
2489 reg.register("k", "laptop", "macos", "dev-k", k_tx, true, false);
2490 reg.report(
2491 &id,
2492 LinkState {
2493 playing: true,
2494 title: Some("Roygbiv".into()),
2495 outputs: Some(Default::default()),
2496 ..Default::default()
2497 },
2498 );
2499 reg.share("j", "dev-phone", "k", true).unwrap();
2500 drain(&mut phone);
2501 (reg, phone, mac, k)
2502 }
2503
2504 #[test]
2505 fn a_grantee_sees_the_shared_device_and_nothing_else_of_the_owner() {
2506 let (_reg, _phone, _mac, mut k) = shared();
2507 let listed = drain(&mut k)
2508 .into_iter()
2509 .rev()
2510 .find_map(|c| match c {
2511 LinkCommand::Devices { devices } => Some(devices),
2512 _ => None,
2513 })
2514 .expect("told of it");
2515 assert_eq!(listed.len(), 1, "the phone, not the Mac");
2516 let phone = &listed[0];
2517 assert_eq!(phone.id, "dev-phone");
2518 assert_eq!(phone.owner.as_deref(), Some("j"));
2519 let state = phone.state.as_ref().expect("its state");
2520 assert_eq!(state.title.as_deref(), Some("Roygbiv"));
2521 assert!(
2522 state.outputs.is_some(),
2523 "outputs, to choose from: the output is in the playback set"
2524 );
2525 }
2526
2527 #[test]
2533 fn a_grantee_controls_playback_as_itself_and_nothing_of_the_owners_account() {
2534 let (reg, mut phone, mut mac, _k) = shared();
2535 let wrapped = |c: LinkCommand| LinkCommand::Shared {
2536 command: Box::new(c),
2537 };
2538 for cmd in [
2539 LinkCommand::Pause,
2540 LinkCommand::SetRendererVolume { volume: 40 },
2541 LinkCommand::SleepTimer {
2542 timer: Some(koan_core::player::state::SleepTimer::After { minutes: 30 }),
2543 },
2544 LinkCommand::HandOff { to: "dev-k".into() },
2545 ] {
2546 reg.relay("k", "dev-phone", cmd.clone()).unwrap();
2547 assert_eq!(drain(&mut phone), [wrapped(cmd)]);
2548 }
2549 for cmd in [
2550 LinkCommand::Sync { full: false },
2551 LinkCommand::Evict { track_ids: vec![] },
2552 LinkCommand::Shared {
2553 command: Box::new(LinkCommand::Pause),
2554 },
2555 ] {
2556 assert!(!cmd.allowed_playback());
2557 assert!(reg.relay("k", "dev-phone", cmd).is_err());
2558 }
2559 assert!(drain(&mut phone).is_empty());
2560 assert!(
2561 reg.relay("k", "dev-mac", LinkCommand::Pause).is_err(),
2562 "not shared, not reachable"
2563 );
2564 assert!(
2565 drain(&mut mac)
2566 .iter()
2567 .all(|c| matches!(c, LinkCommand::Devices { .. } | LinkCommand::Shares { .. })),
2568 "news, and no command"
2569 );
2570 assert!(reg.shared_owner("k", "dev-mac").is_none(), "nor wakeable");
2571 }
2572
2573 #[test]
2577 fn a_hand_off_runs_both_ways_across_a_grant_and_nothing_else_does() {
2578 let (reg, _phone, _mac, mut k) = shared();
2579 drain(&mut k);
2580 let play = LinkCommand::Play {
2581 track_ids: vec!["t".into()],
2582 start_at: 0,
2583 position_ms: 1000,
2584 paused: false,
2585 handoff: true,
2586 };
2587 reg.relay_from("j", Some("dev-phone"), "dev-k", play.clone())
2588 .unwrap();
2589 assert!(drain(&mut k).contains(&LinkCommand::Shared {
2590 command: Box::new(play.clone())
2591 }));
2592 assert!(
2593 reg.relay_from("j", Some("dev-phone"), "dev-k", LinkCommand::Pause)
2594 .is_err(),
2595 "the owner's phone does not command the grantee's devices"
2596 );
2597 assert!(
2598 reg.relay_from("j", Some("dev-mac"), "dev-k", play.clone())
2599 .is_err(),
2600 "only the shared device"
2601 );
2602 let not_a_hand_off = LinkCommand::Play {
2603 track_ids: vec!["t".into()],
2604 start_at: 0,
2605 position_ms: 0,
2606 paused: false,
2607 handoff: false,
2608 };
2609 assert!(
2610 reg.relay_from("j", Some("dev-phone"), "dev-k", not_a_hand_off)
2611 .is_err(),
2612 "a hand-off, not any play"
2613 );
2614 }
2615
2616 #[test]
2617 fn revoking_ends_control_at_once() {
2618 let (reg, mut phone, _mac, mut k) = shared();
2619 reg.share("j", "dev-phone", "k", false).unwrap();
2620 assert!(reg.relay("k", "dev-phone", LinkCommand::Pause).is_err());
2621 assert!(
2622 drain(&mut phone)
2623 .iter()
2624 .all(|c| !matches!(c, LinkCommand::Pause))
2625 );
2626 assert!(reg.shared_owner("k", "dev-phone").is_none());
2627 let listed = drain(&mut k)
2628 .into_iter()
2629 .rev()
2630 .find_map(|c| match c {
2631 LinkCommand::Devices { devices } => Some(devices),
2632 _ => None,
2633 })
2634 .expect("told it is gone");
2635 assert!(listed.is_empty());
2636 }
2637
2638 #[test]
2639 fn a_device_hears_whom_it_is_shared_with_every_time_it_links() {
2640 let reg = Registry::default();
2641 let (tx, mut phone) = tokio::sync::mpsc::unbounded_channel();
2642 reg.register("j", "phone", "ios", "dev-phone", tx, true, false);
2643 let first = drain(&mut phone);
2644 assert!(
2645 first.contains(&LinkCommand::Shares {
2646 grantees: vec![],
2647 error: None,
2648 accounts: vec![],
2649 }),
2650 "an empty list replaces one from another server"
2651 );
2652 reg.send_shares(
2653 "j",
2654 "dev-phone",
2655 Some("There is no account called x".into()),
2656 );
2657 assert!(drain(&mut phone).iter().any(
2658 |c| matches!(c, LinkCommand::Shares { error: Some(e), .. } if e.contains("no account"))
2659 ));
2660 }
2661
2662 #[test]
2663 fn the_owner_is_told_who_it_shares_with() {
2664 let reg = Registry::default();
2665 let (tx, mut phone) = tokio::sync::mpsc::unbounded_channel();
2666 reg.register("j", "phone", "ios", "dev-phone", tx, true, false);
2667 reg.share("j", "dev-phone", "k", true).unwrap();
2668 reg.share("j", "dev-phone", "m", true).unwrap();
2669 let last = drain(&mut phone).into_iter().rev().find_map(|c| match c {
2670 LinkCommand::Shares { grantees, .. } => Some(grantees),
2671 _ => None,
2672 });
2673 assert_eq!(last, Some(vec!["k".to_string(), "m".to_string()]));
2674 assert!(
2675 reg.share("j", "dev-phone", "j", true).is_err(),
2676 "not with itself"
2677 );
2678 }
2679
2680 #[test]
2684 fn a_device_is_woken_for_another_account_on_its_network_or_by_grant() {
2685 let reg = Registry::default();
2686 let home: std::net::IpAddr = "203.0.113.7".parse().unwrap();
2687 let away: std::net::IpAddr = "198.51.100.2".parse().unwrap();
2688 reg.seen_at("dev-ipad", "sarita", home);
2689 reg.seen_at("dev-phone", "admin", home);
2690 assert_eq!(
2691 reg.wake_owner("admin", "dev-phone", "dev-ipad"),
2692 "sarita",
2693 "same address: the iPad's own token"
2694 );
2695 reg.seen_at("dev-phone", "admin", away);
2696 assert_eq!(
2697 reg.wake_owner("admin", "dev-phone", "dev-ipad"),
2698 "admin",
2699 "elsewhere, with no grant: only admin's own, which it is not"
2700 );
2701 reg.share("sarita", "dev-ipad", "admin", true).unwrap();
2702 assert_eq!(
2703 reg.wake_owner("admin", "dev-phone", "dev-ipad"),
2704 "sarita",
2705 "shared: from anywhere"
2706 );
2707 reg.seen_at("dev-other", "mallory", home);
2709 assert_eq!(reg.wake_owner("admin", "dev-other", "dev-tv"), "admin");
2710 }
2711
2712 #[test]
2715 fn where_a_device_last_linked_from_survives_a_restart() {
2716 let conn = rusqlite::Connection::open_in_memory().unwrap();
2717 koan_core::db::schema::create_tables(&conn).unwrap();
2718 for (device, user, seen) in [("dev-ipad", "sarita", 1), ("dev-phone", "admin", 2)] {
2719 conn.execute(
2720 "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES (?1, ?2, ?1, 'ios', ?3)",
2721 rusqlite::params![device, user, seen],
2722 )
2723 .unwrap();
2724 }
2725 let home: std::net::IpAddr = "203.0.113.7".parse().unwrap();
2726 outbox::save_address_in(&conn, "dev-ipad", "sarita", home);
2727 outbox::save_address_in(&conn, "dev-phone", "admin", home);
2728
2729 let reg = Registry::default();
2731 *reg.addresses.lock() = outbox::load_addresses_in(&conn);
2732 assert_eq!(
2733 reg.wake_owner("admin", "dev-phone", "dev-ipad"),
2734 "sarita",
2735 "still woken through its own account after the restart"
2736 );
2737 }
2738
2739 #[test]
2743 fn a_forgotten_device_is_dropped_by_the_accounts_devices() {
2744 let reg = Registry::default();
2745 let (mac_tx, mut mac) = tokio::sync::mpsc::unbounded_channel();
2746 let (other_tx, mut other) = tokio::sync::mpsc::unbounded_channel();
2747 reg.register("j", "mac", "macos", "dev-mac", mac_tx, true, false);
2748 reg.register("k", "laptop", "macos", "dev-k", other_tx, true, false);
2749 while mac.try_recv().is_ok() {}
2750 while other.try_recv().is_ok() {}
2751
2752 reg.forget("j", "dev-phone").unwrap();
2753 let heard: Vec<LinkCommand> =
2754 std::iter::from_fn(|| mac.try_recv().ok().map(|e| e.command)).collect();
2755 assert!(heard.contains(&LinkCommand::Forgotten {
2756 device: "dev-phone".into()
2757 }));
2758 assert!(
2759 std::iter::from_fn(|| other.try_recv().ok().map(|e| e.command))
2760 .all(|c| !matches!(c, LinkCommand::Forgotten { .. })),
2761 "another account hears nothing of it"
2762 );
2763 assert!(
2764 reg.forget("j", "dev-mac").is_err(),
2765 "linked: it would be back"
2766 );
2767 assert!(
2768 reg.relay(
2769 "j",
2770 "dev-mac",
2771 LinkCommand::Forgotten { device: "x".into() }
2772 )
2773 .is_err(),
2774 "news from the server, not a command a device may send"
2775 );
2776 }
2777
2778 #[test]
2781 fn forgetting_a_shared_device_declines_the_share() {
2782 let (reg, mut phone, _mac, mut k) = shared();
2783 drain(&mut k);
2784 reg.forget("k", "dev-phone").unwrap();
2785 assert!(reg.shared_owner("k", "dev-phone").is_none());
2786 assert!(
2787 drain(&mut phone)
2788 .iter()
2789 .any(|c| matches!(c, LinkCommand::Shares { grantees, .. } if grantees.is_empty()))
2790 );
2791 let listed = drain(&mut k).into_iter().rev().find_map(|c| match c {
2792 LinkCommand::Devices { devices } => Some(devices),
2793 _ => None,
2794 });
2795 assert_eq!(listed.map(|d| d.len()), Some(0));
2796 }
2797
2798 #[test]
2801 fn forgetting_a_device_stops_pushes_to_it_and_only_for_its_account() {
2802 let conn = rusqlite::Connection::open_in_memory().unwrap();
2803 koan_core::db::schema::create_tables(&conn).unwrap();
2804 for user in ["j", "k"] {
2805 conn.execute(
2806 "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES ('dev-phone', ?1, 'phone', 'ios', 1)",
2807 [user],
2808 )
2809 .unwrap();
2810 conn.execute(
2811 "INSERT INTO link_push (device, username, token, sandbox, updated_at) VALUES ('dev-phone', ?1, 't', 0, 1)",
2812 [user],
2813 )
2814 .unwrap();
2815 }
2816 outbox::forget_device_in(&conn, "dev-phone", "k");
2817 assert_eq!(
2818 outbox::push_targets_in(&conn, Some("j")).len(),
2819 1,
2820 "k forgetting its own leaves j's alone"
2821 );
2822 outbox::forget_device_in(&conn, "dev-phone", "j");
2823 assert!(outbox::push_targets_in(&conn, Some("j")).is_empty());
2824 assert!(outbox::push_targets_in(&conn, None).is_empty());
2825 }
2826
2827 #[test]
2828 fn a_device_hears_its_peers_and_commands_reach_them_by_device_id() {
2829 let reg = Registry::default();
2830 let (tx1, mut rx1) = tokio::sync::mpsc::unbounded_channel();
2831 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
2832 let (tx3, mut rx3) = tokio::sync::mpsc::unbounded_channel();
2833 reg.register("j", "mac", "macos", "dev-mac", tx1, true, false);
2834 let phone = reg.register("j", "phone", "ios", "dev-phone", tx2, true, false);
2835 reg.register("someone", "laptop", "macos", "dev-other", tx3, true, false);
2836
2837 reg.report(
2838 &phone,
2839 LinkState {
2840 playing: true,
2841 title: Some("Roygbiv".into()),
2842 ..Default::default()
2843 },
2844 );
2845 let devices = drain(&mut rx1)
2846 .into_iter()
2847 .rev()
2848 .find_map(|c| match c {
2849 LinkCommand::Devices { devices } => Some(devices),
2850 _ => None,
2851 })
2852 .expect("the Mac is told of the phone");
2853 assert_eq!(devices.len(), 1, "not itself, not another account's");
2854 assert_eq!(devices[0].id, "dev-phone");
2855 assert_eq!(
2856 devices[0].state.as_ref().and_then(|s| s.title.as_deref()),
2857 Some("Roygbiv")
2858 );
2859 while rx3.try_recv().is_ok() {}
2860 assert!(rx3.try_recv().is_err());
2861
2862 while rx2.try_recv().is_ok() {}
2863 reg.relay("j", "dev-phone", LinkCommand::Pause).unwrap();
2864 assert_eq!(rx2.try_recv().unwrap().command, LinkCommand::Pause);
2865 assert!(
2866 reg.relay("someone", "dev-phone", LinkCommand::Pause)
2867 .is_err(),
2868 "another account cannot reach it"
2869 );
2870 assert!(
2871 reg.relay("j", "dev-phone", LinkCommand::Devices { devices: vec![] })
2872 .is_err()
2873 );
2874 }
2875
2876 #[test]
2877 fn a_device_linked_for_a_month_is_not_forgotten() {
2878 let conn = rusqlite::Connection::open_in_memory().unwrap();
2879 koan_core::db::schema::create_tables(&conn).unwrap();
2880 let now = 100 * 24 * 60 * 60;
2881 let long_ago = now - 40 * 24 * 60 * 60;
2882 for device in ["mac", "old-phone"] {
2883 conn.execute(
2884 "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES (?1, 'j', ?1, 'ios', ?2)",
2885 rusqlite::params![device, long_ago],
2886 )
2887 .unwrap();
2888 conn.execute(
2889 "INSERT INTO link_push (device, username, token, sandbox, updated_at) VALUES (?1, 'j', 't', 0, ?2)",
2890 rusqlite::params![device, long_ago],
2891 )
2892 .unwrap();
2893 }
2894 outbox::forget_stale(&conn, &[("mac".into(), "j".into())], now);
2895 let left = |table: &str| -> Vec<String> {
2896 conn.prepare(&format!("SELECT device FROM {table} ORDER BY device"))
2897 .unwrap()
2898 .query_map([], |r| r.get(0))
2899 .unwrap()
2900 .collect::<Result<_, _>>()
2901 .unwrap()
2902 };
2903 assert_eq!(left("link_devices"), ["mac"]);
2904 assert_eq!(left("link_push"), ["mac"]);
2905 }
2906}