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::HistoryChanged
707 | LinkCommand::Levels { .. }
708 | LinkCommand::WatchLevels { .. }
709 | LinkCommand::Shares { .. }
710 | LinkCommand::Shared { .. }
711 ) {
712 return Err("not a command".into());
713 }
714 if let Some(owner) = self.shared_owner(username, to) {
718 if !command.allowed_playback() {
719 return Err(format!("{to} is shared for playback only"));
720 }
721 let command = LinkCommand::Shared {
722 command: Box::new(command),
723 };
724 return self.send(Some(&owner), Some(to), command);
725 }
726 if let Some(from) = from
730 && let Some(grantee) = self.grantee_of(username, from, to)
731 {
732 if !matches!(command, LinkCommand::Play { handoff: true, .. }) {
733 return Err(format!("{to} is another account's"));
734 }
735 let command = LinkCommand::Shared {
736 command: Box::new(command),
737 };
738 return self.send(Some(&grantee), Some(to), command);
739 }
740 self.send(Some(username), Some(to), command)
741 }
742
743 pub fn set_activity(
746 &self,
747 username: &str,
748 watcher: &str,
749 activity: Option<(String, String, bool)>,
750 ) {
751 let mut activities = self.activities.lock();
752 activities.retain(|a| !(a.username == username && a.watcher == watcher));
753 let Some((token, target, sandbox)) = activity else {
754 return;
755 };
756 activities.push(Activity {
757 username: username.to_string(),
758 watcher: watcher.to_string(),
759 target: target.clone(),
760 token,
761 sandbox,
762 sent: None,
763 });
764 drop(activities);
765 self.update_activities(username, &target);
766 }
767
768 fn update_activities(&self, username: &str, target: &str) {
771 let Some(pusher) = crate::push::pusher() else {
772 return;
773 };
774 let owner = self
777 .shared_owner(username, target)
778 .unwrap_or_else(|| username.to_string());
779 let Some(info) = self
780 .list(Some(&owner))
781 .into_iter()
782 .find(|c| c.device == target)
783 else {
784 return;
785 };
786 let watchers: Vec<String> = self
787 .grants
788 .lock()
789 .iter()
790 .filter(|g| g.owner == owner && g.device == target)
791 .map(|g| g.grantee.clone())
792 .chain(std::iter::once(owner.clone()))
793 .collect();
794 let state = crate::push::ActivityState::of(&info);
795 let mut due = Vec::new();
796 for a in self.activities.lock().iter_mut() {
797 if watchers.contains(&a.username)
798 && a.target == target
799 && a.sent.as_ref().is_none_or(|s| s.differs(&state))
800 {
801 a.sent = Some(state.clone());
802 due.push((
803 a.token.clone(),
804 a.sandbox,
805 a.watcher.clone(),
806 a.username.clone(),
807 ));
808 }
809 }
810 if due.is_empty() {
811 return;
812 }
813 std::thread::spawn(move || {
814 for (token, sandbox, watcher, username) in due {
815 let push = crate::push::Push::Activity(state.clone());
816 match pusher.send(&token, sandbox, &push) {
817 crate::push::Outcome::Sent => {}
818 crate::push::Outcome::Gone => {
819 log::info!("push: a Live Activity on {watcher} has ended");
820 registry().set_activity(&username, &watcher, None);
821 }
822 crate::push::Outcome::Failed(e) => {
823 log::warn!("push: Live Activity on {watcher}: {e}");
824 }
825 }
826 }
827 });
828 }
829
830 pub fn list(&self, username: Option<&str>) -> Vec<ClientInfo> {
832 let mut out: Vec<ClientInfo> = self
833 .entries
834 .lock()
835 .iter()
836 .filter(|e| username.is_none_or(|u| e.info.username == u))
837 .map(|e| e.info.clone())
838 .collect();
839 out.sort_by_key(|c| std::cmp::Reverse(c.connected_at));
840 out
841 }
842
843 pub fn send(
848 &self,
849 username: Option<&str>,
850 id: Option<&str>,
851 cmd: LinkCommand,
852 ) -> Result<ClientInfo, String> {
853 let clients = self.list(username);
854 let target = match id {
855 Some(id) => clients
856 .iter()
857 .find(|c| c.id == id || c.device == id || c.name.eq_ignore_ascii_case(id)),
858 None if clients.is_empty() => None,
859 None => Some(pick(&clients, chrono::Utc::now().timestamp())?),
860 };
861 let Some(target) = target else {
863 return reach_absent(username, id, &cmd).unwrap_or_else(|| {
864 Err(match id {
865 Some(id) => format!("no linked client {id}; see `clients`"),
866 None => "no koan app is linked to this server; open koan on the device".into(),
867 })
868 });
869 };
870 let entries = self.entries.lock();
871 let entry = entries
872 .iter()
873 .find(|e| e.info.id == target.id)
874 .ok_or("that client has just gone")?;
875 entry
876 .tx
877 .send(cmd)
878 .map_err(|_| "that client has just gone".to_string())?;
879 Ok(target.clone())
880 }
881}
882
883impl Registry {
884 pub fn add_order(&self, order: Order) {
885 outbox::save_order(&order);
886 self.orders.lock().push(order);
887 }
888
889 pub fn orders(&self, username: Option<&str>) -> Vec<Order> {
890 self.orders
891 .lock()
892 .iter()
893 .filter(|o| username.is_none() || o.username.as_deref() == username)
894 .cloned()
895 .collect()
896 }
897
898 pub fn cancel_order(&self, username: Option<&str>, id: &str) -> bool {
899 let mut orders = self.orders.lock();
900 let before = orders.len();
901 orders
902 .retain(|o| !(o.id == id && (username.is_none() || o.username.as_deref() == username)));
903 let gone = orders.len() != before;
904 if gone {
905 outbox::drop_order(id);
906 }
907 gone
908 }
909
910 fn done(&self, id: &str) {
911 self.orders.lock().retain(|o| o.id != id);
912 outbox::drop_order(id);
913 }
914
915 pub fn fulfil_orders(&self, find: impl Fn(&Order) -> Option<Vec<String>>) {
918 let now = chrono::Utc::now().timestamp();
919 let pending: Vec<Order> = {
920 let mut orders = self.orders.lock();
921 for o in orders.iter().filter(|o| now - o.created_at >= ORDER_TTL) {
922 outbox::drop_order(&o.id);
923 }
924 orders.retain(|o| now - o.created_at < ORDER_TTL);
925 orders.clone()
926 };
927 for order in pending.into_iter().filter(|o| o.playlist.is_none()) {
928 let Some(ids) = find(&order).filter(|ids| !ids.is_empty()) else {
929 continue;
930 };
931 let track_ids = ids;
932 let cmd = if order.play_next {
933 LinkCommand::PlayNext { track_ids }
934 } else {
935 LinkCommand::Enqueue { track_ids }
936 };
937 match self.send(order.username.as_deref(), order.client.as_deref(), cmd) {
938 Ok(c) => {
939 log::info!(
940 "link: {} — {} arrived; queued on {}",
941 order.artist,
942 order.album,
943 c.name
944 );
945 self.done(&order.id);
946 }
947 Err(e) => log::info!("link: {} — {} arrived but {e}", order.artist, order.album),
949 }
950 }
951 }
952}
953
954pub fn fulfil_from(db_path: &std::path::Path) {
956 let registry = registry();
957 if registry.orders.lock().is_empty() {
958 return;
959 }
960 let Ok(db) = koan_core::db::connection::Database::open_existing(db_path) else {
961 return;
962 };
963 let for_playlists: Vec<Order> = registry
966 .orders
967 .lock()
968 .iter()
969 .filter(|o| o.playlist.is_some())
970 .cloned()
971 .collect();
972 let mut edited = false;
973 for order in for_playlists {
974 let Some(playlist) = order.playlist else {
975 continue;
976 };
977 let Some(ids) = order_tracks(&db.conn, &order) else {
978 continue;
979 };
980 match koan_core::db::queries::add_tracks(&db.conn, playlist, &ids) {
981 Ok(_) => {
982 if !cfg!(test) {
984 koan_core::playlists::push_to_remote(playlist);
985 }
986 log::info!(
987 "link: {} — {} arrived; added {} tracks to playlist {playlist}",
988 order.artist,
989 order.album,
990 ids.len()
991 );
992 registry.done(&order.id);
993 edited = true;
994 }
995 Err(e) => log::warn!("link: could not add to playlist {playlist}: {e}"),
996 }
997 }
998 if edited {
999 changed();
1000 }
1001 registry.fulfil_orders(|order| {
1002 let rows = order_tracks(&db.conn, order)?;
1003 queries::uids_in_order(&db.conn, UidKind::Track, &rows).ok()
1004 });
1005}
1006
1007fn order_tracks(conn: &rusqlite::Connection, order: &Order) -> Option<Vec<i64>> {
1010 let tracks = album_tracks(conn, &order.artist, &order.album)?;
1011 if order.titles.is_empty() {
1012 return Some(tracks.into_iter().map(|(id, _)| id).collect());
1013 }
1014 let picked: Vec<i64> = order
1015 .titles
1016 .iter()
1017 .filter_map(|want| {
1018 let want = want.to_lowercase();
1019 tracks
1020 .iter()
1021 .find(|(_, t)| t.to_lowercase().contains(&want))
1022 .map(|(id, _)| *id)
1023 })
1024 .collect();
1025 (!picked.is_empty()).then_some(picked)
1026}
1027
1028pub fn album_tracks(
1031 conn: &rusqlite::Connection,
1032 artist: &str,
1033 album: &str,
1034) -> Option<Vec<(i64, String)>> {
1035 let like = |s: &str| format!("%{}%", s.replace(['%', '_'], ""));
1036 let album_id: i64 = conn
1037 .query_row(
1038 "SELECT al.id FROM albums al JOIN artists a ON a.id = al.artist_id
1039 WHERE a.name LIKE ?1 COLLATE NOCASE AND al.title LIKE ?2 COLLATE NOCASE
1040 ORDER BY al.id DESC LIMIT 1",
1041 [like(artist), like(album)],
1042 |r| r.get(0),
1043 )
1044 .ok()?;
1045 let mut stmt = conn
1046 .prepare("SELECT id, title FROM tracks WHERE album_id = ?1 ORDER BY disc, track_number, id")
1047 .ok()?;
1048 let tracks = stmt
1049 .query_map([album_id], |r| Ok((r.get(0)?, r.get(1)?)))
1050 .ok()?
1051 .filter_map(Result::ok)
1052 .collect();
1053 Some(tracks)
1054}
1055
1056impl Registry {
1057 pub fn broadcast(&self, username: Option<&str>, cmd: LinkCommand) -> Vec<String> {
1060 let ids: Vec<String> = self.list(username).into_iter().map(|c| c.id).collect();
1061 let entries = self.entries.lock();
1062 entries
1063 .iter()
1064 .filter(|e| ids.contains(&e.info.id) && e.tx.send(cmd.clone()).is_ok())
1065 .map(|e| e.info.name.clone())
1066 .collect()
1067 }
1068}
1069
1070impl Registry {
1071 pub fn deliver(&self, username: Option<&str>, cmd: LinkCommand) -> (Vec<String>, Vec<String>) {
1076 let (sent, queued) = self.link_or_queue(username, &cmd);
1077 wake(&queued.iter().map(Absent::key).collect::<Vec<_>>());
1078 (sent, queued.into_iter().map(|q| q.name).collect())
1079 }
1080
1081 fn link_or_queue(
1082 &self,
1083 username: Option<&str>,
1084 cmd: &LinkCommand,
1085 ) -> (Vec<String>, Vec<Absent>) {
1086 let sent = self.broadcast(username, cmd.clone());
1087 let queued = outbox::queue_for_absent(username, &self.live(), cmd);
1088 (sent, queued)
1089 }
1090
1091 fn live(&self) -> Vec<(String, String)> {
1093 self.entries
1094 .lock()
1095 .iter()
1096 .map(|e| (e.device.clone(), e.info.username.clone()))
1097 .collect()
1098 }
1099}
1100
1101impl Absent {
1102 fn key(&self) -> (String, String) {
1103 (self.device.clone(), self.username.clone())
1104 }
1105}
1106
1107fn wake(devices: &[(String, String)]) {
1110 let Some(pusher) = crate::push::pusher() else {
1111 return;
1112 };
1113 let targets: Vec<outbox::PushTarget> = outbox::push_targets(None)
1114 .into_iter()
1115 .filter(|t| {
1116 devices
1117 .iter()
1118 .any(|(device, username)| *device == t.device && *username == t.username)
1119 })
1120 .collect();
1121 if targets.is_empty() {
1122 return;
1123 }
1124 std::thread::spawn(move || {
1125 for t in targets {
1126 deliver_push(pusher, &t, &crate::push::Push::Wake);
1127 }
1128 });
1129}
1130
1131fn reach_absent(
1140 username: Option<&str>,
1141 id: Option<&str>,
1142 cmd: &LinkCommand,
1143) -> Option<Result<ClientInfo, String>> {
1144 let pusher = crate::push::pusher()?;
1145 let target = outbox::push_targets(username)
1148 .into_iter()
1149 .find(|t| id.is_none_or(|id| t.device == id || t.name.eq_ignore_ascii_case(id)))?;
1150 let info = ClientInfo {
1151 id: target.device.clone(),
1152 device: target.device.clone(),
1153 name: target.name.clone(),
1154 platform: target.platform.clone(),
1155 username: target.username.clone(),
1156 connected_at: 0,
1157 state: LinkState::default(),
1158 last_played_at: None,
1159 state_at: 0,
1160 reports: false,
1161 notified: true,
1162 };
1163 let asked = match cmd {
1165 LinkCommand::Shared { command } => command.as_ref(),
1166 cmd => cmd,
1167 };
1168 let verb = match asked {
1169 LinkCommand::Play { .. } | LinkCommand::JumpTo { .. } | LinkCommand::Resume => Some("Play"),
1170 _ => None,
1171 };
1172 let push = match verb {
1173 Some(verb) => crate::push::Push::Notify {
1174 title: format!("{verb} on {}", target.name),
1175 body: outbox::describe(asked).unwrap_or_else(|| "From your koan server".into()),
1176 command: serde_json::to_value(cmd).ok()?,
1177 image: cover_track(asked)
1178 .and_then(outbox::track_row)
1179 .and_then(|t| pusher.cover_link(t)),
1180 },
1181 None => {
1182 outbox::queue_for(&target.device, &target.username, cmd);
1183 crate::push::Push::Wake
1184 }
1185 };
1186 std::thread::spawn(move || deliver_push(pusher, &target, &push));
1187 Some(Ok(info))
1188}
1189
1190fn cover_track(cmd: &LinkCommand) -> Option<&str> {
1193 match cmd {
1194 LinkCommand::Play {
1195 track_ids,
1196 start_at,
1197 ..
1198 } => track_ids
1199 .get(*start_at as usize)
1200 .or(track_ids.first())
1201 .map(String::as_str),
1202 LinkCommand::JumpTo { track_id } => Some(track_id),
1203 _ => None,
1204 }
1205}
1206
1207fn deliver_push(
1209 pusher: &crate::push::Pusher,
1210 target: &outbox::PushTarget,
1211 push: &crate::push::Push,
1212) {
1213 use crate::push::Outcome;
1214 match pusher.send(&target.token, target.sandbox, push) {
1215 Outcome::Sent => log::info!("push: sent to {}", target.name),
1216 Outcome::Gone => {
1217 log::info!(
1218 "push: {}'s token is no longer valid; forgotten",
1219 target.name
1220 );
1221 outbox::forget_push(&target.username, &target.device);
1222 }
1223 Outcome::Failed(e) => log::warn!("push: to {} failed: {e}", target.name),
1224 }
1225}
1226
1227type Woken = std::collections::HashMap<(String, String), (std::time::Instant, &'static str)>;
1231
1232static WOKEN: LazyLock<Mutex<Woken>> = LazyLock::new(Default::default);
1233
1234pub fn smart_activity(
1238 db: &koan_core::db::connection::Database,
1239 user: i64,
1240 fields: &[koan_core::smart::Field],
1241) {
1242 match koan_core::db::queries::smart::refresh_after_activity(&db.conn, user, fields) {
1243 Ok(moved) if !moved.is_empty() => changed(),
1244 Ok(_) => {}
1245 Err(e) => log::warn!("smart playlists not refreshed after activity: {e}"),
1246 }
1247}
1248
1249pub fn changed() {
1254 let (_, queued) = registry().link_or_queue(None, &LinkCommand::Sync { full: false });
1255 if queued.is_empty() || crate::push::pusher().is_none() {
1256 return;
1257 }
1258 let mut wakes = WAKES.lock();
1259 wakes.add(queued.iter().map(Absent::key), std::time::Instant::now());
1260 if !wakes.timer {
1261 wakes.timer = true;
1262 std::thread::spawn(send_wakes);
1263 }
1264}
1265
1266const QUIET: std::time::Duration = std::time::Duration::from_secs(30);
1270
1271const LONGEST_WAIT: std::time::Duration = std::time::Duration::from_secs(5 * 60);
1273
1274static WAKES: Mutex<Wakes> = Mutex::new(Wakes {
1275 pending: Vec::new(),
1276 timer: false,
1277});
1278
1279struct Wakes {
1282 pending: Vec<((String, String), std::time::Instant, std::time::Instant)>,
1284 timer: bool,
1286}
1287
1288impl Wakes {
1289 fn add(
1290 &mut self,
1291 devices: impl IntoIterator<Item = (String, String)>,
1292 now: std::time::Instant,
1293 ) {
1294 for key in devices {
1295 match self.pending.iter_mut().find(|(k, _, _)| *k == key) {
1296 Some((_, _, last)) => *last = now,
1297 None => self.pending.push((key, now, now)),
1298 }
1299 }
1300 }
1301
1302 fn due_at(first: std::time::Instant, last: std::time::Instant) -> std::time::Instant {
1303 (last + QUIET).min(first + LONGEST_WAIT)
1304 }
1305
1306 fn next(&self) -> Option<std::time::Instant> {
1308 self.pending
1309 .iter()
1310 .map(|(_, first, last)| Self::due_at(*first, *last))
1311 .min()
1312 }
1313
1314 fn take_due(&mut self, now: std::time::Instant) -> Vec<(String, String)> {
1316 let (due, waiting) = std::mem::take(&mut self.pending)
1317 .into_iter()
1318 .partition(|(_, first, last)| Self::due_at(*first, *last) <= now);
1319 self.pending = waiting;
1320 due.into_iter().map(|(key, _, _)| key).collect()
1321 }
1322}
1323
1324fn send_wakes() {
1327 loop {
1328 let (due, next) = {
1329 let mut wakes = WAKES.lock();
1330 let due = wakes.take_due(std::time::Instant::now());
1331 let next = wakes.next();
1332 if due.is_empty() && next.is_none() {
1333 wakes.timer = false;
1334 return;
1335 }
1336 (due, next)
1337 };
1338 let live = registry().live();
1339 let absent: Vec<(String, String)> = due.into_iter().filter(|d| !live.contains(d)).collect();
1340 if !absent.is_empty() {
1341 wake(&absent);
1342 }
1343 if let Some(next) = next {
1344 std::thread::sleep(next.saturating_duration_since(std::time::Instant::now()));
1345 }
1346 }
1347}
1348
1349pub fn changed_if_library_moved(conn: &rusqlite::Connection) {
1353 static LAST: parking_lot::Mutex<Option<(i64, i64, i64)>> = parking_lot::Mutex::new(None);
1354 let Ok(now) = conn.query_row(
1355 "SELECT (SELECT COUNT(*) FROM tracks), (SELECT COALESCE(MAX(id), 0) FROM tracks),
1356 (SELECT COUNT(*) FROM albums)",
1357 [],
1358 |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
1359 ) else {
1360 return;
1361 };
1362 let before = LAST.lock().replace(now);
1363 if before.is_some_and(|b| b != now) {
1364 changed();
1365 }
1366}
1367
1368mod outbox {
1371 use koan_core::db::queries;
1372 use koan_core::remote::link::LinkCommand;
1373
1374 const KEEP_SECS: i64 = 30 * 24 * 60 * 60;
1379
1380 fn db() -> Option<koan_core::db::pool::Handle<'static>> {
1385 if cfg!(test) {
1386 return None;
1387 }
1388 koan_core::db::pool::shared().get().ok()
1389 }
1390
1391 pub fn load_orders() -> Vec<super::Order> {
1392 let Some(db) = db() else { return Vec::new() };
1393 db.conn
1394 .prepare("SELECT body FROM link_orders ORDER BY created_at")
1395 .and_then(|mut s| {
1396 s.query_map([], |r| r.get::<_, String>(0))?
1397 .collect::<Result<Vec<_>, _>>()
1398 })
1399 .unwrap_or_default()
1400 .into_iter()
1401 .filter_map(|b| serde_json::from_str(&b).ok())
1402 .collect()
1403 }
1404
1405 pub fn save_order(order: &super::Order) {
1406 let (Some(db), Ok(body)) = (db(), serde_json::to_string(order)) else {
1407 return;
1408 };
1409 let _ = db.conn.execute(
1410 "INSERT OR REPLACE INTO link_orders (id, body, created_at) VALUES (?1, ?2, ?3)",
1411 rusqlite::params![order.id, body, order.created_at],
1412 );
1413 }
1414
1415 pub fn drop_order(id: &str) {
1416 if let Some(db) = db() {
1417 let _ = db
1418 .conn
1419 .execute("DELETE FROM link_orders WHERE id = ?1", [id]);
1420 }
1421 }
1422
1423 pub fn take_and_remember(
1424 username: &str,
1425 device: &str,
1426 name: &str,
1427 platform: &str,
1428 live: &[(String, String)],
1429 ) -> Vec<LinkCommand> {
1430 let Some(db) = db() else { return Vec::new() };
1431 let now = chrono::Utc::now().timestamp();
1432 let _ = db.conn.execute(
1433 "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES (?1, ?2, ?3, ?4, ?5)
1434 ON CONFLICT (device, username) DO UPDATE SET name = ?3, platform = ?4, last_seen = ?5",
1435 rusqlite::params![device, username, name, platform, now],
1436 );
1437 forget_stale(&db.conn, live, now);
1438 let waiting: Vec<(i64, String)> = db
1439 .conn
1440 .prepare("SELECT id, command FROM link_outbox WHERE device = ?1 AND username = ?2 ORDER BY id")
1441 .and_then(|mut s| {
1442 s.query_map([device, username], |r| Ok((r.get(0)?, r.get(1)?)))?
1443 .collect()
1444 })
1445 .unwrap_or_default();
1446 let _ = db.conn.execute(
1447 "DELETE FROM link_outbox WHERE device = ?1 AND username = ?2",
1448 [device, username],
1449 );
1450 if !waiting.is_empty() {
1451 log::info!("link: {} waiting commands for {name}", waiting.len());
1452 }
1453 waiting
1454 .into_iter()
1455 .filter_map(|(_, c)| serde_json::from_str(&c).ok())
1456 .collect()
1457 }
1458
1459 pub(super) fn forget_stale(conn: &rusqlite::Connection, live: &[(String, String)], now: i64) {
1462 touch_with(conn, live, now);
1463 let _ = conn.execute(
1464 "DELETE FROM link_outbox WHERE created_at < ?1",
1465 [now - KEEP_SECS],
1466 );
1467 let _ = conn.execute(
1468 "DELETE FROM link_push WHERE (device, username) IN
1469 (SELECT device, username FROM link_devices WHERE last_seen < ?1)",
1470 [now - KEEP_SECS],
1471 );
1472 let _ = conn.execute(
1473 "DELETE FROM link_devices WHERE last_seen < ?1",
1474 [now - KEEP_SECS],
1475 );
1476 }
1477
1478 pub fn touch(devices: &[(String, String)]) {
1480 if let Some(db) = db() {
1481 touch_with(&db.conn, devices, chrono::Utc::now().timestamp());
1482 }
1483 }
1484
1485 fn touch_with(conn: &rusqlite::Connection, devices: &[(String, String)], now: i64) {
1486 for (device, username) in devices {
1487 let _ = conn.execute(
1488 "UPDATE link_devices SET last_seen = ?1 WHERE device = ?2 AND username = ?3",
1489 rusqlite::params![now, device, username],
1490 );
1491 }
1492 }
1493
1494 pub struct Absent {
1496 pub device: String,
1497 pub username: String,
1498 pub name: String,
1499 }
1500
1501 pub fn queue_for_absent(
1503 username: Option<&str>,
1504 live: &[(String, String)],
1505 cmd: &LinkCommand,
1506 ) -> Vec<Absent> {
1507 let Some(db) = db() else { return Vec::new() };
1508 let known: Vec<(String, String, String)> = db
1509 .conn
1510 .prepare("SELECT device, username, name FROM link_devices")
1511 .and_then(|mut s| {
1512 s.query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))?
1513 .collect()
1514 })
1515 .unwrap_or_default();
1516 let Ok(text) = serde_json::to_string(cmd) else {
1517 return Vec::new();
1518 };
1519 let absent: Vec<_> = known
1520 .into_iter()
1521 .filter(|(device, user, _)| {
1522 !username.is_some_and(|u| u != user)
1523 && !live.iter().any(|(d, u)| d == device && u == user)
1524 })
1525 .collect();
1526 if absent.is_empty() {
1527 return Vec::new();
1528 }
1529 let is_sync = matches!(cmd, LinkCommand::Sync { .. });
1530 let now = chrono::Utc::now().timestamp();
1531 koan_core::db::queries::atomically(&db.conn, || {
1534 let mut queued = Vec::new();
1535 for (device, user, name) in absent {
1536 if is_sync {
1537 let _ = db.conn.execute(
1539 "DELETE FROM link_outbox WHERE device = ?1 AND username = ?2 AND command LIKE '{\"type\":\"sync\"%'",
1540 [&device, &user],
1541 );
1542 }
1543 if db
1544 .conn
1545 .execute(
1546 "INSERT INTO link_outbox (device, username, command, created_at) VALUES (?1, ?2, ?3, ?4)",
1547 rusqlite::params![device, user, text, now],
1548 )
1549 .is_ok()
1550 {
1551 queued.push(Absent {
1552 device,
1553 username: user,
1554 name,
1555 });
1556 }
1557 }
1558 Ok::<_, rusqlite::Error>(queued)
1559 })
1560 .unwrap_or_default()
1561 }
1562
1563 #[derive(Clone)]
1565 pub struct PushTarget {
1566 pub device: String,
1567 pub username: String,
1568 pub name: String,
1569 pub platform: String,
1570 pub token: String,
1571 pub sandbox: bool,
1572 pub last_seen: i64,
1574 }
1575
1576 pub fn queue_for(device: &str, username: &str, cmd: &LinkCommand) {
1578 let (Some(db), Ok(text)) = (db(), serde_json::to_string(cmd)) else {
1579 return;
1580 };
1581 let _ = db.conn.execute(
1582 "INSERT INTO link_outbox (device, username, command, created_at) VALUES (?1, ?2, ?3, ?4)",
1583 rusqlite::params![device, username, text, chrono::Utc::now().timestamp()],
1584 );
1585 }
1586
1587 pub fn accounts() -> Vec<String> {
1591 let Some(db) = db() else { return Vec::new() };
1592 koan_core::db::queries::auth::list_users(&db.conn)
1593 .map(|users| users.into_iter().map(|u| u.username).collect())
1594 .unwrap_or_default()
1595 }
1596
1597 pub fn load_addresses() -> std::collections::HashMap<String, (String, std::net::IpAddr)> {
1600 let Some(db) = db() else {
1601 return Default::default();
1602 };
1603 load_addresses_in(&db.conn)
1604 }
1605
1606 pub(super) fn load_addresses_in(
1607 conn: &rusqlite::Connection,
1608 ) -> std::collections::HashMap<String, (String, std::net::IpAddr)> {
1609 conn.prepare(
1610 "SELECT device, username, addr FROM link_devices WHERE addr IS NOT NULL ORDER BY last_seen",
1611 )
1612 .and_then(|mut s| {
1613 s.query_map([], |r| {
1614 Ok((
1615 r.get::<_, String>(0)?,
1616 r.get::<_, String>(1)?,
1617 r.get::<_, String>(2)?,
1618 ))
1619 })?
1620 .collect::<Result<Vec<_>, _>>()
1621 })
1622 .unwrap_or_default()
1623 .into_iter()
1624 .filter_map(|(device, username, addr)| Some((device, (username, addr.parse().ok()?))))
1625 .collect()
1626 }
1627
1628 pub fn save_address(device: &str, username: &str, addr: std::net::IpAddr) {
1629 if let Some(db) = db() {
1630 save_address_in(&db.conn, device, username, addr);
1631 }
1632 }
1633
1634 pub(super) fn save_address_in(
1635 conn: &rusqlite::Connection,
1636 device: &str,
1637 username: &str,
1638 addr: std::net::IpAddr,
1639 ) {
1640 let _ = conn.execute(
1641 "UPDATE link_devices SET addr = ?1 WHERE device = ?2 AND username = ?3",
1642 rusqlite::params![addr.to_string(), device, username],
1643 );
1644 }
1645
1646 pub fn load_grants() -> Vec<super::Grant> {
1647 let Some(db) = db() else { return Vec::new() };
1648 db.conn
1649 .prepare("SELECT device, owner, grantee FROM link_grants ORDER BY created_at")
1650 .and_then(|mut s| {
1651 s.query_map([], |r| {
1652 Ok(super::Grant {
1653 device: r.get(0)?,
1654 owner: r.get(1)?,
1655 grantee: r.get(2)?,
1656 })
1657 })?
1658 .collect()
1659 })
1660 .unwrap_or_default()
1661 }
1662
1663 pub fn save_grant(g: &super::Grant) {
1664 if let Some(db) = db() {
1665 let _ = db.conn.execute(
1666 "INSERT OR IGNORE INTO link_grants (device, owner, grantee, created_at) VALUES (?1, ?2, ?3, ?4)",
1667 rusqlite::params![g.device, g.owner, g.grantee, chrono::Utc::now().timestamp()],
1668 );
1669 }
1670 }
1671
1672 pub fn drop_grant(g: &super::Grant) {
1673 if let Some(db) = db() {
1674 let _ = db.conn.execute(
1675 "DELETE FROM link_grants WHERE device = ?1 AND owner = ?2 AND grantee = ?3",
1676 [&g.device, &g.owner, &g.grantee],
1677 );
1678 }
1679 }
1680
1681 pub fn save_push(username: &str, device: &str, token: &str, sandbox: bool) {
1682 let Some(db) = db() else { return };
1683 let _ = db.conn.execute(
1684 "INSERT INTO link_push (device, username, token, sandbox, updated_at) VALUES (?1, ?2, ?3, ?4, ?5)
1685 ON CONFLICT (device, username) DO UPDATE SET token = ?3, sandbox = ?4, updated_at = ?5",
1686 rusqlite::params![device, username, token, sandbox, chrono::Utc::now().timestamp()],
1687 );
1688 }
1689
1690 pub fn forget_push(username: &str, device: &str) {
1691 if let Some(db) = db() {
1692 let _ = db.conn.execute(
1693 "DELETE FROM link_push WHERE device = ?1 AND username = ?2",
1694 [device, username],
1695 );
1696 }
1697 }
1698
1699 pub fn push_targets(username: Option<&str>) -> Vec<PushTarget> {
1701 let Some(db) = db() else { return Vec::new() };
1702 push_targets_in(&db.conn, username)
1703 }
1704
1705 pub(super) fn push_targets_in(
1706 conn: &rusqlite::Connection,
1707 username: Option<&str>,
1708 ) -> Vec<PushTarget> {
1709 conn
1710 .prepare(
1711 "SELECT p.device, p.username, d.name, d.platform, p.token, p.sandbox, d.last_seen
1712 FROM link_push p JOIN link_devices d ON d.device = p.device AND d.username = p.username
1713 WHERE ?1 IS NULL OR p.username = ?1
1714 ORDER BY d.last_seen DESC",
1715 )
1716 .and_then(|mut s| {
1717 s.query_map([username], |r| {
1718 Ok(PushTarget {
1719 device: r.get(0)?,
1720 username: r.get(1)?,
1721 name: r.get(2)?,
1722 platform: r.get(3)?,
1723 token: r.get(4)?,
1724 sandbox: r.get(5)?,
1725 last_seen: r.get(6)?,
1726 })
1727 })?
1728 .collect()
1729 })
1730 .unwrap_or_default()
1731 }
1732
1733 pub fn forget_device(device: &str, username: &str) {
1736 if let Some(db) = db() {
1737 forget_device_in(&db.conn, device, username);
1738 }
1739 }
1740
1741 pub(super) fn forget_device_in(conn: &rusqlite::Connection, device: &str, username: &str) {
1742 for table in ["link_push", "link_outbox", "link_devices"] {
1743 let _ = conn.execute(
1744 &format!("DELETE FROM {table} WHERE device = ?1 AND username = ?2"),
1745 [device, username],
1746 );
1747 }
1748 }
1749
1750 pub fn track_row(id: &str) -> Option<i64> {
1752 let db = db()?;
1753 queries::resolve_id(&db.conn, queries::UidKind::Track, id)
1754 .ok()
1755 .flatten()
1756 }
1757
1758 pub fn describe(cmd: &LinkCommand) -> Option<String> {
1761 let ids: Vec<&String> = match cmd {
1762 LinkCommand::Play { track_ids, .. }
1763 | LinkCommand::Enqueue { track_ids }
1764 | LinkCommand::PlayNext { track_ids } => track_ids.iter().collect(),
1765 LinkCommand::JumpTo { track_id } => vec![track_id],
1766 _ => return None,
1767 };
1768 let db = db()?;
1769 let ids: Vec<i64> = ids
1770 .into_iter()
1771 .filter_map(|t| {
1772 queries::resolve_id(&db.conn, queries::UidKind::Track, t)
1773 .ok()
1774 .flatten()
1775 })
1776 .collect();
1777 let row = |id: i64| {
1778 db.conn
1779 .query_row(
1780 "SELECT t.title, COALESCE(a.name, ''), COALESCE(al.title, ''), t.album_id
1781 FROM tracks t LEFT JOIN artists a ON a.id = t.artist_id
1782 LEFT JOIN albums al ON al.id = t.album_id WHERE t.id = ?1",
1783 [id],
1784 |r| {
1785 Ok((
1786 r.get::<_, String>(0)?,
1787 r.get::<_, String>(1)?,
1788 r.get::<_, String>(2)?,
1789 r.get::<_, Option<i64>>(3)?,
1790 ))
1791 },
1792 )
1793 .ok()
1794 };
1795 let (title, artist, album, album_id) = row(*ids.first()?)?;
1796 let one_album = ids.len() > 1
1797 && album_id.is_some()
1798 && ids
1799 .iter()
1800 .all(|id| row(*id).is_some_and(|r| r.3 == album_id));
1801 Some(match (one_album, ids.len()) {
1802 (true, _) => format!("{album} — {artist}"),
1803 (false, 1) => format!("{title} — {artist}"),
1804 (false, n) => format!("{title} — {artist}, and {} more", n - 1),
1805 })
1806 }
1807}
1808
1809const RECENT: i64 = 6 * 60 * 60;
1812
1813fn pick(clients: &[ClientInfo], now: i64) -> Result<&ClientInfo, String> {
1814 if let Some(c) = clients.iter().find(|c| c.state.playing) {
1815 return Ok(c);
1816 }
1817 if let Some(c) = clients
1818 .iter()
1819 .filter(|c| c.last_played_at.is_some_and(|t| now - t < RECENT))
1820 .max_by_key(|c| c.last_played_at)
1821 {
1822 return Ok(c);
1823 }
1824 match clients {
1825 [] => Err("no koan app is linked to this server; open koan on the device".into()),
1826 [only] => Ok(only),
1827 several => Err(format!(
1828 "several koan apps are linked and none has played recently: {}. Ask which, then pass `client`",
1829 several
1830 .iter()
1831 .map(|c| format!("{} ({}, id {})", c.name, c.platform, c.device))
1832 .collect::<Vec<_>>()
1833 .join(", ")
1834 )),
1835 }
1836}
1837
1838#[cfg(test)]
1839mod tests {
1840
1841 #[test]
1842 fn a_watch_is_never_queued_and_is_renewed_when_the_target_relinks() {
1843 let reg = Registry::default();
1844 let watches = |rx: &mut tokio::sync::mpsc::UnboundedReceiver<LinkCommand>| {
1845 std::iter::from_fn(|| rx.try_recv().ok())
1846 .filter(|c| matches!(c, LinkCommand::WatchLevels { .. }))
1847 .collect::<Vec<_>>()
1848 };
1849 assert!(
1851 reg.relay("rl", "rl-phone", LinkCommand::WatchLevels { on: true })
1852 .is_err()
1853 );
1854 reg.watch_levels("rl", "rl-mac", "rl-phone", true);
1856 let (tx, mut phone) = tokio::sync::mpsc::unbounded_channel();
1857 reg.register("rl", "phone", "ios", "rl-phone", tx, false);
1858 assert_eq!(
1859 watches(&mut phone),
1860 vec![LinkCommand::WatchLevels { on: true }],
1861 "told once, on linking, because it is watched now"
1862 );
1863 let (tx, mut phone) = tokio::sync::mpsc::unbounded_channel();
1865 reg.register("rl", "phone", "ios", "rl-phone", tx, false);
1866 assert_eq!(
1867 watches(&mut phone),
1868 vec![LinkCommand::WatchLevels { on: true }]
1869 );
1870 }
1871
1872 #[test]
1873 fn levels_reach_a_watcher_only_while_it_watches() {
1874 use koan_core::remote::levels::Frame;
1875 let reg = Registry::default();
1876 let (tx_mac, mut mac) = tokio::sync::mpsc::unbounded_channel();
1877 let (tx_phone, mut phone) = tokio::sync::mpsc::unbounded_channel();
1878 let mac_id = reg.register("lv", "mac", "macos", "lv-mac", tx_mac, false);
1879 reg.register("lv", "phone", "ios", "lv-phone", tx_phone, false);
1880 let levels = |rx: &mut tokio::sync::mpsc::UnboundedReceiver<LinkCommand>| {
1881 std::iter::from_fn(|| rx.try_recv().ok())
1882 .filter(|c| {
1883 matches!(
1884 c,
1885 LinkCommand::Levels { .. } | LinkCommand::WatchLevels { .. }
1886 )
1887 })
1888 .collect::<Vec<_>>()
1889 };
1890 let f = Frame(1_000, 1, 2, 3);
1891
1892 reg.levels("lv", "lv-phone", f);
1893 assert!(levels(&mut mac).is_empty(), "nobody watching");
1894
1895 reg.watch_levels("lv", "lv-mac", "lv-phone", true);
1896 assert_eq!(
1897 levels(&mut phone),
1898 vec![LinkCommand::WatchLevels { on: true }]
1899 );
1900 reg.levels("lv", "lv-phone", f);
1901 assert_eq!(
1902 levels(&mut mac),
1903 vec![LinkCommand::Levels {
1904 from: "lv-phone".into(),
1905 f
1906 }]
1907 );
1908
1909 reg.unregister(&mac_id);
1911 assert_eq!(
1912 levels(&mut phone),
1913 vec![LinkCommand::WatchLevels { on: false }]
1914 );
1915 reg.watch_levels("lv", "lv-mac", "lv-phone", true);
1916 reg.watch_levels("lv", "lv-mac", "lv-phone", false);
1917 assert_eq!(
1918 levels(&mut phone),
1919 vec![
1920 LinkCommand::WatchLevels { on: true },
1921 LinkCommand::WatchLevels { on: false }
1922 ]
1923 );
1924 }
1925
1926 #[test]
1927 fn a_playlist_order_adds_the_named_tracks_once_they_arrive() {
1928 let dir = tempfile::tempdir().unwrap();
1929 let path = dir.path().join("koan.db");
1930 let db = koan_core::db::connection::Database::open(&path).unwrap();
1931 let playlist = koan_core::db::queries::create_playlist(
1932 &db.conn,
1933 koan_core::db::queries::LOCAL_USER,
1934 "cyberpunk",
1935 None,
1936 )
1937 .unwrap();
1938 let order = Order {
1939 id: "o1".into(),
1940 username: None,
1941 client: None,
1942 artist: "Perturbator".into(),
1943 album: "Dangerous Days".into(),
1944 play_next: false,
1945 playlist: Some(playlist),
1946 titles: vec!["Future Club".into()],
1947 created_at: chrono::Utc::now().timestamp(),
1948 };
1949 registry().add_order(order);
1950
1951 fulfil_from(&path);
1953 assert!(registry().orders(None).iter().any(|o| o.id == "o1"));
1954
1955 db.conn
1956 .execute_batch(
1957 "INSERT INTO artists (id, name) VALUES (1, 'Perturbator');
1958 INSERT INTO albums (id, title, artist_id) VALUES (1, 'Dangerous Days', 1);
1959 INSERT INTO tracks (id, title, album_id, artist_id, track_number, path) VALUES
1960 (1, 'Welcome Back', 1, 1, 1, '/1.flac'), (2, 'Future Club', 1, 1, 2, '/2.flac');",
1961 )
1962 .unwrap();
1963 fulfil_from(&path);
1964 let held: Vec<i64> = db
1965 .conn
1966 .prepare("SELECT track_id FROM playlist_tracks WHERE playlist_id = ?1")
1967 .unwrap()
1968 .query_map([playlist], |r| r.get(0))
1969 .unwrap()
1970 .collect::<Result<_, _>>()
1971 .unwrap();
1972 assert_eq!(held, [2]);
1973 assert!(!registry().orders(None).iter().any(|o| o.id == "o1"));
1974 }
1975
1976 use super::*;
1977
1978 #[test]
1979 fn a_notification_cover_is_named_by_the_uid_the_command_carries() {
1980 let uid = "0190a5b2-7c3d-7e4f-8a1b-2c3d4e5f6a7b".to_string();
1982 let play = LinkCommand::Play {
1983 track_ids: vec!["x".into(), uid.clone()],
1984 start_at: 1,
1985 position_ms: 0,
1986 paused: false,
1987 handoff: false,
1988 };
1989 assert_eq!(cover_track(&play), Some(uid.as_str()));
1990 }
1991
1992 #[test]
1993 fn a_reconnect_replaces_the_device_and_commands_reach_it() {
1994 let reg = Registry::default();
1995 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
1996 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
1997 let (tx3, _rx3) = tokio::sync::mpsc::unbounded_channel();
1998 reg.register("j", "phone", "ios", "dev-1", tx1, false);
1999 let id = reg.register("j", "phone", "ios", "dev-1", tx2, false);
2000 reg.register("someone", "laptop", "macos", "dev-2", tx3, false);
2001
2002 assert_eq!(reg.list(Some("j")).len(), 1);
2003 assert_eq!(reg.list(None).len(), 2);
2004
2005 let sent = reg.send(Some("j"), None, LinkCommand::Pause).unwrap();
2006 assert_eq!(sent.id, id);
2007 assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
2008
2009 assert!(
2011 reg.send(Some("j"), Some("laptop"), LinkCommand::Pause)
2012 .is_err()
2013 );
2014
2015 reg.unregister(&id);
2016 assert!(reg.send(Some("j"), None, LinkCommand::Pause).is_err());
2017 }
2018
2019 #[test]
2020 fn disconnecting_an_account_closes_only_its_links() {
2021 let reg = Registry::default();
2022 let (tx1, mut rx1) = tokio::sync::mpsc::unbounded_channel();
2023 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
2024 reg.register("j", "phone", "ios", "dev-1", tx1, false);
2025 reg.register("someone", "laptop", "macos", "dev-2", tx2, false);
2026
2027 reg.disconnect("j");
2028 assert!(reg.list(Some("j")).is_empty());
2029 assert!(matches!(
2031 rx1.try_recv(),
2032 Err(tokio::sync::mpsc::error::TryRecvError::Disconnected)
2033 ));
2034 assert_eq!(reg.list(Some("someone")).len(), 1);
2035 assert!(matches!(
2036 rx2.try_recv(),
2037 Err(tokio::sync::mpsc::error::TryRecvError::Empty)
2038 ));
2039 }
2040
2041 #[test]
2042 fn a_burst_of_changes_wakes_each_device_once_when_it_goes_quiet() {
2043 let t0 = std::time::Instant::now();
2044 let s = std::time::Duration::from_secs;
2045 let phone = || ("dev-1".to_string(), "j".to_string());
2046 let ipad = || ("dev-2".to_string(), "j".to_string());
2047 let mut wakes = Wakes {
2048 pending: Vec::new(),
2049 timer: false,
2050 };
2051
2052 wakes.add([phone()], t0);
2053 wakes.add([phone(), ipad()], t0 + s(10));
2054 wakes.add([phone()], t0 + s(20));
2055 assert_eq!(wakes.pending.len(), 2);
2056
2057 assert!(wakes.take_due(t0 + s(39)).is_empty());
2059 assert_eq!(wakes.take_due(t0 + s(40)), [ipad()]);
2060 assert_eq!(wakes.next(), Some(t0 + s(50)));
2061 assert_eq!(wakes.take_due(t0 + s(50)), [phone()]);
2062 assert_eq!(wakes.next(), None);
2063 }
2064
2065 #[test]
2066 fn a_library_that_never_goes_quiet_still_wakes_devices() {
2067 let t0 = std::time::Instant::now();
2068 let phone = || ("dev-1".to_string(), "j".to_string());
2069 let mut wakes = Wakes {
2070 pending: Vec::new(),
2071 timer: false,
2072 };
2073 let mut sent = 0;
2074 for i in 0..40 {
2075 let now = t0 + std::time::Duration::from_secs(i * 10);
2076 sent += wakes.take_due(now).len();
2077 wakes.add([phone()], now);
2078 }
2079 assert_eq!(sent, 1);
2082 assert_eq!(wakes.pending.len(), 1);
2083 }
2084
2085 #[test]
2086 fn the_device_playing_is_the_one_meant() {
2087 let reg = Registry::default();
2088 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
2089 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
2090 let mac = reg.register("j", "mac", "macos", "dev-1", tx1, false);
2091 let phone = reg.register("j", "phone", "ios", "dev-2", tx2, false);
2092
2093 let err = reg.send(Some("j"), None, LinkCommand::Pause).unwrap_err();
2095 assert!(err.contains("mac") && err.contains("phone"), "{err}");
2096
2097 reg.report(
2098 &phone,
2099 LinkState {
2100 playing: true,
2101 ..Default::default()
2102 },
2103 );
2104 assert_eq!(
2105 reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
2106 phone
2107 );
2108 assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
2109
2110 reg.report(&phone, LinkState::default());
2112 assert_eq!(
2113 reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
2114 phone
2115 );
2116 let _ = mac;
2117 }
2118
2119 fn drain(rx: &mut tokio::sync::mpsc::UnboundedReceiver<LinkCommand>) -> Vec<LinkCommand> {
2120 std::iter::from_fn(|| rx.try_recv().ok()).collect()
2121 }
2122
2123 fn shared() -> (
2125 Registry,
2126 tokio::sync::mpsc::UnboundedReceiver<LinkCommand>,
2127 tokio::sync::mpsc::UnboundedReceiver<LinkCommand>,
2128 tokio::sync::mpsc::UnboundedReceiver<LinkCommand>,
2129 ) {
2130 let reg = Registry::default();
2131 let (phone_tx, mut phone) = tokio::sync::mpsc::unbounded_channel();
2132 let (mac_tx, mac) = tokio::sync::mpsc::unbounded_channel();
2133 let (k_tx, k) = tokio::sync::mpsc::unbounded_channel();
2134 let id = reg.register("j", "phone", "ios", "dev-phone", phone_tx, true);
2135 reg.register("j", "mac", "macos", "dev-mac", mac_tx, true);
2136 reg.register("k", "laptop", "macos", "dev-k", k_tx, true);
2137 reg.report(
2138 &id,
2139 LinkState {
2140 playing: true,
2141 title: Some("Roygbiv".into()),
2142 outputs: Some(Default::default()),
2143 ..Default::default()
2144 },
2145 );
2146 reg.share("j", "dev-phone", "k", true).unwrap();
2147 drain(&mut phone);
2148 (reg, phone, mac, k)
2149 }
2150
2151 #[test]
2152 fn a_grantee_sees_the_shared_device_and_nothing_else_of_the_owner() {
2153 let (_reg, _phone, _mac, mut k) = shared();
2154 let listed = drain(&mut k)
2155 .into_iter()
2156 .rev()
2157 .find_map(|c| match c {
2158 LinkCommand::Devices { devices } => Some(devices),
2159 _ => None,
2160 })
2161 .expect("told of it");
2162 assert_eq!(listed.len(), 1, "the phone, not the Mac");
2163 let phone = &listed[0];
2164 assert_eq!(phone.id, "dev-phone");
2165 assert_eq!(phone.owner.as_deref(), Some("j"));
2166 let state = phone.state.as_ref().expect("its state");
2167 assert_eq!(state.title.as_deref(), Some("Roygbiv"));
2168 assert!(
2169 state.outputs.is_some(),
2170 "outputs, to choose from: the output is in the playback set"
2171 );
2172 }
2173
2174 #[test]
2180 fn a_grantee_controls_playback_as_itself_and_nothing_of_the_owners_account() {
2181 let (reg, mut phone, mut mac, _k) = shared();
2182 let wrapped = |c: LinkCommand| LinkCommand::Shared {
2183 command: Box::new(c),
2184 };
2185 for cmd in [
2186 LinkCommand::Pause,
2187 LinkCommand::SetRendererVolume { volume: 40 },
2188 LinkCommand::SleepTimer {
2189 timer: Some(koan_core::player::state::SleepTimer::After { minutes: 30 }),
2190 },
2191 LinkCommand::HandOff { to: "dev-k".into() },
2192 ] {
2193 reg.relay("k", "dev-phone", cmd.clone()).unwrap();
2194 assert_eq!(drain(&mut phone), [wrapped(cmd)]);
2195 }
2196 for cmd in [
2197 LinkCommand::Sync { full: false },
2198 LinkCommand::Evict { track_ids: vec![] },
2199 LinkCommand::Shared {
2200 command: Box::new(LinkCommand::Pause),
2201 },
2202 ] {
2203 assert!(!cmd.allowed_playback());
2204 assert!(reg.relay("k", "dev-phone", cmd).is_err());
2205 }
2206 assert!(drain(&mut phone).is_empty());
2207 assert!(
2208 reg.relay("k", "dev-mac", LinkCommand::Pause).is_err(),
2209 "not shared, not reachable"
2210 );
2211 assert!(
2212 drain(&mut mac)
2213 .iter()
2214 .all(|c| matches!(c, LinkCommand::Devices { .. } | LinkCommand::Shares { .. })),
2215 "news, and no command"
2216 );
2217 assert!(reg.shared_owner("k", "dev-mac").is_none(), "nor wakeable");
2218 }
2219
2220 #[test]
2224 fn a_hand_off_runs_both_ways_across_a_grant_and_nothing_else_does() {
2225 let (reg, _phone, _mac, mut k) = shared();
2226 drain(&mut k);
2227 let play = LinkCommand::Play {
2228 track_ids: vec!["t".into()],
2229 start_at: 0,
2230 position_ms: 1000,
2231 paused: false,
2232 handoff: true,
2233 };
2234 reg.relay_from("j", Some("dev-phone"), "dev-k", play.clone())
2235 .unwrap();
2236 assert!(drain(&mut k).contains(&LinkCommand::Shared {
2237 command: Box::new(play.clone())
2238 }));
2239 assert!(
2240 reg.relay_from("j", Some("dev-phone"), "dev-k", LinkCommand::Pause)
2241 .is_err(),
2242 "the owner's phone does not command the grantee's devices"
2243 );
2244 assert!(
2245 reg.relay_from("j", Some("dev-mac"), "dev-k", play.clone())
2246 .is_err(),
2247 "only the shared device"
2248 );
2249 let not_a_hand_off = LinkCommand::Play {
2250 track_ids: vec!["t".into()],
2251 start_at: 0,
2252 position_ms: 0,
2253 paused: false,
2254 handoff: false,
2255 };
2256 assert!(
2257 reg.relay_from("j", Some("dev-phone"), "dev-k", not_a_hand_off)
2258 .is_err(),
2259 "a hand-off, not any play"
2260 );
2261 }
2262
2263 #[test]
2264 fn revoking_ends_control_at_once() {
2265 let (reg, mut phone, _mac, mut k) = shared();
2266 reg.share("j", "dev-phone", "k", false).unwrap();
2267 assert!(reg.relay("k", "dev-phone", LinkCommand::Pause).is_err());
2268 assert!(
2269 drain(&mut phone)
2270 .iter()
2271 .all(|c| !matches!(c, LinkCommand::Pause))
2272 );
2273 assert!(reg.shared_owner("k", "dev-phone").is_none());
2274 let listed = drain(&mut k)
2275 .into_iter()
2276 .rev()
2277 .find_map(|c| match c {
2278 LinkCommand::Devices { devices } => Some(devices),
2279 _ => None,
2280 })
2281 .expect("told it is gone");
2282 assert!(listed.is_empty());
2283 }
2284
2285 #[test]
2286 fn a_device_hears_whom_it_is_shared_with_every_time_it_links() {
2287 let reg = Registry::default();
2288 let (tx, mut phone) = tokio::sync::mpsc::unbounded_channel();
2289 reg.register("j", "phone", "ios", "dev-phone", tx, true);
2290 let first = drain(&mut phone);
2291 assert!(
2292 first.contains(&LinkCommand::Shares {
2293 grantees: vec![],
2294 error: None,
2295 accounts: vec![],
2296 }),
2297 "an empty list replaces one from another server"
2298 );
2299 reg.send_shares(
2300 "j",
2301 "dev-phone",
2302 Some("There is no account called x".into()),
2303 );
2304 assert!(drain(&mut phone).iter().any(
2305 |c| matches!(c, LinkCommand::Shares { error: Some(e), .. } if e.contains("no account"))
2306 ));
2307 }
2308
2309 #[test]
2310 fn the_owner_is_told_who_it_shares_with() {
2311 let reg = Registry::default();
2312 let (tx, mut phone) = tokio::sync::mpsc::unbounded_channel();
2313 reg.register("j", "phone", "ios", "dev-phone", tx, true);
2314 reg.share("j", "dev-phone", "k", true).unwrap();
2315 reg.share("j", "dev-phone", "m", true).unwrap();
2316 let last = drain(&mut phone).into_iter().rev().find_map(|c| match c {
2317 LinkCommand::Shares { grantees, .. } => Some(grantees),
2318 _ => None,
2319 });
2320 assert_eq!(last, Some(vec!["k".to_string(), "m".to_string()]));
2321 assert!(
2322 reg.share("j", "dev-phone", "j", true).is_err(),
2323 "not with itself"
2324 );
2325 }
2326
2327 #[test]
2331 fn a_device_is_woken_for_another_account_on_its_network_or_by_grant() {
2332 let reg = Registry::default();
2333 let home: std::net::IpAddr = "203.0.113.7".parse().unwrap();
2334 let away: std::net::IpAddr = "198.51.100.2".parse().unwrap();
2335 reg.seen_at("dev-ipad", "sarita", home);
2336 reg.seen_at("dev-phone", "admin", home);
2337 assert_eq!(
2338 reg.wake_owner("admin", "dev-phone", "dev-ipad"),
2339 "sarita",
2340 "same address: the iPad's own token"
2341 );
2342 reg.seen_at("dev-phone", "admin", away);
2343 assert_eq!(
2344 reg.wake_owner("admin", "dev-phone", "dev-ipad"),
2345 "admin",
2346 "elsewhere, with no grant: only admin's own, which it is not"
2347 );
2348 reg.share("sarita", "dev-ipad", "admin", true).unwrap();
2349 assert_eq!(
2350 reg.wake_owner("admin", "dev-phone", "dev-ipad"),
2351 "sarita",
2352 "shared: from anywhere"
2353 );
2354 reg.seen_at("dev-other", "mallory", home);
2356 assert_eq!(reg.wake_owner("admin", "dev-other", "dev-tv"), "admin");
2357 }
2358
2359 #[test]
2362 fn where_a_device_last_linked_from_survives_a_restart() {
2363 let conn = rusqlite::Connection::open_in_memory().unwrap();
2364 koan_core::db::schema::create_tables(&conn).unwrap();
2365 for (device, user, seen) in [("dev-ipad", "sarita", 1), ("dev-phone", "admin", 2)] {
2366 conn.execute(
2367 "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES (?1, ?2, ?1, 'ios', ?3)",
2368 rusqlite::params![device, user, seen],
2369 )
2370 .unwrap();
2371 }
2372 let home: std::net::IpAddr = "203.0.113.7".parse().unwrap();
2373 outbox::save_address_in(&conn, "dev-ipad", "sarita", home);
2374 outbox::save_address_in(&conn, "dev-phone", "admin", home);
2375
2376 let reg = Registry::default();
2378 *reg.addresses.lock() = outbox::load_addresses_in(&conn);
2379 assert_eq!(
2380 reg.wake_owner("admin", "dev-phone", "dev-ipad"),
2381 "sarita",
2382 "still woken through its own account after the restart"
2383 );
2384 }
2385
2386 #[test]
2390 fn a_forgotten_device_is_dropped_by_the_accounts_devices() {
2391 let reg = Registry::default();
2392 let (mac_tx, mut mac) = tokio::sync::mpsc::unbounded_channel();
2393 let (other_tx, mut other) = tokio::sync::mpsc::unbounded_channel();
2394 reg.register("j", "mac", "macos", "dev-mac", mac_tx, true);
2395 reg.register("k", "laptop", "macos", "dev-k", other_tx, true);
2396 while mac.try_recv().is_ok() {}
2397 while other.try_recv().is_ok() {}
2398
2399 reg.forget("j", "dev-phone").unwrap();
2400 let heard: Vec<LinkCommand> = std::iter::from_fn(|| mac.try_recv().ok()).collect();
2401 assert!(heard.contains(&LinkCommand::Forgotten {
2402 device: "dev-phone".into()
2403 }));
2404 assert!(
2405 std::iter::from_fn(|| other.try_recv().ok())
2406 .all(|c| !matches!(c, LinkCommand::Forgotten { .. })),
2407 "another account hears nothing of it"
2408 );
2409 assert!(
2410 reg.forget("j", "dev-mac").is_err(),
2411 "linked: it would be back"
2412 );
2413 assert!(
2414 reg.relay(
2415 "j",
2416 "dev-mac",
2417 LinkCommand::Forgotten { device: "x".into() }
2418 )
2419 .is_err(),
2420 "news from the server, not a command a device may send"
2421 );
2422 }
2423
2424 #[test]
2427 fn forgetting_a_shared_device_declines_the_share() {
2428 let (reg, mut phone, _mac, mut k) = shared();
2429 drain(&mut k);
2430 reg.forget("k", "dev-phone").unwrap();
2431 assert!(reg.shared_owner("k", "dev-phone").is_none());
2432 assert!(
2433 drain(&mut phone)
2434 .iter()
2435 .any(|c| matches!(c, LinkCommand::Shares { grantees, .. } if grantees.is_empty()))
2436 );
2437 let listed = drain(&mut k).into_iter().rev().find_map(|c| match c {
2438 LinkCommand::Devices { devices } => Some(devices),
2439 _ => None,
2440 });
2441 assert_eq!(listed.map(|d| d.len()), Some(0));
2442 }
2443
2444 #[test]
2447 fn forgetting_a_device_stops_pushes_to_it_and_only_for_its_account() {
2448 let conn = rusqlite::Connection::open_in_memory().unwrap();
2449 koan_core::db::schema::create_tables(&conn).unwrap();
2450 for user in ["j", "k"] {
2451 conn.execute(
2452 "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES ('dev-phone', ?1, 'phone', 'ios', 1)",
2453 [user],
2454 )
2455 .unwrap();
2456 conn.execute(
2457 "INSERT INTO link_push (device, username, token, sandbox, updated_at) VALUES ('dev-phone', ?1, 't', 0, 1)",
2458 [user],
2459 )
2460 .unwrap();
2461 }
2462 outbox::forget_device_in(&conn, "dev-phone", "k");
2463 assert_eq!(
2464 outbox::push_targets_in(&conn, Some("j")).len(),
2465 1,
2466 "k forgetting its own leaves j's alone"
2467 );
2468 outbox::forget_device_in(&conn, "dev-phone", "j");
2469 assert!(outbox::push_targets_in(&conn, Some("j")).is_empty());
2470 assert!(outbox::push_targets_in(&conn, None).is_empty());
2471 }
2472
2473 #[test]
2474 fn a_device_hears_its_peers_and_commands_reach_them_by_device_id() {
2475 let reg = Registry::default();
2476 let (tx1, mut rx1) = tokio::sync::mpsc::unbounded_channel();
2477 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
2478 let (tx3, mut rx3) = tokio::sync::mpsc::unbounded_channel();
2479 reg.register("j", "mac", "macos", "dev-mac", tx1, true);
2480 let phone = reg.register("j", "phone", "ios", "dev-phone", tx2, true);
2481 reg.register("someone", "laptop", "macos", "dev-other", tx3, true);
2482
2483 reg.report(
2484 &phone,
2485 LinkState {
2486 playing: true,
2487 title: Some("Roygbiv".into()),
2488 ..Default::default()
2489 },
2490 );
2491 let devices = drain(&mut rx1)
2492 .into_iter()
2493 .rev()
2494 .find_map(|c| match c {
2495 LinkCommand::Devices { devices } => Some(devices),
2496 _ => None,
2497 })
2498 .expect("the Mac is told of the phone");
2499 assert_eq!(devices.len(), 1, "not itself, not another account's");
2500 assert_eq!(devices[0].id, "dev-phone");
2501 assert_eq!(
2502 devices[0].state.as_ref().and_then(|s| s.title.as_deref()),
2503 Some("Roygbiv")
2504 );
2505 while rx3.try_recv().is_ok() {}
2506 assert!(rx3.try_recv().is_err());
2507
2508 while rx2.try_recv().is_ok() {}
2509 reg.relay("j", "dev-phone", LinkCommand::Pause).unwrap();
2510 assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
2511 assert!(
2512 reg.relay("someone", "dev-phone", LinkCommand::Pause)
2513 .is_err(),
2514 "another account cannot reach it"
2515 );
2516 assert!(
2517 reg.relay("j", "dev-phone", LinkCommand::Devices { devices: vec![] })
2518 .is_err()
2519 );
2520 }
2521
2522 #[test]
2523 fn a_device_linked_for_a_month_is_not_forgotten() {
2524 let conn = rusqlite::Connection::open_in_memory().unwrap();
2525 koan_core::db::schema::create_tables(&conn).unwrap();
2526 let now = 100 * 24 * 60 * 60;
2527 let long_ago = now - 40 * 24 * 60 * 60;
2528 for device in ["mac", "old-phone"] {
2529 conn.execute(
2530 "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES (?1, 'j', ?1, 'ios', ?2)",
2531 rusqlite::params![device, long_ago],
2532 )
2533 .unwrap();
2534 conn.execute(
2535 "INSERT INTO link_push (device, username, token, sandbox, updated_at) VALUES (?1, 'j', 't', 0, ?2)",
2536 rusqlite::params![device, long_ago],
2537 )
2538 .unwrap();
2539 }
2540 outbox::forget_stale(&conn, &[("mac".into(), "j".into())], now);
2541 let left = |table: &str| -> Vec<String> {
2542 conn.prepare(&format!("SELECT device FROM {table} ORDER BY device"))
2543 .unwrap()
2544 .query_map([], |r| r.get(0))
2545 .unwrap()
2546 .collect::<Result<_, _>>()
2547 .unwrap()
2548 };
2549 assert_eq!(left("link_devices"), ["mac"]);
2550 assert_eq!(left("link_push"), ["mac"]);
2551 }
2552}