1use std::sync::LazyLock;
10
11use koan_core::db::queries::{self, UidKind};
12use koan_core::remote::link::{LinkCommand, LinkDevice, LinkState};
13use outbox::Absent;
14use parking_lot::Mutex;
15use tokio::sync::mpsc::UnboundedSender;
16
17#[derive(Debug, Clone)]
19pub struct ClientInfo {
20 pub id: String,
21 pub device: String,
23 pub name: String,
24 pub platform: String,
25 pub username: String,
26 pub connected_at: i64,
28 pub state: LinkState,
30 pub last_played_at: Option<i64>,
32 pub state_at: i64,
34 pub reports: bool,
37 pub notified: bool,
40}
41
42impl ClientInfo {
43 pub fn position_ms(&self) -> u64 {
45 let pos = self.state.position_ms;
46 if !self.state.playing {
47 return pos;
48 }
49 let run = (chrono::Utc::now().timestamp_millis() - self.state_at).max(0) as u64;
50 let pos = pos + run;
51 if self.state.duration_ms > 0 {
52 pos.min(self.state.duration_ms)
53 } else {
54 pos
55 }
56 }
57}
58
59struct Entry {
60 info: ClientInfo,
61 device: String,
63 tx: UnboundedSender<LinkCommand>,
64 wants_devices: bool,
67}
68
69struct Activity {
72 username: String,
73 watcher: String,
75 target: String,
77 token: String,
78 sandbox: bool,
79 sent: Option<crate::push::ActivityState>,
81}
82
83#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
86pub struct Order {
87 pub id: String,
88 pub username: Option<String>,
90 pub client: Option<String>,
92 pub artist: String,
93 pub album: String,
94 pub play_next: bool,
96 #[serde(default)]
98 pub playlist: Option<i64>,
99 #[serde(default)]
102 pub titles: Vec<String>,
103 pub created_at: i64,
105}
106
107const ORDER_TTL: i64 = 24 * 60 * 60;
109
110#[derive(Debug, Clone, PartialEq, Eq)]
114pub struct Grant {
115 pub device: String,
116 pub owner: String,
117 pub grantee: String,
118}
119
120#[derive(Default)]
121pub struct Registry {
122 entries: Mutex<Vec<Entry>>,
123 orders: Mutex<Vec<Order>>,
124 activities: Mutex<Vec<Activity>>,
125 grants: Mutex<Vec<Grant>>,
126 level_watches: Mutex<Vec<LevelWatch>>,
129 addresses: Mutex<std::collections::HashMap<String, (String, std::net::IpAddr)>>,
132}
133
134struct LevelWatch {
138 owner: String,
139 watcher_user: String,
140 watcher: String,
141 target: String,
142}
143
144pub fn registry() -> &'static Registry {
147 static REGISTRY: LazyLock<Registry> = LazyLock::new(|| {
148 let registry = Registry::default();
149 *registry.orders.lock() = outbox::load_orders();
150 *registry.grants.lock() = outbox::load_grants();
151 *registry.addresses.lock() = outbox::load_addresses();
152 registry
153 });
154 ®ISTRY
155}
156
157impl Registry {
158 pub fn register(
161 &self,
162 username: &str,
163 name: &str,
164 platform: &str,
165 device: &str,
166 tx: UnboundedSender<LinkCommand>,
167 wants_devices: bool,
168 ) -> String {
169 let id = uuid::Uuid::now_v7().to_string();
170 if let Some((sent, how)) = WOKEN
171 .lock()
172 .remove(&(device.to_string(), username.to_string()))
173 {
174 log::info!(
175 "wake: {name} linked {}ms after its {how}",
176 sent.elapsed().as_millis()
177 );
178 }
179 let live = self.live();
182 for cmd in outbox::take_and_remember(username, device, name, platform, &live) {
183 let _ = tx.send(cmd);
184 }
185 let mut entries = self.entries.lock();
186 entries.retain(|e| !(e.device == device && e.info.username == username));
187 entries.push(Entry {
188 info: ClientInfo {
189 id: id.clone(),
190 device: device.to_string(),
191 name: name.to_string(),
192 platform: platform.to_string(),
193 username: username.to_string(),
194 connected_at: chrono::Utc::now().timestamp(),
195 state: LinkState::default(),
196 last_played_at: None,
197 state_at: chrono::Utc::now().timestamp_millis(),
198 reports: false,
199 notified: false,
200 },
201 device: device.to_string(),
202 tx,
203 wants_devices,
204 });
205 drop(entries);
206 self.send_shares(username, device, None);
209 self.announce(username);
210 if self
213 .level_watches
214 .lock()
215 .iter()
216 .any(|w| w.owner == username && w.target == device)
217 {
218 self.send_live(username, device, LinkCommand::WatchLevels { on: true });
219 }
220 id
221 }
222
223 pub fn report(&self, id: &str, state: LinkState) {
225 let mut entries = self.entries.lock();
226 if let Some(e) = entries.iter_mut().find(|e| e.info.id == id) {
227 if state.playing || e.info.state.playing {
228 e.info.last_played_at = Some(chrono::Utc::now().timestamp());
229 }
230 e.info.state = state;
231 e.info.state_at = chrono::Utc::now().timestamp_millis();
232 e.info.reports = true;
233 let (username, device) = (e.info.username.clone(), e.device.clone());
234 drop(entries);
235 self.announce(&username);
236 self.update_activities(&username, &device);
237 }
238 }
239
240 pub fn set_push(&self, username: &str, device: &str, token: &str, sandbox: bool) {
242 outbox::save_push(username, device, token, sandbox);
243 }
244
245 pub fn disconnect(&self, username: &str) {
248 self.entries.lock().retain(|e| e.info.username != username);
249 }
250
251 pub fn unregister(&self, id: &str) {
252 let mut entries = self.entries.lock();
253 let gone = entries
254 .iter()
255 .find(|e| e.info.id == id)
256 .map(|e| (e.device.clone(), e.info.username.clone()));
257 entries.retain(|e| e.info.id != id);
258 drop(entries);
259 if let Some((device, username)) = gone {
260 outbox::touch(&[(device.clone(), username.clone())]);
261 self.announce(&username);
262 let targets: Vec<String> = self
264 .level_watches
265 .lock()
266 .iter()
267 .filter(|w| w.watcher_user == username && w.watcher == device)
268 .map(|w| w.target.clone())
269 .collect();
270 for target in targets {
271 self.watch_levels(&username, &device, &target, false);
272 }
273 }
274 }
275
276 pub fn watch_levels(&self, username: &str, watcher: &str, target: &str, on: bool) {
281 let owner = self
283 .shared_owner(username, target)
284 .unwrap_or_else(|| username.to_string());
285 let mut watches = self.level_watches.lock();
286 watches.retain(|w| {
287 !(w.watcher_user == username && w.watcher == watcher && w.target == target)
288 });
289 if on {
290 watches.push(LevelWatch {
291 owner: owner.clone(),
292 watcher_user: username.into(),
293 watcher: watcher.into(),
294 target: target.into(),
295 });
296 }
297 let watched = watches
298 .iter()
299 .any(|w| w.owner == owner && w.target == target);
300 drop(watches);
301 if on || !watched {
302 self.send_live(&owner, target, LinkCommand::WatchLevels { on: watched });
303 }
304 }
305
306 pub fn levels(&self, username: &str, from: &str, f: koan_core::remote::levels::Frame) {
309 let watchers: Vec<(String, String)> = self
310 .level_watches
311 .lock()
312 .iter()
313 .filter(|w| w.owner == username && w.target == from)
314 .map(|w| (w.watcher_user.clone(), w.watcher.clone()))
315 .collect();
316 for (watcher_user, watcher) in watchers {
317 self.send_live(
318 &watcher_user,
319 &watcher,
320 LinkCommand::Levels {
321 from: from.to_string(),
322 f,
323 },
324 );
325 }
326 }
327
328 fn send_live(&self, username: &str, device: &str, cmd: LinkCommand) {
330 if let Some(e) = self
331 .entries
332 .lock()
333 .iter()
334 .find(|e| e.info.username == username && e.device == device)
335 {
336 let _ = e.tx.send(cmd);
337 }
338 }
339
340 fn announce(&self, username: &str) {
344 self.announce_to(username);
345 let grantees: Vec<String> = {
347 let mut g: Vec<String> = self
348 .grants
349 .lock()
350 .iter()
351 .filter(|g| g.owner == username)
352 .map(|g| g.grantee.clone())
353 .collect();
354 g.sort();
355 g.dedup();
356 g
357 };
358 for grantee in grantees {
359 self.announce_to(&grantee);
360 }
361 }
362
363 fn announce_to(&self, username: &str) {
365 let listening = |e: &Entry| e.info.username == username && e.wants_devices;
366 if !self.entries.lock().iter().any(listening) {
367 return;
368 }
369 let asleep = outbox::push_targets(Some(username));
370 let can_push = crate::push::pusher().is_some();
372 let entries = self.entries.lock();
373 let ours: Vec<&Entry> = entries
374 .iter()
375 .filter(|e| e.info.username == username)
376 .collect();
377 if !ours.iter().any(|e| e.wants_devices) {
378 return;
379 }
380 let mut all: Vec<LinkDevice> = ours
381 .iter()
382 .map(|e| LinkDevice {
383 id: e.device.clone(),
384 name: e.info.name.clone(),
385 platform: e.info.platform.clone(),
386 linked: true,
387 state: e.info.reports.then(|| LinkState {
388 position_ms: e.info.position_ms(),
389 ..e.info.state.clone()
390 }),
391 last_seen: None,
392 wakeable: Some(can_push && asleep.iter().any(|t| t.device == e.device)),
393 owner: None,
394 })
395 .collect();
396 for t in asleep {
397 if !all.iter().any(|d| d.id == t.device) {
398 all.push(LinkDevice {
399 id: t.device,
400 name: t.name,
401 platform: t.platform,
402 linked: false,
403 state: None,
404 last_seen: Some(t.last_seen),
405 wakeable: Some(can_push),
406 owner: None,
407 });
408 }
409 }
410 all.extend(self.shared_devices(username, &entries, can_push));
411 for e in ours.iter().filter(|e| e.wants_devices) {
412 let devices = all.iter().filter(|d| d.id != e.device).cloned().collect();
413 let _ = e.tx.send(LinkCommand::Devices { devices });
414 }
415 }
416
417 fn shared_devices(&self, grantee: &str, entries: &[Entry], can_push: bool) -> Vec<LinkDevice> {
421 let grants: Vec<Grant> = self
422 .grants
423 .lock()
424 .iter()
425 .filter(|g| g.grantee == grantee)
426 .cloned()
427 .collect();
428 let mut out = Vec::new();
429 for g in grants {
430 let linked = entries
431 .iter()
432 .find(|e| e.device == g.device && e.info.username == g.owner);
433 match linked {
434 Some(e) => out.push(LinkDevice {
435 id: e.device.clone(),
436 name: e.info.name.clone(),
437 platform: e.info.platform.clone(),
438 linked: true,
439 state: e.info.reports.then(|| LinkState {
440 position_ms: e.info.position_ms(),
441 ..e.info.state.clone()
442 }),
443 last_seen: None,
444 wakeable: Some(can_push),
445 owner: Some(g.owner.clone()),
446 }),
447 None => {
448 if let Some(t) = outbox::push_targets(Some(&g.owner))
449 .into_iter()
450 .find(|t| t.device == g.device)
451 {
452 out.push(LinkDevice {
453 id: t.device,
454 name: t.name,
455 platform: t.platform,
456 linked: false,
457 state: None,
458 last_seen: Some(t.last_seen),
459 wakeable: Some(can_push),
460 owner: Some(g.owner.clone()),
461 });
462 }
463 }
464 }
465 }
466 out
467 }
468
469 pub fn share(
473 &self,
474 owner: &str,
475 device: &str,
476 grantee: &str,
477 allow: bool,
478 ) -> Result<(), String> {
479 if grantee == owner {
480 return Err("a device is already its own account's".into());
481 }
482 let grant = Grant {
483 device: device.to_string(),
484 owner: owner.to_string(),
485 grantee: grantee.to_string(),
486 };
487 {
488 let mut grants = self.grants.lock();
489 grants.retain(|g| *g != grant);
490 if allow {
491 grants.push(grant.clone());
492 }
493 }
494 if allow {
495 outbox::save_grant(&grant);
496 } else {
497 outbox::drop_grant(&grant);
498 }
499 log::info!(
500 "share: {owner}'s {device} {} {grantee}",
501 if allow {
502 "shared with"
503 } else {
504 "no longer shared with"
505 }
506 );
507 self.send_shares(owner, device, None);
508 self.announce_to(grantee);
509 Ok(())
510 }
511
512 pub fn send_shares(&self, owner: &str, device: &str, error: Option<String>) {
515 let mut grantees: Vec<String> = self
516 .grants
517 .lock()
518 .iter()
519 .filter(|g| g.owner == owner && g.device == device)
520 .map(|g| g.grantee.clone())
521 .collect();
522 grantees.sort();
523 if let Some(e) = self
524 .entries
525 .lock()
526 .iter()
527 .find(|e| e.info.username == owner && e.device == device && e.wants_devices)
528 {
529 let accounts = outbox::accounts()
530 .into_iter()
531 .filter(|a| a != owner)
532 .collect();
533 let _ = e.tx.send(LinkCommand::Shares {
534 grantees,
535 error,
536 accounts,
537 });
538 }
539 }
540
541 pub fn seen_at(&self, device: &str, username: &str, addr: std::net::IpAddr) {
543 self.addresses
544 .lock()
545 .insert(device.to_string(), (username.to_string(), addr));
546 outbox::save_address(device, username, addr);
547 }
548
549 fn wake_owner(&self, username: &str, from: &str, to: &str) -> String {
556 if let Some(owner) = self.shared_owner(username, to) {
557 return owner;
558 }
559 let addresses = self.addresses.lock();
560 let here = addresses
561 .get(from)
562 .filter(|(user, _)| user == username)
563 .map(|(_, addr)| *addr);
564 match (here, addresses.get(to)) {
565 (Some(here), Some((owner, there))) if *there == here => owner.clone(),
566 _ => username.to_string(),
567 }
568 }
569
570 fn grantee_of(&self, owner: &str, from: &str, to: &str) -> Option<String> {
573 let to_user = self
574 .entries
575 .lock()
576 .iter()
577 .find(|e| e.device == to && e.info.username != owner)
578 .map(|e| e.info.username.clone())?;
579 self.grants
580 .lock()
581 .iter()
582 .any(|g| g.owner == owner && g.device == from && g.grantee == to_user)
583 .then_some(to_user)
584 }
585
586 fn shared_owner(&self, username: &str, to: &str) -> Option<String> {
589 self.grants
590 .lock()
591 .iter()
592 .find(|g| g.grantee == username && g.device == to)
593 .map(|g| g.owner.clone())
594 }
595
596 pub fn wake(&self, username: &str, from: &str, to: &str, notify: bool) {
602 let from_name = self
603 .list(Some(username))
604 .into_iter()
605 .find(|c| c.device == from)
606 .map_or_else(|| "Another device".to_string(), |c| c.name);
607 let owner = self.wake_owner(username, from, to);
608 if owner != username {
609 log::info!("wake: {to} is {owner}'s, woken for {username}");
610 }
611 let username = &owner;
612 if self.list(Some(username)).iter().any(|c| c.device == to) {
613 log::info!("wake: {to} is already linked");
614 return;
615 }
616 let Some(pusher) = crate::push::pusher() else {
617 log::info!("wake: no push key, so {to} cannot be woken");
618 return;
619 };
620 let Some(target) = outbox::push_targets(Some(username))
621 .into_iter()
622 .find(|t| t.device == to)
623 else {
624 log::info!("wake: {to} has given no push token");
625 return;
626 };
627 let (push, how) = if notify {
628 (
629 crate::push::Push::Summon {
630 title: format!("{from_name} wants to play here"),
631 body: "Tap to open kōan and let it.".into(),
632 },
633 "notification",
634 )
635 } else {
636 (crate::push::Push::WakeNow, "wake push")
637 };
638 WOKEN.lock().insert(
639 (target.device.clone(), target.username.clone()),
640 (std::time::Instant::now(), how),
641 );
642 std::thread::spawn(move || {
643 let started = std::time::Instant::now();
644 deliver_push(pusher, &target, &push);
645 log::info!(
646 "wake: {how} for {} answered by APNs in {}ms",
647 target.name,
648 started.elapsed().as_millis()
649 );
650 });
651 }
652
653 pub fn forget(&self, username: &str, device: &str) -> Result<(), String> {
660 if let Some(owner) = self.shared_owner(username, device) {
663 return self.share(&owner, device, username, false);
664 }
665 if self.list(Some(username)).iter().any(|c| c.device == device) {
666 return Err(format!("{device} is linked; it would be back at once"));
667 }
668 outbox::forget_device(device, username);
669 log::info!("devices: {username} forgot {device}");
670 self.broadcast(
671 Some(username),
672 LinkCommand::Forgotten {
673 device: device.to_string(),
674 },
675 );
676 self.announce(username);
677 Ok(())
678 }
679
680 pub fn relay(
682 &self,
683 username: &str,
684 to: &str,
685 command: LinkCommand,
686 ) -> Result<ClientInfo, String> {
687 self.relay_from(username, None, to, command)
688 }
689
690 pub fn relay_from(
694 &self,
695 username: &str,
696 from: Option<&str>,
697 to: &str,
698 command: LinkCommand,
699 ) -> Result<ClientInfo, String> {
700 if matches!(
703 command,
704 LinkCommand::Devices { .. }
705 | LinkCommand::Forgotten { .. }
706 | LinkCommand::Levels { .. }
707 | LinkCommand::WatchLevels { .. }
708 | LinkCommand::Shares { .. }
709 | LinkCommand::Shared { .. }
710 ) {
711 return Err("not a command".into());
712 }
713 if let Some(owner) = self.shared_owner(username, to) {
717 if !command.allowed_playback() {
718 return Err(format!("{to} is shared for playback only"));
719 }
720 let command = LinkCommand::Shared {
721 command: Box::new(command),
722 };
723 return self.send(Some(&owner), Some(to), command);
724 }
725 if let Some(from) = from
729 && let Some(grantee) = self.grantee_of(username, from, to)
730 {
731 if !matches!(command, LinkCommand::Play { handoff: true, .. }) {
732 return Err(format!("{to} is another account's"));
733 }
734 let command = LinkCommand::Shared {
735 command: Box::new(command),
736 };
737 return self.send(Some(&grantee), Some(to), command);
738 }
739 self.send(Some(username), Some(to), command)
740 }
741
742 pub fn set_activity(
745 &self,
746 username: &str,
747 watcher: &str,
748 activity: Option<(String, String, bool)>,
749 ) {
750 let mut activities = self.activities.lock();
751 activities.retain(|a| !(a.username == username && a.watcher == watcher));
752 let Some((token, target, sandbox)) = activity else {
753 return;
754 };
755 activities.push(Activity {
756 username: username.to_string(),
757 watcher: watcher.to_string(),
758 target: target.clone(),
759 token,
760 sandbox,
761 sent: None,
762 });
763 drop(activities);
764 self.update_activities(username, &target);
765 }
766
767 fn update_activities(&self, username: &str, target: &str) {
770 let Some(pusher) = crate::push::pusher() else {
771 return;
772 };
773 let owner = self
776 .shared_owner(username, target)
777 .unwrap_or_else(|| username.to_string());
778 let Some(info) = self
779 .list(Some(&owner))
780 .into_iter()
781 .find(|c| c.device == target)
782 else {
783 return;
784 };
785 let watchers: Vec<String> = self
786 .grants
787 .lock()
788 .iter()
789 .filter(|g| g.owner == owner && g.device == target)
790 .map(|g| g.grantee.clone())
791 .chain(std::iter::once(owner.clone()))
792 .collect();
793 let state = crate::push::ActivityState::of(&info);
794 let mut due = Vec::new();
795 for a in self.activities.lock().iter_mut() {
796 if watchers.contains(&a.username)
797 && a.target == target
798 && a.sent.as_ref().is_none_or(|s| s.differs(&state))
799 {
800 a.sent = Some(state.clone());
801 due.push((
802 a.token.clone(),
803 a.sandbox,
804 a.watcher.clone(),
805 a.username.clone(),
806 ));
807 }
808 }
809 if due.is_empty() {
810 return;
811 }
812 std::thread::spawn(move || {
813 for (token, sandbox, watcher, username) in due {
814 let push = crate::push::Push::Activity(state.clone());
815 match pusher.send(&token, sandbox, &push) {
816 crate::push::Outcome::Sent => {}
817 crate::push::Outcome::Gone => {
818 log::info!("push: a Live Activity on {watcher} has ended");
819 registry().set_activity(&username, &watcher, None);
820 }
821 crate::push::Outcome::Failed(e) => {
822 log::warn!("push: Live Activity on {watcher}: {e}");
823 }
824 }
825 }
826 });
827 }
828
829 pub fn list(&self, username: Option<&str>) -> Vec<ClientInfo> {
831 let mut out: Vec<ClientInfo> = self
832 .entries
833 .lock()
834 .iter()
835 .filter(|e| username.is_none_or(|u| e.info.username == u))
836 .map(|e| e.info.clone())
837 .collect();
838 out.sort_by_key(|c| std::cmp::Reverse(c.connected_at));
839 out
840 }
841
842 pub fn send(
847 &self,
848 username: Option<&str>,
849 id: Option<&str>,
850 cmd: LinkCommand,
851 ) -> Result<ClientInfo, String> {
852 let clients = self.list(username);
853 let target = match id {
854 Some(id) => clients
855 .iter()
856 .find(|c| c.id == id || c.device == id || c.name.eq_ignore_ascii_case(id)),
857 None if clients.is_empty() => None,
858 None => Some(pick(&clients, chrono::Utc::now().timestamp())?),
859 };
860 let Some(target) = target else {
862 return reach_absent(username, id, &cmd).unwrap_or_else(|| {
863 Err(match id {
864 Some(id) => format!("no linked client {id}; see `clients`"),
865 None => "no koan app is linked to this server; open koan on the device".into(),
866 })
867 });
868 };
869 let entries = self.entries.lock();
870 let entry = entries
871 .iter()
872 .find(|e| e.info.id == target.id)
873 .ok_or("that client has just gone")?;
874 entry
875 .tx
876 .send(cmd)
877 .map_err(|_| "that client has just gone".to_string())?;
878 Ok(target.clone())
879 }
880}
881
882impl Registry {
883 pub fn add_order(&self, order: Order) {
884 outbox::save_order(&order);
885 self.orders.lock().push(order);
886 }
887
888 pub fn orders(&self, username: Option<&str>) -> Vec<Order> {
889 self.orders
890 .lock()
891 .iter()
892 .filter(|o| username.is_none() || o.username.as_deref() == username)
893 .cloned()
894 .collect()
895 }
896
897 pub fn cancel_order(&self, username: Option<&str>, id: &str) -> bool {
898 let mut orders = self.orders.lock();
899 let before = orders.len();
900 orders
901 .retain(|o| !(o.id == id && (username.is_none() || o.username.as_deref() == username)));
902 let gone = orders.len() != before;
903 if gone {
904 outbox::drop_order(id);
905 }
906 gone
907 }
908
909 fn done(&self, id: &str) {
910 self.orders.lock().retain(|o| o.id != id);
911 outbox::drop_order(id);
912 }
913
914 pub fn fulfil_orders(&self, find: impl Fn(&Order) -> Option<Vec<String>>) {
917 let now = chrono::Utc::now().timestamp();
918 let pending: Vec<Order> = {
919 let mut orders = self.orders.lock();
920 for o in orders.iter().filter(|o| now - o.created_at >= ORDER_TTL) {
921 outbox::drop_order(&o.id);
922 }
923 orders.retain(|o| now - o.created_at < ORDER_TTL);
924 orders.clone()
925 };
926 for order in pending.into_iter().filter(|o| o.playlist.is_none()) {
927 let Some(ids) = find(&order).filter(|ids| !ids.is_empty()) else {
928 continue;
929 };
930 let track_ids = ids;
931 let cmd = if order.play_next {
932 LinkCommand::PlayNext { track_ids }
933 } else {
934 LinkCommand::Enqueue { track_ids }
935 };
936 match self.send(order.username.as_deref(), order.client.as_deref(), cmd) {
937 Ok(c) => {
938 log::info!(
939 "link: {} — {} arrived; queued on {}",
940 order.artist,
941 order.album,
942 c.name
943 );
944 self.done(&order.id);
945 }
946 Err(e) => log::info!("link: {} — {} arrived but {e}", order.artist, order.album),
948 }
949 }
950 }
951}
952
953pub fn fulfil_from(db_path: &std::path::Path) {
955 let registry = registry();
956 if registry.orders.lock().is_empty() {
957 return;
958 }
959 let Ok(db) = koan_core::db::connection::Database::open_existing(db_path) else {
960 return;
961 };
962 let for_playlists: Vec<Order> = registry
965 .orders
966 .lock()
967 .iter()
968 .filter(|o| o.playlist.is_some())
969 .cloned()
970 .collect();
971 let mut edited = false;
972 for order in for_playlists {
973 let Some(playlist) = order.playlist else {
974 continue;
975 };
976 let Some(ids) = order_tracks(&db.conn, &order) else {
977 continue;
978 };
979 match koan_core::db::queries::add_tracks(&db.conn, playlist, &ids) {
980 Ok(_) => {
981 if !cfg!(test) {
983 koan_core::playlists::push_to_remote(playlist);
984 }
985 log::info!(
986 "link: {} — {} arrived; added {} tracks to playlist {playlist}",
987 order.artist,
988 order.album,
989 ids.len()
990 );
991 registry.done(&order.id);
992 edited = true;
993 }
994 Err(e) => log::warn!("link: could not add to playlist {playlist}: {e}"),
995 }
996 }
997 if edited {
998 changed();
999 }
1000 registry.fulfil_orders(|order| {
1001 let rows = order_tracks(&db.conn, order)?;
1002 queries::uids_in_order(&db.conn, UidKind::Track, &rows).ok()
1003 });
1004}
1005
1006fn order_tracks(conn: &rusqlite::Connection, order: &Order) -> Option<Vec<i64>> {
1009 let tracks = album_tracks(conn, &order.artist, &order.album)?;
1010 if order.titles.is_empty() {
1011 return Some(tracks.into_iter().map(|(id, _)| id).collect());
1012 }
1013 let picked: Vec<i64> = order
1014 .titles
1015 .iter()
1016 .filter_map(|want| {
1017 let want = want.to_lowercase();
1018 tracks
1019 .iter()
1020 .find(|(_, t)| t.to_lowercase().contains(&want))
1021 .map(|(id, _)| *id)
1022 })
1023 .collect();
1024 (!picked.is_empty()).then_some(picked)
1025}
1026
1027pub fn album_tracks(
1030 conn: &rusqlite::Connection,
1031 artist: &str,
1032 album: &str,
1033) -> Option<Vec<(i64, String)>> {
1034 let like = |s: &str| format!("%{}%", s.replace(['%', '_'], ""));
1035 let album_id: i64 = conn
1036 .query_row(
1037 "SELECT al.id FROM albums al JOIN artists a ON a.id = al.artist_id
1038 WHERE a.name LIKE ?1 COLLATE NOCASE AND al.title LIKE ?2 COLLATE NOCASE
1039 ORDER BY al.id DESC LIMIT 1",
1040 [like(artist), like(album)],
1041 |r| r.get(0),
1042 )
1043 .ok()?;
1044 let mut stmt = conn
1045 .prepare("SELECT id, title FROM tracks WHERE album_id = ?1 ORDER BY disc, track_number, id")
1046 .ok()?;
1047 let tracks = stmt
1048 .query_map([album_id], |r| Ok((r.get(0)?, r.get(1)?)))
1049 .ok()?
1050 .filter_map(Result::ok)
1051 .collect();
1052 Some(tracks)
1053}
1054
1055impl Registry {
1056 pub fn broadcast(&self, username: Option<&str>, cmd: LinkCommand) -> Vec<String> {
1059 let ids: Vec<String> = self.list(username).into_iter().map(|c| c.id).collect();
1060 let entries = self.entries.lock();
1061 entries
1062 .iter()
1063 .filter(|e| ids.contains(&e.info.id) && e.tx.send(cmd.clone()).is_ok())
1064 .map(|e| e.info.name.clone())
1065 .collect()
1066 }
1067}
1068
1069impl Registry {
1070 pub fn deliver(&self, username: Option<&str>, cmd: LinkCommand) -> (Vec<String>, Vec<String>) {
1075 let (sent, queued) = self.link_or_queue(username, &cmd);
1076 wake(&queued.iter().map(Absent::key).collect::<Vec<_>>());
1077 (sent, queued.into_iter().map(|q| q.name).collect())
1078 }
1079
1080 fn link_or_queue(
1081 &self,
1082 username: Option<&str>,
1083 cmd: &LinkCommand,
1084 ) -> (Vec<String>, Vec<Absent>) {
1085 let sent = self.broadcast(username, cmd.clone());
1086 let queued = outbox::queue_for_absent(username, &self.live(), cmd);
1087 (sent, queued)
1088 }
1089
1090 fn live(&self) -> Vec<(String, String)> {
1092 self.entries
1093 .lock()
1094 .iter()
1095 .map(|e| (e.device.clone(), e.info.username.clone()))
1096 .collect()
1097 }
1098}
1099
1100impl Absent {
1101 fn key(&self) -> (String, String) {
1102 (self.device.clone(), self.username.clone())
1103 }
1104}
1105
1106fn wake(devices: &[(String, String)]) {
1109 let Some(pusher) = crate::push::pusher() else {
1110 return;
1111 };
1112 let targets: Vec<outbox::PushTarget> = outbox::push_targets(None)
1113 .into_iter()
1114 .filter(|t| {
1115 devices
1116 .iter()
1117 .any(|(device, username)| *device == t.device && *username == t.username)
1118 })
1119 .collect();
1120 if targets.is_empty() {
1121 return;
1122 }
1123 std::thread::spawn(move || {
1124 for t in targets {
1125 deliver_push(pusher, &t, &crate::push::Push::Wake);
1126 }
1127 });
1128}
1129
1130fn reach_absent(
1139 username: Option<&str>,
1140 id: Option<&str>,
1141 cmd: &LinkCommand,
1142) -> Option<Result<ClientInfo, String>> {
1143 let pusher = crate::push::pusher()?;
1144 let target = outbox::push_targets(username)
1147 .into_iter()
1148 .find(|t| id.is_none_or(|id| t.device == id || t.name.eq_ignore_ascii_case(id)))?;
1149 let info = ClientInfo {
1150 id: target.device.clone(),
1151 device: target.device.clone(),
1152 name: target.name.clone(),
1153 platform: target.platform.clone(),
1154 username: target.username.clone(),
1155 connected_at: 0,
1156 state: LinkState::default(),
1157 last_played_at: None,
1158 state_at: 0,
1159 reports: false,
1160 notified: true,
1161 };
1162 let asked = match cmd {
1164 LinkCommand::Shared { command } => command.as_ref(),
1165 cmd => cmd,
1166 };
1167 let verb = match asked {
1168 LinkCommand::Play { .. } | LinkCommand::JumpTo { .. } | LinkCommand::Resume => Some("Play"),
1169 _ => None,
1170 };
1171 let push = match verb {
1172 Some(verb) => crate::push::Push::Notify {
1173 title: format!("{verb} on {}", target.name),
1174 body: outbox::describe(asked).unwrap_or_else(|| "From your koan server".into()),
1175 command: serde_json::to_value(cmd).ok()?,
1176 image: cover_track(asked)
1177 .and_then(outbox::track_row)
1178 .and_then(|t| pusher.cover_link(t)),
1179 },
1180 None => {
1181 outbox::queue_for(&target.device, &target.username, cmd);
1182 crate::push::Push::Wake
1183 }
1184 };
1185 std::thread::spawn(move || deliver_push(pusher, &target, &push));
1186 Some(Ok(info))
1187}
1188
1189fn cover_track(cmd: &LinkCommand) -> Option<&str> {
1192 match cmd {
1193 LinkCommand::Play {
1194 track_ids,
1195 start_at,
1196 ..
1197 } => track_ids
1198 .get(*start_at as usize)
1199 .or(track_ids.first())
1200 .map(String::as_str),
1201 LinkCommand::JumpTo { track_id } => Some(track_id),
1202 _ => None,
1203 }
1204}
1205
1206fn deliver_push(
1208 pusher: &crate::push::Pusher,
1209 target: &outbox::PushTarget,
1210 push: &crate::push::Push,
1211) {
1212 use crate::push::Outcome;
1213 match pusher.send(&target.token, target.sandbox, push) {
1214 Outcome::Sent => log::info!("push: sent to {}", target.name),
1215 Outcome::Gone => {
1216 log::info!(
1217 "push: {}'s token is no longer valid; forgotten",
1218 target.name
1219 );
1220 outbox::forget_push(&target.username, &target.device);
1221 }
1222 Outcome::Failed(e) => log::warn!("push: to {} failed: {e}", target.name),
1223 }
1224}
1225
1226type Woken = std::collections::HashMap<(String, String), (std::time::Instant, &'static str)>;
1230
1231static WOKEN: LazyLock<Mutex<Woken>> = LazyLock::new(Default::default);
1232
1233pub fn changed() {
1238 let (_, queued) = registry().link_or_queue(None, &LinkCommand::Sync { full: false });
1239 if queued.is_empty() || crate::push::pusher().is_none() {
1240 return;
1241 }
1242 let mut wakes = WAKES.lock();
1243 wakes.add(queued.iter().map(Absent::key), std::time::Instant::now());
1244 if !wakes.timer {
1245 wakes.timer = true;
1246 std::thread::spawn(send_wakes);
1247 }
1248}
1249
1250const QUIET: std::time::Duration = std::time::Duration::from_secs(30);
1254
1255const LONGEST_WAIT: std::time::Duration = std::time::Duration::from_secs(5 * 60);
1257
1258static WAKES: Mutex<Wakes> = Mutex::new(Wakes {
1259 pending: Vec::new(),
1260 timer: false,
1261});
1262
1263struct Wakes {
1266 pending: Vec<((String, String), std::time::Instant, std::time::Instant)>,
1268 timer: bool,
1270}
1271
1272impl Wakes {
1273 fn add(
1274 &mut self,
1275 devices: impl IntoIterator<Item = (String, String)>,
1276 now: std::time::Instant,
1277 ) {
1278 for key in devices {
1279 match self.pending.iter_mut().find(|(k, _, _)| *k == key) {
1280 Some((_, _, last)) => *last = now,
1281 None => self.pending.push((key, now, now)),
1282 }
1283 }
1284 }
1285
1286 fn due_at(first: std::time::Instant, last: std::time::Instant) -> std::time::Instant {
1287 (last + QUIET).min(first + LONGEST_WAIT)
1288 }
1289
1290 fn next(&self) -> Option<std::time::Instant> {
1292 self.pending
1293 .iter()
1294 .map(|(_, first, last)| Self::due_at(*first, *last))
1295 .min()
1296 }
1297
1298 fn take_due(&mut self, now: std::time::Instant) -> Vec<(String, String)> {
1300 let (due, waiting) = std::mem::take(&mut self.pending)
1301 .into_iter()
1302 .partition(|(_, first, last)| Self::due_at(*first, *last) <= now);
1303 self.pending = waiting;
1304 due.into_iter().map(|(key, _, _)| key).collect()
1305 }
1306}
1307
1308fn send_wakes() {
1311 loop {
1312 let (due, next) = {
1313 let mut wakes = WAKES.lock();
1314 let due = wakes.take_due(std::time::Instant::now());
1315 let next = wakes.next();
1316 if due.is_empty() && next.is_none() {
1317 wakes.timer = false;
1318 return;
1319 }
1320 (due, next)
1321 };
1322 let live = registry().live();
1323 let absent: Vec<(String, String)> = due.into_iter().filter(|d| !live.contains(d)).collect();
1324 if !absent.is_empty() {
1325 wake(&absent);
1326 }
1327 if let Some(next) = next {
1328 std::thread::sleep(next.saturating_duration_since(std::time::Instant::now()));
1329 }
1330 }
1331}
1332
1333pub fn changed_if_library_moved(conn: &rusqlite::Connection) {
1337 static LAST: parking_lot::Mutex<Option<(i64, i64, i64)>> = parking_lot::Mutex::new(None);
1338 let Ok(now) = conn.query_row(
1339 "SELECT (SELECT COUNT(*) FROM tracks), (SELECT COALESCE(MAX(id), 0) FROM tracks),
1340 (SELECT COUNT(*) FROM albums)",
1341 [],
1342 |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
1343 ) else {
1344 return;
1345 };
1346 let before = LAST.lock().replace(now);
1347 if before.is_some_and(|b| b != now) {
1348 changed();
1349 }
1350}
1351
1352mod outbox {
1355 use koan_core::db::queries;
1356 use koan_core::remote::link::LinkCommand;
1357
1358 const KEEP_SECS: i64 = 30 * 24 * 60 * 60;
1363
1364 fn db() -> Option<koan_core::db::pool::Handle<'static>> {
1369 if cfg!(test) {
1370 return None;
1371 }
1372 koan_core::db::pool::shared().get().ok()
1373 }
1374
1375 pub fn load_orders() -> Vec<super::Order> {
1376 let Some(db) = db() else { return Vec::new() };
1377 db.conn
1378 .prepare("SELECT body FROM link_orders ORDER BY created_at")
1379 .and_then(|mut s| {
1380 s.query_map([], |r| r.get::<_, String>(0))?
1381 .collect::<Result<Vec<_>, _>>()
1382 })
1383 .unwrap_or_default()
1384 .into_iter()
1385 .filter_map(|b| serde_json::from_str(&b).ok())
1386 .collect()
1387 }
1388
1389 pub fn save_order(order: &super::Order) {
1390 let (Some(db), Ok(body)) = (db(), serde_json::to_string(order)) else {
1391 return;
1392 };
1393 let _ = db.conn.execute(
1394 "INSERT OR REPLACE INTO link_orders (id, body, created_at) VALUES (?1, ?2, ?3)",
1395 rusqlite::params![order.id, body, order.created_at],
1396 );
1397 }
1398
1399 pub fn drop_order(id: &str) {
1400 if let Some(db) = db() {
1401 let _ = db
1402 .conn
1403 .execute("DELETE FROM link_orders WHERE id = ?1", [id]);
1404 }
1405 }
1406
1407 pub fn take_and_remember(
1408 username: &str,
1409 device: &str,
1410 name: &str,
1411 platform: &str,
1412 live: &[(String, String)],
1413 ) -> Vec<LinkCommand> {
1414 let Some(db) = db() else { return Vec::new() };
1415 let now = chrono::Utc::now().timestamp();
1416 let _ = db.conn.execute(
1417 "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES (?1, ?2, ?3, ?4, ?5)
1418 ON CONFLICT (device, username) DO UPDATE SET name = ?3, platform = ?4, last_seen = ?5",
1419 rusqlite::params![device, username, name, platform, now],
1420 );
1421 forget_stale(&db.conn, live, now);
1422 let waiting: Vec<(i64, String)> = db
1423 .conn
1424 .prepare("SELECT id, command FROM link_outbox WHERE device = ?1 AND username = ?2 ORDER BY id")
1425 .and_then(|mut s| {
1426 s.query_map([device, username], |r| Ok((r.get(0)?, r.get(1)?)))?
1427 .collect()
1428 })
1429 .unwrap_or_default();
1430 let _ = db.conn.execute(
1431 "DELETE FROM link_outbox WHERE device = ?1 AND username = ?2",
1432 [device, username],
1433 );
1434 if !waiting.is_empty() {
1435 log::info!("link: {} waiting commands for {name}", waiting.len());
1436 }
1437 waiting
1438 .into_iter()
1439 .filter_map(|(_, c)| serde_json::from_str(&c).ok())
1440 .collect()
1441 }
1442
1443 pub(super) fn forget_stale(conn: &rusqlite::Connection, live: &[(String, String)], now: i64) {
1446 touch_with(conn, live, now);
1447 let _ = conn.execute(
1448 "DELETE FROM link_outbox WHERE created_at < ?1",
1449 [now - KEEP_SECS],
1450 );
1451 let _ = conn.execute(
1452 "DELETE FROM link_push WHERE (device, username) IN
1453 (SELECT device, username FROM link_devices WHERE last_seen < ?1)",
1454 [now - KEEP_SECS],
1455 );
1456 let _ = conn.execute(
1457 "DELETE FROM link_devices WHERE last_seen < ?1",
1458 [now - KEEP_SECS],
1459 );
1460 }
1461
1462 pub fn touch(devices: &[(String, String)]) {
1464 if let Some(db) = db() {
1465 touch_with(&db.conn, devices, chrono::Utc::now().timestamp());
1466 }
1467 }
1468
1469 fn touch_with(conn: &rusqlite::Connection, devices: &[(String, String)], now: i64) {
1470 for (device, username) in devices {
1471 let _ = conn.execute(
1472 "UPDATE link_devices SET last_seen = ?1 WHERE device = ?2 AND username = ?3",
1473 rusqlite::params![now, device, username],
1474 );
1475 }
1476 }
1477
1478 pub struct Absent {
1480 pub device: String,
1481 pub username: String,
1482 pub name: String,
1483 }
1484
1485 pub fn queue_for_absent(
1487 username: Option<&str>,
1488 live: &[(String, String)],
1489 cmd: &LinkCommand,
1490 ) -> Vec<Absent> {
1491 let Some(db) = db() else { return Vec::new() };
1492 let known: Vec<(String, String, String)> = db
1493 .conn
1494 .prepare("SELECT device, username, name FROM link_devices")
1495 .and_then(|mut s| {
1496 s.query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))?
1497 .collect()
1498 })
1499 .unwrap_or_default();
1500 let Ok(text) = serde_json::to_string(cmd) else {
1501 return Vec::new();
1502 };
1503 let absent: Vec<_> = known
1504 .into_iter()
1505 .filter(|(device, user, _)| {
1506 !username.is_some_and(|u| u != user)
1507 && !live.iter().any(|(d, u)| d == device && u == user)
1508 })
1509 .collect();
1510 if absent.is_empty() {
1511 return Vec::new();
1512 }
1513 let is_sync = matches!(cmd, LinkCommand::Sync { .. });
1514 let now = chrono::Utc::now().timestamp();
1515 koan_core::db::queries::atomically(&db.conn, || {
1518 let mut queued = Vec::new();
1519 for (device, user, name) in absent {
1520 if is_sync {
1521 let _ = db.conn.execute(
1523 "DELETE FROM link_outbox WHERE device = ?1 AND username = ?2 AND command LIKE '{\"type\":\"sync\"%'",
1524 [&device, &user],
1525 );
1526 }
1527 if db
1528 .conn
1529 .execute(
1530 "INSERT INTO link_outbox (device, username, command, created_at) VALUES (?1, ?2, ?3, ?4)",
1531 rusqlite::params![device, user, text, now],
1532 )
1533 .is_ok()
1534 {
1535 queued.push(Absent {
1536 device,
1537 username: user,
1538 name,
1539 });
1540 }
1541 }
1542 Ok::<_, rusqlite::Error>(queued)
1543 })
1544 .unwrap_or_default()
1545 }
1546
1547 #[derive(Clone)]
1549 pub struct PushTarget {
1550 pub device: String,
1551 pub username: String,
1552 pub name: String,
1553 pub platform: String,
1554 pub token: String,
1555 pub sandbox: bool,
1556 pub last_seen: i64,
1558 }
1559
1560 pub fn queue_for(device: &str, username: &str, cmd: &LinkCommand) {
1562 let (Some(db), Ok(text)) = (db(), serde_json::to_string(cmd)) else {
1563 return;
1564 };
1565 let _ = db.conn.execute(
1566 "INSERT INTO link_outbox (device, username, command, created_at) VALUES (?1, ?2, ?3, ?4)",
1567 rusqlite::params![device, username, text, chrono::Utc::now().timestamp()],
1568 );
1569 }
1570
1571 pub fn accounts() -> Vec<String> {
1575 let Some(db) = db() else { return Vec::new() };
1576 koan_core::db::queries::auth::list_users(&db.conn)
1577 .map(|users| users.into_iter().map(|u| u.username).collect())
1578 .unwrap_or_default()
1579 }
1580
1581 pub fn load_addresses() -> std::collections::HashMap<String, (String, std::net::IpAddr)> {
1584 let Some(db) = db() else {
1585 return Default::default();
1586 };
1587 load_addresses_in(&db.conn)
1588 }
1589
1590 pub(super) fn load_addresses_in(
1591 conn: &rusqlite::Connection,
1592 ) -> std::collections::HashMap<String, (String, std::net::IpAddr)> {
1593 conn.prepare(
1594 "SELECT device, username, addr FROM link_devices WHERE addr IS NOT NULL ORDER BY last_seen",
1595 )
1596 .and_then(|mut s| {
1597 s.query_map([], |r| {
1598 Ok((
1599 r.get::<_, String>(0)?,
1600 r.get::<_, String>(1)?,
1601 r.get::<_, String>(2)?,
1602 ))
1603 })?
1604 .collect::<Result<Vec<_>, _>>()
1605 })
1606 .unwrap_or_default()
1607 .into_iter()
1608 .filter_map(|(device, username, addr)| Some((device, (username, addr.parse().ok()?))))
1609 .collect()
1610 }
1611
1612 pub fn save_address(device: &str, username: &str, addr: std::net::IpAddr) {
1613 if let Some(db) = db() {
1614 save_address_in(&db.conn, device, username, addr);
1615 }
1616 }
1617
1618 pub(super) fn save_address_in(
1619 conn: &rusqlite::Connection,
1620 device: &str,
1621 username: &str,
1622 addr: std::net::IpAddr,
1623 ) {
1624 let _ = conn.execute(
1625 "UPDATE link_devices SET addr = ?1 WHERE device = ?2 AND username = ?3",
1626 rusqlite::params![addr.to_string(), device, username],
1627 );
1628 }
1629
1630 pub fn load_grants() -> Vec<super::Grant> {
1631 let Some(db) = db() else { return Vec::new() };
1632 db.conn
1633 .prepare("SELECT device, owner, grantee FROM link_grants ORDER BY created_at")
1634 .and_then(|mut s| {
1635 s.query_map([], |r| {
1636 Ok(super::Grant {
1637 device: r.get(0)?,
1638 owner: r.get(1)?,
1639 grantee: r.get(2)?,
1640 })
1641 })?
1642 .collect()
1643 })
1644 .unwrap_or_default()
1645 }
1646
1647 pub fn save_grant(g: &super::Grant) {
1648 if let Some(db) = db() {
1649 let _ = db.conn.execute(
1650 "INSERT OR IGNORE INTO link_grants (device, owner, grantee, created_at) VALUES (?1, ?2, ?3, ?4)",
1651 rusqlite::params![g.device, g.owner, g.grantee, chrono::Utc::now().timestamp()],
1652 );
1653 }
1654 }
1655
1656 pub fn drop_grant(g: &super::Grant) {
1657 if let Some(db) = db() {
1658 let _ = db.conn.execute(
1659 "DELETE FROM link_grants WHERE device = ?1 AND owner = ?2 AND grantee = ?3",
1660 [&g.device, &g.owner, &g.grantee],
1661 );
1662 }
1663 }
1664
1665 pub fn save_push(username: &str, device: &str, token: &str, sandbox: bool) {
1666 let Some(db) = db() else { return };
1667 let _ = db.conn.execute(
1668 "INSERT INTO link_push (device, username, token, sandbox, updated_at) VALUES (?1, ?2, ?3, ?4, ?5)
1669 ON CONFLICT (device, username) DO UPDATE SET token = ?3, sandbox = ?4, updated_at = ?5",
1670 rusqlite::params![device, username, token, sandbox, chrono::Utc::now().timestamp()],
1671 );
1672 }
1673
1674 pub fn forget_push(username: &str, device: &str) {
1675 if let Some(db) = db() {
1676 let _ = db.conn.execute(
1677 "DELETE FROM link_push WHERE device = ?1 AND username = ?2",
1678 [device, username],
1679 );
1680 }
1681 }
1682
1683 pub fn push_targets(username: Option<&str>) -> Vec<PushTarget> {
1685 let Some(db) = db() else { return Vec::new() };
1686 push_targets_in(&db.conn, username)
1687 }
1688
1689 pub(super) fn push_targets_in(
1690 conn: &rusqlite::Connection,
1691 username: Option<&str>,
1692 ) -> Vec<PushTarget> {
1693 conn
1694 .prepare(
1695 "SELECT p.device, p.username, d.name, d.platform, p.token, p.sandbox, d.last_seen
1696 FROM link_push p JOIN link_devices d ON d.device = p.device AND d.username = p.username
1697 WHERE ?1 IS NULL OR p.username = ?1
1698 ORDER BY d.last_seen DESC",
1699 )
1700 .and_then(|mut s| {
1701 s.query_map([username], |r| {
1702 Ok(PushTarget {
1703 device: r.get(0)?,
1704 username: r.get(1)?,
1705 name: r.get(2)?,
1706 platform: r.get(3)?,
1707 token: r.get(4)?,
1708 sandbox: r.get(5)?,
1709 last_seen: r.get(6)?,
1710 })
1711 })?
1712 .collect()
1713 })
1714 .unwrap_or_default()
1715 }
1716
1717 pub fn forget_device(device: &str, username: &str) {
1720 if let Some(db) = db() {
1721 forget_device_in(&db.conn, device, username);
1722 }
1723 }
1724
1725 pub(super) fn forget_device_in(conn: &rusqlite::Connection, device: &str, username: &str) {
1726 for table in ["link_push", "link_outbox", "link_devices"] {
1727 let _ = conn.execute(
1728 &format!("DELETE FROM {table} WHERE device = ?1 AND username = ?2"),
1729 [device, username],
1730 );
1731 }
1732 }
1733
1734 pub fn track_row(id: &str) -> Option<i64> {
1736 let db = db()?;
1737 queries::resolve_id(&db.conn, queries::UidKind::Track, id)
1738 .ok()
1739 .flatten()
1740 }
1741
1742 pub fn describe(cmd: &LinkCommand) -> Option<String> {
1745 let ids: Vec<&String> = match cmd {
1746 LinkCommand::Play { track_ids, .. }
1747 | LinkCommand::Enqueue { track_ids }
1748 | LinkCommand::PlayNext { track_ids } => track_ids.iter().collect(),
1749 LinkCommand::JumpTo { track_id } => vec![track_id],
1750 _ => return None,
1751 };
1752 let db = db()?;
1753 let ids: Vec<i64> = ids
1754 .into_iter()
1755 .filter_map(|t| {
1756 queries::resolve_id(&db.conn, queries::UidKind::Track, t)
1757 .ok()
1758 .flatten()
1759 })
1760 .collect();
1761 let row = |id: i64| {
1762 db.conn
1763 .query_row(
1764 "SELECT t.title, COALESCE(a.name, ''), COALESCE(al.title, ''), t.album_id
1765 FROM tracks t LEFT JOIN artists a ON a.id = t.artist_id
1766 LEFT JOIN albums al ON al.id = t.album_id WHERE t.id = ?1",
1767 [id],
1768 |r| {
1769 Ok((
1770 r.get::<_, String>(0)?,
1771 r.get::<_, String>(1)?,
1772 r.get::<_, String>(2)?,
1773 r.get::<_, Option<i64>>(3)?,
1774 ))
1775 },
1776 )
1777 .ok()
1778 };
1779 let (title, artist, album, album_id) = row(*ids.first()?)?;
1780 let one_album = ids.len() > 1
1781 && album_id.is_some()
1782 && ids
1783 .iter()
1784 .all(|id| row(*id).is_some_and(|r| r.3 == album_id));
1785 Some(match (one_album, ids.len()) {
1786 (true, _) => format!("{album} — {artist}"),
1787 (false, 1) => format!("{title} — {artist}"),
1788 (false, n) => format!("{title} — {artist}, and {} more", n - 1),
1789 })
1790 }
1791}
1792
1793const RECENT: i64 = 6 * 60 * 60;
1796
1797fn pick(clients: &[ClientInfo], now: i64) -> Result<&ClientInfo, String> {
1798 if let Some(c) = clients.iter().find(|c| c.state.playing) {
1799 return Ok(c);
1800 }
1801 if let Some(c) = clients
1802 .iter()
1803 .filter(|c| c.last_played_at.is_some_and(|t| now - t < RECENT))
1804 .max_by_key(|c| c.last_played_at)
1805 {
1806 return Ok(c);
1807 }
1808 match clients {
1809 [] => Err("no koan app is linked to this server; open koan on the device".into()),
1810 [only] => Ok(only),
1811 several => Err(format!(
1812 "several koan apps are linked and none has played recently: {}. Ask which, then pass `client`",
1813 several
1814 .iter()
1815 .map(|c| format!("{} ({}, id {})", c.name, c.platform, c.device))
1816 .collect::<Vec<_>>()
1817 .join(", ")
1818 )),
1819 }
1820}
1821
1822#[cfg(test)]
1823mod tests {
1824
1825 #[test]
1826 fn a_watch_is_never_queued_and_is_renewed_when_the_target_relinks() {
1827 let reg = Registry::default();
1828 let watches = |rx: &mut tokio::sync::mpsc::UnboundedReceiver<LinkCommand>| {
1829 std::iter::from_fn(|| rx.try_recv().ok())
1830 .filter(|c| matches!(c, LinkCommand::WatchLevels { .. }))
1831 .collect::<Vec<_>>()
1832 };
1833 assert!(
1835 reg.relay("rl", "rl-phone", LinkCommand::WatchLevels { on: true })
1836 .is_err()
1837 );
1838 reg.watch_levels("rl", "rl-mac", "rl-phone", true);
1840 let (tx, mut phone) = tokio::sync::mpsc::unbounded_channel();
1841 reg.register("rl", "phone", "ios", "rl-phone", tx, false);
1842 assert_eq!(
1843 watches(&mut phone),
1844 vec![LinkCommand::WatchLevels { on: true }],
1845 "told once, on linking, because it is watched now"
1846 );
1847 let (tx, mut phone) = tokio::sync::mpsc::unbounded_channel();
1849 reg.register("rl", "phone", "ios", "rl-phone", tx, false);
1850 assert_eq!(
1851 watches(&mut phone),
1852 vec![LinkCommand::WatchLevels { on: true }]
1853 );
1854 }
1855
1856 #[test]
1857 fn levels_reach_a_watcher_only_while_it_watches() {
1858 use koan_core::remote::levels::Frame;
1859 let reg = Registry::default();
1860 let (tx_mac, mut mac) = tokio::sync::mpsc::unbounded_channel();
1861 let (tx_phone, mut phone) = tokio::sync::mpsc::unbounded_channel();
1862 let mac_id = reg.register("lv", "mac", "macos", "lv-mac", tx_mac, false);
1863 reg.register("lv", "phone", "ios", "lv-phone", tx_phone, false);
1864 let levels = |rx: &mut tokio::sync::mpsc::UnboundedReceiver<LinkCommand>| {
1865 std::iter::from_fn(|| rx.try_recv().ok())
1866 .filter(|c| {
1867 matches!(
1868 c,
1869 LinkCommand::Levels { .. } | LinkCommand::WatchLevels { .. }
1870 )
1871 })
1872 .collect::<Vec<_>>()
1873 };
1874 let f = Frame(1_000, 1, 2, 3);
1875
1876 reg.levels("lv", "lv-phone", f);
1877 assert!(levels(&mut mac).is_empty(), "nobody watching");
1878
1879 reg.watch_levels("lv", "lv-mac", "lv-phone", true);
1880 assert_eq!(
1881 levels(&mut phone),
1882 vec![LinkCommand::WatchLevels { on: true }]
1883 );
1884 reg.levels("lv", "lv-phone", f);
1885 assert_eq!(
1886 levels(&mut mac),
1887 vec![LinkCommand::Levels {
1888 from: "lv-phone".into(),
1889 f
1890 }]
1891 );
1892
1893 reg.unregister(&mac_id);
1895 assert_eq!(
1896 levels(&mut phone),
1897 vec![LinkCommand::WatchLevels { on: false }]
1898 );
1899 reg.watch_levels("lv", "lv-mac", "lv-phone", true);
1900 reg.watch_levels("lv", "lv-mac", "lv-phone", false);
1901 assert_eq!(
1902 levels(&mut phone),
1903 vec![
1904 LinkCommand::WatchLevels { on: true },
1905 LinkCommand::WatchLevels { on: false }
1906 ]
1907 );
1908 }
1909
1910 #[test]
1911 fn a_playlist_order_adds_the_named_tracks_once_they_arrive() {
1912 let dir = tempfile::tempdir().unwrap();
1913 let path = dir.path().join("koan.db");
1914 let db = koan_core::db::connection::Database::open(&path).unwrap();
1915 let playlist = koan_core::db::queries::create_playlist(
1916 &db.conn,
1917 koan_core::db::queries::LOCAL_USER,
1918 "cyberpunk",
1919 None,
1920 )
1921 .unwrap();
1922 let order = Order {
1923 id: "o1".into(),
1924 username: None,
1925 client: None,
1926 artist: "Perturbator".into(),
1927 album: "Dangerous Days".into(),
1928 play_next: false,
1929 playlist: Some(playlist),
1930 titles: vec!["Future Club".into()],
1931 created_at: chrono::Utc::now().timestamp(),
1932 };
1933 registry().add_order(order);
1934
1935 fulfil_from(&path);
1937 assert!(registry().orders(None).iter().any(|o| o.id == "o1"));
1938
1939 db.conn
1940 .execute_batch(
1941 "INSERT INTO artists (id, name) VALUES (1, 'Perturbator');
1942 INSERT INTO albums (id, title, artist_id) VALUES (1, 'Dangerous Days', 1);
1943 INSERT INTO tracks (id, title, album_id, artist_id, track_number, path) VALUES
1944 (1, 'Welcome Back', 1, 1, 1, '/1.flac'), (2, 'Future Club', 1, 1, 2, '/2.flac');",
1945 )
1946 .unwrap();
1947 fulfil_from(&path);
1948 let held: Vec<i64> = db
1949 .conn
1950 .prepare("SELECT track_id FROM playlist_tracks WHERE playlist_id = ?1")
1951 .unwrap()
1952 .query_map([playlist], |r| r.get(0))
1953 .unwrap()
1954 .collect::<Result<_, _>>()
1955 .unwrap();
1956 assert_eq!(held, [2]);
1957 assert!(!registry().orders(None).iter().any(|o| o.id == "o1"));
1958 }
1959
1960 use super::*;
1961
1962 #[test]
1963 fn a_notification_cover_is_named_by_the_uid_the_command_carries() {
1964 let uid = "0190a5b2-7c3d-7e4f-8a1b-2c3d4e5f6a7b".to_string();
1966 let play = LinkCommand::Play {
1967 track_ids: vec!["x".into(), uid.clone()],
1968 start_at: 1,
1969 position_ms: 0,
1970 paused: false,
1971 handoff: false,
1972 };
1973 assert_eq!(cover_track(&play), Some(uid.as_str()));
1974 }
1975
1976 #[test]
1977 fn a_reconnect_replaces_the_device_and_commands_reach_it() {
1978 let reg = Registry::default();
1979 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
1980 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
1981 let (tx3, _rx3) = tokio::sync::mpsc::unbounded_channel();
1982 reg.register("j", "phone", "ios", "dev-1", tx1, false);
1983 let id = reg.register("j", "phone", "ios", "dev-1", tx2, false);
1984 reg.register("someone", "laptop", "macos", "dev-2", tx3, false);
1985
1986 assert_eq!(reg.list(Some("j")).len(), 1);
1987 assert_eq!(reg.list(None).len(), 2);
1988
1989 let sent = reg.send(Some("j"), None, LinkCommand::Pause).unwrap();
1990 assert_eq!(sent.id, id);
1991 assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
1992
1993 assert!(
1995 reg.send(Some("j"), Some("laptop"), LinkCommand::Pause)
1996 .is_err()
1997 );
1998
1999 reg.unregister(&id);
2000 assert!(reg.send(Some("j"), None, LinkCommand::Pause).is_err());
2001 }
2002
2003 #[test]
2004 fn disconnecting_an_account_closes_only_its_links() {
2005 let reg = Registry::default();
2006 let (tx1, mut rx1) = tokio::sync::mpsc::unbounded_channel();
2007 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
2008 reg.register("j", "phone", "ios", "dev-1", tx1, false);
2009 reg.register("someone", "laptop", "macos", "dev-2", tx2, false);
2010
2011 reg.disconnect("j");
2012 assert!(reg.list(Some("j")).is_empty());
2013 assert!(matches!(
2015 rx1.try_recv(),
2016 Err(tokio::sync::mpsc::error::TryRecvError::Disconnected)
2017 ));
2018 assert_eq!(reg.list(Some("someone")).len(), 1);
2019 assert!(matches!(
2020 rx2.try_recv(),
2021 Err(tokio::sync::mpsc::error::TryRecvError::Empty)
2022 ));
2023 }
2024
2025 #[test]
2026 fn a_burst_of_changes_wakes_each_device_once_when_it_goes_quiet() {
2027 let t0 = std::time::Instant::now();
2028 let s = std::time::Duration::from_secs;
2029 let phone = || ("dev-1".to_string(), "j".to_string());
2030 let ipad = || ("dev-2".to_string(), "j".to_string());
2031 let mut wakes = Wakes {
2032 pending: Vec::new(),
2033 timer: false,
2034 };
2035
2036 wakes.add([phone()], t0);
2037 wakes.add([phone(), ipad()], t0 + s(10));
2038 wakes.add([phone()], t0 + s(20));
2039 assert_eq!(wakes.pending.len(), 2);
2040
2041 assert!(wakes.take_due(t0 + s(39)).is_empty());
2043 assert_eq!(wakes.take_due(t0 + s(40)), [ipad()]);
2044 assert_eq!(wakes.next(), Some(t0 + s(50)));
2045 assert_eq!(wakes.take_due(t0 + s(50)), [phone()]);
2046 assert_eq!(wakes.next(), None);
2047 }
2048
2049 #[test]
2050 fn a_library_that_never_goes_quiet_still_wakes_devices() {
2051 let t0 = std::time::Instant::now();
2052 let phone = || ("dev-1".to_string(), "j".to_string());
2053 let mut wakes = Wakes {
2054 pending: Vec::new(),
2055 timer: false,
2056 };
2057 let mut sent = 0;
2058 for i in 0..40 {
2059 let now = t0 + std::time::Duration::from_secs(i * 10);
2060 sent += wakes.take_due(now).len();
2061 wakes.add([phone()], now);
2062 }
2063 assert_eq!(sent, 1);
2066 assert_eq!(wakes.pending.len(), 1);
2067 }
2068
2069 #[test]
2070 fn the_device_playing_is_the_one_meant() {
2071 let reg = Registry::default();
2072 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
2073 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
2074 let mac = reg.register("j", "mac", "macos", "dev-1", tx1, false);
2075 let phone = reg.register("j", "phone", "ios", "dev-2", tx2, false);
2076
2077 let err = reg.send(Some("j"), None, LinkCommand::Pause).unwrap_err();
2079 assert!(err.contains("mac") && err.contains("phone"), "{err}");
2080
2081 reg.report(
2082 &phone,
2083 LinkState {
2084 playing: true,
2085 ..Default::default()
2086 },
2087 );
2088 assert_eq!(
2089 reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
2090 phone
2091 );
2092 assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
2093
2094 reg.report(&phone, LinkState::default());
2096 assert_eq!(
2097 reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
2098 phone
2099 );
2100 let _ = mac;
2101 }
2102
2103 fn drain(rx: &mut tokio::sync::mpsc::UnboundedReceiver<LinkCommand>) -> Vec<LinkCommand> {
2104 std::iter::from_fn(|| rx.try_recv().ok()).collect()
2105 }
2106
2107 fn shared() -> (
2109 Registry,
2110 tokio::sync::mpsc::UnboundedReceiver<LinkCommand>,
2111 tokio::sync::mpsc::UnboundedReceiver<LinkCommand>,
2112 tokio::sync::mpsc::UnboundedReceiver<LinkCommand>,
2113 ) {
2114 let reg = Registry::default();
2115 let (phone_tx, mut phone) = tokio::sync::mpsc::unbounded_channel();
2116 let (mac_tx, mac) = tokio::sync::mpsc::unbounded_channel();
2117 let (k_tx, k) = tokio::sync::mpsc::unbounded_channel();
2118 let id = reg.register("j", "phone", "ios", "dev-phone", phone_tx, true);
2119 reg.register("j", "mac", "macos", "dev-mac", mac_tx, true);
2120 reg.register("k", "laptop", "macos", "dev-k", k_tx, true);
2121 reg.report(
2122 &id,
2123 LinkState {
2124 playing: true,
2125 title: Some("Roygbiv".into()),
2126 outputs: Some(Default::default()),
2127 ..Default::default()
2128 },
2129 );
2130 reg.share("j", "dev-phone", "k", true).unwrap();
2131 drain(&mut phone);
2132 (reg, phone, mac, k)
2133 }
2134
2135 #[test]
2136 fn a_grantee_sees_the_shared_device_and_nothing_else_of_the_owner() {
2137 let (_reg, _phone, _mac, mut k) = shared();
2138 let listed = drain(&mut k)
2139 .into_iter()
2140 .rev()
2141 .find_map(|c| match c {
2142 LinkCommand::Devices { devices } => Some(devices),
2143 _ => None,
2144 })
2145 .expect("told of it");
2146 assert_eq!(listed.len(), 1, "the phone, not the Mac");
2147 let phone = &listed[0];
2148 assert_eq!(phone.id, "dev-phone");
2149 assert_eq!(phone.owner.as_deref(), Some("j"));
2150 let state = phone.state.as_ref().expect("its state");
2151 assert_eq!(state.title.as_deref(), Some("Roygbiv"));
2152 assert!(
2153 state.outputs.is_some(),
2154 "outputs, to choose from: the output is in the playback set"
2155 );
2156 }
2157
2158 #[test]
2164 fn a_grantee_controls_playback_as_itself_and_nothing_of_the_owners_account() {
2165 let (reg, mut phone, mut mac, _k) = shared();
2166 let wrapped = |c: LinkCommand| LinkCommand::Shared {
2167 command: Box::new(c),
2168 };
2169 for cmd in [
2170 LinkCommand::Pause,
2171 LinkCommand::SetRendererVolume { volume: 40 },
2172 LinkCommand::HandOff { to: "dev-k".into() },
2173 ] {
2174 reg.relay("k", "dev-phone", cmd.clone()).unwrap();
2175 assert_eq!(drain(&mut phone), [wrapped(cmd)]);
2176 }
2177 for cmd in [
2178 LinkCommand::Sync { full: false },
2179 LinkCommand::Evict { track_ids: vec![] },
2180 LinkCommand::Shared {
2181 command: Box::new(LinkCommand::Pause),
2182 },
2183 ] {
2184 assert!(!cmd.allowed_playback());
2185 assert!(reg.relay("k", "dev-phone", cmd).is_err());
2186 }
2187 assert!(drain(&mut phone).is_empty());
2188 assert!(
2189 reg.relay("k", "dev-mac", LinkCommand::Pause).is_err(),
2190 "not shared, not reachable"
2191 );
2192 assert!(
2193 drain(&mut mac)
2194 .iter()
2195 .all(|c| matches!(c, LinkCommand::Devices { .. } | LinkCommand::Shares { .. })),
2196 "news, and no command"
2197 );
2198 assert!(reg.shared_owner("k", "dev-mac").is_none(), "nor wakeable");
2199 }
2200
2201 #[test]
2205 fn a_hand_off_runs_both_ways_across_a_grant_and_nothing_else_does() {
2206 let (reg, _phone, _mac, mut k) = shared();
2207 drain(&mut k);
2208 let play = LinkCommand::Play {
2209 track_ids: vec!["t".into()],
2210 start_at: 0,
2211 position_ms: 1000,
2212 paused: false,
2213 handoff: true,
2214 };
2215 reg.relay_from("j", Some("dev-phone"), "dev-k", play.clone())
2216 .unwrap();
2217 assert!(drain(&mut k).contains(&LinkCommand::Shared {
2218 command: Box::new(play.clone())
2219 }));
2220 assert!(
2221 reg.relay_from("j", Some("dev-phone"), "dev-k", LinkCommand::Pause)
2222 .is_err(),
2223 "the owner's phone does not command the grantee's devices"
2224 );
2225 assert!(
2226 reg.relay_from("j", Some("dev-mac"), "dev-k", play.clone())
2227 .is_err(),
2228 "only the shared device"
2229 );
2230 let not_a_hand_off = LinkCommand::Play {
2231 track_ids: vec!["t".into()],
2232 start_at: 0,
2233 position_ms: 0,
2234 paused: false,
2235 handoff: false,
2236 };
2237 assert!(
2238 reg.relay_from("j", Some("dev-phone"), "dev-k", not_a_hand_off)
2239 .is_err(),
2240 "a hand-off, not any play"
2241 );
2242 }
2243
2244 #[test]
2245 fn revoking_ends_control_at_once() {
2246 let (reg, mut phone, _mac, mut k) = shared();
2247 reg.share("j", "dev-phone", "k", false).unwrap();
2248 assert!(reg.relay("k", "dev-phone", LinkCommand::Pause).is_err());
2249 assert!(
2250 drain(&mut phone)
2251 .iter()
2252 .all(|c| !matches!(c, LinkCommand::Pause))
2253 );
2254 assert!(reg.shared_owner("k", "dev-phone").is_none());
2255 let listed = drain(&mut k)
2256 .into_iter()
2257 .rev()
2258 .find_map(|c| match c {
2259 LinkCommand::Devices { devices } => Some(devices),
2260 _ => None,
2261 })
2262 .expect("told it is gone");
2263 assert!(listed.is_empty());
2264 }
2265
2266 #[test]
2267 fn a_device_hears_whom_it_is_shared_with_every_time_it_links() {
2268 let reg = Registry::default();
2269 let (tx, mut phone) = tokio::sync::mpsc::unbounded_channel();
2270 reg.register("j", "phone", "ios", "dev-phone", tx, true);
2271 let first = drain(&mut phone);
2272 assert!(
2273 first.contains(&LinkCommand::Shares {
2274 grantees: vec![],
2275 error: None,
2276 accounts: vec![],
2277 }),
2278 "an empty list replaces one from another server"
2279 );
2280 reg.send_shares(
2281 "j",
2282 "dev-phone",
2283 Some("There is no account called x".into()),
2284 );
2285 assert!(drain(&mut phone).iter().any(
2286 |c| matches!(c, LinkCommand::Shares { error: Some(e), .. } if e.contains("no account"))
2287 ));
2288 }
2289
2290 #[test]
2291 fn the_owner_is_told_who_it_shares_with() {
2292 let reg = Registry::default();
2293 let (tx, mut phone) = tokio::sync::mpsc::unbounded_channel();
2294 reg.register("j", "phone", "ios", "dev-phone", tx, true);
2295 reg.share("j", "dev-phone", "k", true).unwrap();
2296 reg.share("j", "dev-phone", "m", true).unwrap();
2297 let last = drain(&mut phone).into_iter().rev().find_map(|c| match c {
2298 LinkCommand::Shares { grantees, .. } => Some(grantees),
2299 _ => None,
2300 });
2301 assert_eq!(last, Some(vec!["k".to_string(), "m".to_string()]));
2302 assert!(
2303 reg.share("j", "dev-phone", "j", true).is_err(),
2304 "not with itself"
2305 );
2306 }
2307
2308 #[test]
2312 fn a_device_is_woken_for_another_account_on_its_network_or_by_grant() {
2313 let reg = Registry::default();
2314 let home: std::net::IpAddr = "203.0.113.7".parse().unwrap();
2315 let away: std::net::IpAddr = "198.51.100.2".parse().unwrap();
2316 reg.seen_at("dev-ipad", "sarita", home);
2317 reg.seen_at("dev-phone", "admin", home);
2318 assert_eq!(
2319 reg.wake_owner("admin", "dev-phone", "dev-ipad"),
2320 "sarita",
2321 "same address: the iPad's own token"
2322 );
2323 reg.seen_at("dev-phone", "admin", away);
2324 assert_eq!(
2325 reg.wake_owner("admin", "dev-phone", "dev-ipad"),
2326 "admin",
2327 "elsewhere, with no grant: only admin's own, which it is not"
2328 );
2329 reg.share("sarita", "dev-ipad", "admin", true).unwrap();
2330 assert_eq!(
2331 reg.wake_owner("admin", "dev-phone", "dev-ipad"),
2332 "sarita",
2333 "shared: from anywhere"
2334 );
2335 reg.seen_at("dev-other", "mallory", home);
2337 assert_eq!(reg.wake_owner("admin", "dev-other", "dev-tv"), "admin");
2338 }
2339
2340 #[test]
2343 fn where_a_device_last_linked_from_survives_a_restart() {
2344 let conn = rusqlite::Connection::open_in_memory().unwrap();
2345 koan_core::db::schema::create_tables(&conn).unwrap();
2346 for (device, user, seen) in [("dev-ipad", "sarita", 1), ("dev-phone", "admin", 2)] {
2347 conn.execute(
2348 "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES (?1, ?2, ?1, 'ios', ?3)",
2349 rusqlite::params![device, user, seen],
2350 )
2351 .unwrap();
2352 }
2353 let home: std::net::IpAddr = "203.0.113.7".parse().unwrap();
2354 outbox::save_address_in(&conn, "dev-ipad", "sarita", home);
2355 outbox::save_address_in(&conn, "dev-phone", "admin", home);
2356
2357 let reg = Registry::default();
2359 *reg.addresses.lock() = outbox::load_addresses_in(&conn);
2360 assert_eq!(
2361 reg.wake_owner("admin", "dev-phone", "dev-ipad"),
2362 "sarita",
2363 "still woken through its own account after the restart"
2364 );
2365 }
2366
2367 #[test]
2371 fn a_forgotten_device_is_dropped_by_the_accounts_devices() {
2372 let reg = Registry::default();
2373 let (mac_tx, mut mac) = tokio::sync::mpsc::unbounded_channel();
2374 let (other_tx, mut other) = tokio::sync::mpsc::unbounded_channel();
2375 reg.register("j", "mac", "macos", "dev-mac", mac_tx, true);
2376 reg.register("k", "laptop", "macos", "dev-k", other_tx, true);
2377 while mac.try_recv().is_ok() {}
2378 while other.try_recv().is_ok() {}
2379
2380 reg.forget("j", "dev-phone").unwrap();
2381 let heard: Vec<LinkCommand> = std::iter::from_fn(|| mac.try_recv().ok()).collect();
2382 assert!(heard.contains(&LinkCommand::Forgotten {
2383 device: "dev-phone".into()
2384 }));
2385 assert!(
2386 std::iter::from_fn(|| other.try_recv().ok())
2387 .all(|c| !matches!(c, LinkCommand::Forgotten { .. })),
2388 "another account hears nothing of it"
2389 );
2390 assert!(
2391 reg.forget("j", "dev-mac").is_err(),
2392 "linked: it would be back"
2393 );
2394 assert!(
2395 reg.relay(
2396 "j",
2397 "dev-mac",
2398 LinkCommand::Forgotten { device: "x".into() }
2399 )
2400 .is_err(),
2401 "news from the server, not a command a device may send"
2402 );
2403 }
2404
2405 #[test]
2408 fn forgetting_a_shared_device_declines_the_share() {
2409 let (reg, mut phone, _mac, mut k) = shared();
2410 drain(&mut k);
2411 reg.forget("k", "dev-phone").unwrap();
2412 assert!(reg.shared_owner("k", "dev-phone").is_none());
2413 assert!(
2414 drain(&mut phone)
2415 .iter()
2416 .any(|c| matches!(c, LinkCommand::Shares { grantees, .. } if grantees.is_empty()))
2417 );
2418 let listed = drain(&mut k).into_iter().rev().find_map(|c| match c {
2419 LinkCommand::Devices { devices } => Some(devices),
2420 _ => None,
2421 });
2422 assert_eq!(listed.map(|d| d.len()), Some(0));
2423 }
2424
2425 #[test]
2428 fn forgetting_a_device_stops_pushes_to_it_and_only_for_its_account() {
2429 let conn = rusqlite::Connection::open_in_memory().unwrap();
2430 koan_core::db::schema::create_tables(&conn).unwrap();
2431 for user in ["j", "k"] {
2432 conn.execute(
2433 "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES ('dev-phone', ?1, 'phone', 'ios', 1)",
2434 [user],
2435 )
2436 .unwrap();
2437 conn.execute(
2438 "INSERT INTO link_push (device, username, token, sandbox, updated_at) VALUES ('dev-phone', ?1, 't', 0, 1)",
2439 [user],
2440 )
2441 .unwrap();
2442 }
2443 outbox::forget_device_in(&conn, "dev-phone", "k");
2444 assert_eq!(
2445 outbox::push_targets_in(&conn, Some("j")).len(),
2446 1,
2447 "k forgetting its own leaves j's alone"
2448 );
2449 outbox::forget_device_in(&conn, "dev-phone", "j");
2450 assert!(outbox::push_targets_in(&conn, Some("j")).is_empty());
2451 assert!(outbox::push_targets_in(&conn, None).is_empty());
2452 }
2453
2454 #[test]
2455 fn a_device_hears_its_peers_and_commands_reach_them_by_device_id() {
2456 let reg = Registry::default();
2457 let (tx1, mut rx1) = tokio::sync::mpsc::unbounded_channel();
2458 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
2459 let (tx3, mut rx3) = tokio::sync::mpsc::unbounded_channel();
2460 reg.register("j", "mac", "macos", "dev-mac", tx1, true);
2461 let phone = reg.register("j", "phone", "ios", "dev-phone", tx2, true);
2462 reg.register("someone", "laptop", "macos", "dev-other", tx3, true);
2463
2464 reg.report(
2465 &phone,
2466 LinkState {
2467 playing: true,
2468 title: Some("Roygbiv".into()),
2469 ..Default::default()
2470 },
2471 );
2472 let devices = drain(&mut rx1)
2473 .into_iter()
2474 .rev()
2475 .find_map(|c| match c {
2476 LinkCommand::Devices { devices } => Some(devices),
2477 _ => None,
2478 })
2479 .expect("the Mac is told of the phone");
2480 assert_eq!(devices.len(), 1, "not itself, not another account's");
2481 assert_eq!(devices[0].id, "dev-phone");
2482 assert_eq!(
2483 devices[0].state.as_ref().and_then(|s| s.title.as_deref()),
2484 Some("Roygbiv")
2485 );
2486 while rx3.try_recv().is_ok() {}
2487 assert!(rx3.try_recv().is_err());
2488
2489 while rx2.try_recv().is_ok() {}
2490 reg.relay("j", "dev-phone", LinkCommand::Pause).unwrap();
2491 assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
2492 assert!(
2493 reg.relay("someone", "dev-phone", LinkCommand::Pause)
2494 .is_err(),
2495 "another account cannot reach it"
2496 );
2497 assert!(
2498 reg.relay("j", "dev-phone", LinkCommand::Devices { devices: vec![] })
2499 .is_err()
2500 );
2501 }
2502
2503 #[test]
2504 fn a_device_linked_for_a_month_is_not_forgotten() {
2505 let conn = rusqlite::Connection::open_in_memory().unwrap();
2506 koan_core::db::schema::create_tables(&conn).unwrap();
2507 let now = 100 * 24 * 60 * 60;
2508 let long_ago = now - 40 * 24 * 60 * 60;
2509 for device in ["mac", "old-phone"] {
2510 conn.execute(
2511 "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES (?1, 'j', ?1, 'ios', ?2)",
2512 rusqlite::params![device, long_ago],
2513 )
2514 .unwrap();
2515 conn.execute(
2516 "INSERT INTO link_push (device, username, token, sandbox, updated_at) VALUES (?1, 'j', 't', 0, ?2)",
2517 rusqlite::params![device, long_ago],
2518 )
2519 .unwrap();
2520 }
2521 outbox::forget_stale(&conn, &[("mac".into(), "j".into())], now);
2522 let left = |table: &str| -> Vec<String> {
2523 conn.prepare(&format!("SELECT device FROM {table} ORDER BY device"))
2524 .unwrap()
2525 .query_map([], |r| r.get(0))
2526 .unwrap()
2527 .collect::<Result<_, _>>()
2528 .unwrap()
2529 };
2530 assert_eq!(left("link_devices"), ["mac"]);
2531 assert_eq!(left("link_push"), ["mac"]);
2532 }
2533}