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