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(Default)]
111pub struct Registry {
112 entries: Mutex<Vec<Entry>>,
113 orders: Mutex<Vec<Order>>,
114 activities: Mutex<Vec<Activity>>,
115}
116
117pub fn registry() -> &'static Registry {
120 static REGISTRY: LazyLock<Registry> = LazyLock::new(|| {
121 let registry = Registry::default();
122 *registry.orders.lock() = outbox::load_orders();
123 registry
124 });
125 ®ISTRY
126}
127
128impl Registry {
129 pub fn register(
132 &self,
133 username: &str,
134 name: &str,
135 platform: &str,
136 device: &str,
137 tx: UnboundedSender<LinkCommand>,
138 wants_devices: bool,
139 ) -> String {
140 let id = uuid::Uuid::now_v7().to_string();
141 let live = self.live();
144 for cmd in outbox::take_and_remember(username, device, name, platform, &live) {
145 let _ = tx.send(cmd);
146 }
147 let mut entries = self.entries.lock();
148 entries.retain(|e| !(e.device == device && e.info.username == username));
149 entries.push(Entry {
150 info: ClientInfo {
151 id: id.clone(),
152 device: device.to_string(),
153 name: name.to_string(),
154 platform: platform.to_string(),
155 username: username.to_string(),
156 connected_at: chrono::Utc::now().timestamp(),
157 state: LinkState::default(),
158 last_played_at: None,
159 state_at: chrono::Utc::now().timestamp_millis(),
160 reports: false,
161 notified: false,
162 },
163 device: device.to_string(),
164 tx,
165 wants_devices,
166 });
167 drop(entries);
168 self.announce(username);
169 id
170 }
171
172 pub fn report(&self, id: &str, state: LinkState) {
174 let mut entries = self.entries.lock();
175 if let Some(e) = entries.iter_mut().find(|e| e.info.id == id) {
176 if state.playing || e.info.state.playing {
177 e.info.last_played_at = Some(chrono::Utc::now().timestamp());
178 }
179 e.info.state = state;
180 e.info.state_at = chrono::Utc::now().timestamp_millis();
181 e.info.reports = true;
182 let (username, device) = (e.info.username.clone(), e.device.clone());
183 drop(entries);
184 self.announce(&username);
185 self.update_activities(&username, &device);
186 }
187 }
188
189 pub fn set_push(&self, username: &str, device: &str, token: &str, sandbox: bool) {
191 outbox::save_push(username, device, token, sandbox);
192 }
193
194 pub fn disconnect(&self, username: &str) {
197 self.entries.lock().retain(|e| e.info.username != username);
198 }
199
200 pub fn unregister(&self, id: &str) {
201 let mut entries = self.entries.lock();
202 let gone = entries
203 .iter()
204 .find(|e| e.info.id == id)
205 .map(|e| (e.device.clone(), e.info.username.clone()));
206 entries.retain(|e| e.info.id != id);
207 drop(entries);
208 if let Some((device, username)) = gone {
209 outbox::touch(&[(device, username.clone())]);
210 self.announce(&username);
211 }
212 }
213
214 fn announce(&self, username: &str) {
218 let listening = |e: &Entry| e.info.username == username && e.wants_devices;
219 if !self.entries.lock().iter().any(listening) {
220 return;
221 }
222 let asleep = outbox::push_targets(Some(username));
223 let entries = self.entries.lock();
224 let ours: Vec<&Entry> = entries
225 .iter()
226 .filter(|e| e.info.username == username)
227 .collect();
228 if !ours.iter().any(|e| e.wants_devices) {
229 return;
230 }
231 let mut all: Vec<LinkDevice> = ours
232 .iter()
233 .map(|e| LinkDevice {
234 id: e.device.clone(),
235 name: e.info.name.clone(),
236 platform: e.info.platform.clone(),
237 linked: true,
238 state: e.info.reports.then(|| LinkState {
239 position_ms: e.info.position_ms(),
240 ..e.info.state.clone()
241 }),
242 })
243 .collect();
244 for t in asleep {
245 if !all.iter().any(|d| d.id == t.device) {
246 all.push(LinkDevice {
247 id: t.device,
248 name: t.name,
249 platform: t.platform,
250 linked: false,
251 state: None,
252 });
253 }
254 }
255 for e in ours.iter().filter(|e| e.wants_devices) {
256 let devices = all.iter().filter(|d| d.id != e.device).cloned().collect();
257 let _ = e.tx.send(LinkCommand::Devices { devices });
258 }
259 }
260
261 pub fn relay(
263 &self,
264 username: &str,
265 to: &str,
266 command: LinkCommand,
267 ) -> Result<ClientInfo, String> {
268 if matches!(command, LinkCommand::Devices { .. }) {
269 return Err("not a command".into());
270 }
271 self.send(Some(username), Some(to), command)
272 }
273
274 pub fn set_activity(
277 &self,
278 username: &str,
279 watcher: &str,
280 activity: Option<(String, String, bool)>,
281 ) {
282 let mut activities = self.activities.lock();
283 activities.retain(|a| !(a.username == username && a.watcher == watcher));
284 let Some((token, target, sandbox)) = activity else {
285 return;
286 };
287 activities.push(Activity {
288 username: username.to_string(),
289 watcher: watcher.to_string(),
290 target: target.clone(),
291 token,
292 sandbox,
293 sent: None,
294 });
295 drop(activities);
296 self.update_activities(username, &target);
297 }
298
299 fn update_activities(&self, username: &str, target: &str) {
302 let Some(pusher) = crate::push::pusher() else {
303 return;
304 };
305 let Some(info) = self
306 .list(Some(username))
307 .into_iter()
308 .find(|c| c.device == target)
309 else {
310 return;
311 };
312 let state = crate::push::ActivityState::of(&info);
313 let mut due = Vec::new();
314 for a in self.activities.lock().iter_mut() {
315 if a.username == username
316 && a.target == target
317 && a.sent.as_ref().is_none_or(|s| s.differs(&state))
318 {
319 a.sent = Some(state.clone());
320 due.push((a.token.clone(), a.sandbox, a.watcher.clone()));
321 }
322 }
323 if due.is_empty() {
324 return;
325 }
326 let username = username.to_string();
327 std::thread::spawn(move || {
328 for (token, sandbox, watcher) in due {
329 let push = crate::push::Push::Activity(state.clone());
330 match pusher.send(&token, sandbox, &push) {
331 crate::push::Outcome::Sent => {}
332 crate::push::Outcome::Gone => {
333 log::info!("push: a Live Activity on {watcher} has ended");
334 registry().set_activity(&username, &watcher, None);
335 }
336 crate::push::Outcome::Failed(e) => {
337 log::warn!("push: Live Activity on {watcher}: {e}");
338 }
339 }
340 }
341 });
342 }
343
344 pub fn list(&self, username: Option<&str>) -> Vec<ClientInfo> {
346 let mut out: Vec<ClientInfo> = self
347 .entries
348 .lock()
349 .iter()
350 .filter(|e| username.is_none_or(|u| e.info.username == u))
351 .map(|e| e.info.clone())
352 .collect();
353 out.sort_by_key(|c| std::cmp::Reverse(c.connected_at));
354 out
355 }
356
357 pub fn send(
362 &self,
363 username: Option<&str>,
364 id: Option<&str>,
365 cmd: LinkCommand,
366 ) -> Result<ClientInfo, String> {
367 let clients = self.list(username);
368 let target = match id {
369 Some(id) => clients
370 .iter()
371 .find(|c| c.id == id || c.device == id || c.name.eq_ignore_ascii_case(id)),
372 None if clients.is_empty() => None,
373 None => Some(pick(&clients, chrono::Utc::now().timestamp())?),
374 };
375 let Some(target) = target else {
377 return reach_absent(username, id, &cmd).unwrap_or_else(|| {
378 Err(match id {
379 Some(id) => format!("no linked client {id}; see `clients`"),
380 None => "no koan app is linked to this server; open koan on the device".into(),
381 })
382 });
383 };
384 let entries = self.entries.lock();
385 let entry = entries
386 .iter()
387 .find(|e| e.info.id == target.id)
388 .ok_or("that client has just gone")?;
389 entry
390 .tx
391 .send(cmd)
392 .map_err(|_| "that client has just gone".to_string())?;
393 Ok(target.clone())
394 }
395}
396
397impl Registry {
398 pub fn add_order(&self, order: Order) {
399 outbox::save_order(&order);
400 self.orders.lock().push(order);
401 }
402
403 pub fn orders(&self, username: Option<&str>) -> Vec<Order> {
404 self.orders
405 .lock()
406 .iter()
407 .filter(|o| username.is_none() || o.username.as_deref() == username)
408 .cloned()
409 .collect()
410 }
411
412 pub fn cancel_order(&self, username: Option<&str>, id: &str) -> bool {
413 let mut orders = self.orders.lock();
414 let before = orders.len();
415 orders
416 .retain(|o| !(o.id == id && (username.is_none() || o.username.as_deref() == username)));
417 let gone = orders.len() != before;
418 if gone {
419 outbox::drop_order(id);
420 }
421 gone
422 }
423
424 fn done(&self, id: &str) {
425 self.orders.lock().retain(|o| o.id != id);
426 outbox::drop_order(id);
427 }
428
429 pub fn fulfil_orders(&self, find: impl Fn(&Order) -> Option<Vec<String>>) {
432 let now = chrono::Utc::now().timestamp();
433 let pending: Vec<Order> = {
434 let mut orders = self.orders.lock();
435 for o in orders.iter().filter(|o| now - o.created_at >= ORDER_TTL) {
436 outbox::drop_order(&o.id);
437 }
438 orders.retain(|o| now - o.created_at < ORDER_TTL);
439 orders.clone()
440 };
441 for order in pending.into_iter().filter(|o| o.playlist.is_none()) {
442 let Some(ids) = find(&order).filter(|ids| !ids.is_empty()) else {
443 continue;
444 };
445 let track_ids = ids;
446 let cmd = if order.play_next {
447 LinkCommand::PlayNext { track_ids }
448 } else {
449 LinkCommand::Enqueue { track_ids }
450 };
451 match self.send(order.username.as_deref(), order.client.as_deref(), cmd) {
452 Ok(c) => {
453 log::info!(
454 "link: {} — {} arrived; queued on {}",
455 order.artist,
456 order.album,
457 c.name
458 );
459 self.done(&order.id);
460 }
461 Err(e) => log::info!("link: {} — {} arrived but {e}", order.artist, order.album),
463 }
464 }
465 }
466}
467
468pub fn fulfil_from(db_path: &std::path::Path) {
470 let registry = registry();
471 if registry.orders.lock().is_empty() {
472 return;
473 }
474 let Ok(db) = koan_core::db::connection::Database::open_existing(db_path) else {
475 return;
476 };
477 let for_playlists: Vec<Order> = registry
480 .orders
481 .lock()
482 .iter()
483 .filter(|o| o.playlist.is_some())
484 .cloned()
485 .collect();
486 let mut edited = false;
487 for order in for_playlists {
488 let Some(playlist) = order.playlist else {
489 continue;
490 };
491 let Some(ids) = order_tracks(&db.conn, &order) else {
492 continue;
493 };
494 match koan_core::db::queries::add_tracks(&db.conn, playlist, &ids) {
495 Ok(_) => {
496 if !cfg!(test) {
498 koan_core::playlists::push_to_remote(playlist);
499 }
500 log::info!(
501 "link: {} — {} arrived; added {} tracks to playlist {playlist}",
502 order.artist,
503 order.album,
504 ids.len()
505 );
506 registry.done(&order.id);
507 edited = true;
508 }
509 Err(e) => log::warn!("link: could not add to playlist {playlist}: {e}"),
510 }
511 }
512 if edited {
513 changed();
514 }
515 registry.fulfil_orders(|order| {
516 let rows = order_tracks(&db.conn, order)?;
517 queries::uids_in_order(&db.conn, UidKind::Track, &rows).ok()
518 });
519}
520
521fn order_tracks(conn: &rusqlite::Connection, order: &Order) -> Option<Vec<i64>> {
524 let tracks = album_tracks(conn, &order.artist, &order.album)?;
525 if order.titles.is_empty() {
526 return Some(tracks.into_iter().map(|(id, _)| id).collect());
527 }
528 let picked: Vec<i64> = order
529 .titles
530 .iter()
531 .filter_map(|want| {
532 let want = want.to_lowercase();
533 tracks
534 .iter()
535 .find(|(_, t)| t.to_lowercase().contains(&want))
536 .map(|(id, _)| *id)
537 })
538 .collect();
539 (!picked.is_empty()).then_some(picked)
540}
541
542pub fn album_tracks(
545 conn: &rusqlite::Connection,
546 artist: &str,
547 album: &str,
548) -> Option<Vec<(i64, String)>> {
549 let like = |s: &str| format!("%{}%", s.replace(['%', '_'], ""));
550 let album_id: i64 = conn
551 .query_row(
552 "SELECT al.id FROM albums al JOIN artists a ON a.id = al.artist_id
553 WHERE a.name LIKE ?1 COLLATE NOCASE AND al.title LIKE ?2 COLLATE NOCASE
554 ORDER BY al.id DESC LIMIT 1",
555 [like(artist), like(album)],
556 |r| r.get(0),
557 )
558 .ok()?;
559 let mut stmt = conn
560 .prepare("SELECT id, title FROM tracks WHERE album_id = ?1 ORDER BY disc, track_number, id")
561 .ok()?;
562 let tracks = stmt
563 .query_map([album_id], |r| Ok((r.get(0)?, r.get(1)?)))
564 .ok()?
565 .filter_map(Result::ok)
566 .collect();
567 Some(tracks)
568}
569
570impl Registry {
571 pub fn broadcast(&self, username: Option<&str>, cmd: LinkCommand) -> Vec<String> {
574 let ids: Vec<String> = self.list(username).into_iter().map(|c| c.id).collect();
575 let entries = self.entries.lock();
576 entries
577 .iter()
578 .filter(|e| ids.contains(&e.info.id) && e.tx.send(cmd.clone()).is_ok())
579 .map(|e| e.info.name.clone())
580 .collect()
581 }
582}
583
584impl Registry {
585 pub fn deliver(&self, username: Option<&str>, cmd: LinkCommand) -> (Vec<String>, Vec<String>) {
590 let (sent, queued) = self.link_or_queue(username, &cmd);
591 wake(&queued.iter().map(Absent::key).collect::<Vec<_>>());
592 (sent, queued.into_iter().map(|q| q.name).collect())
593 }
594
595 fn link_or_queue(
596 &self,
597 username: Option<&str>,
598 cmd: &LinkCommand,
599 ) -> (Vec<String>, Vec<Absent>) {
600 let sent = self.broadcast(username, cmd.clone());
601 let queued = outbox::queue_for_absent(username, &self.live(), cmd);
602 (sent, queued)
603 }
604
605 fn live(&self) -> Vec<(String, String)> {
607 self.entries
608 .lock()
609 .iter()
610 .map(|e| (e.device.clone(), e.info.username.clone()))
611 .collect()
612 }
613}
614
615impl Absent {
616 fn key(&self) -> (String, String) {
617 (self.device.clone(), self.username.clone())
618 }
619}
620
621fn wake(devices: &[(String, String)]) {
624 let Some(pusher) = crate::push::pusher() else {
625 return;
626 };
627 let targets: Vec<outbox::PushTarget> = outbox::push_targets(None)
628 .into_iter()
629 .filter(|t| {
630 devices
631 .iter()
632 .any(|(device, username)| *device == t.device && *username == t.username)
633 })
634 .collect();
635 if targets.is_empty() {
636 return;
637 }
638 std::thread::spawn(move || {
639 for t in targets {
640 deliver_push(pusher, &t, &crate::push::Push::Wake);
641 }
642 });
643}
644
645fn reach_absent(
654 username: Option<&str>,
655 id: Option<&str>,
656 cmd: &LinkCommand,
657) -> Option<Result<ClientInfo, String>> {
658 let pusher = crate::push::pusher()?;
659 let target = outbox::push_targets(username)
662 .into_iter()
663 .find(|t| id.is_none_or(|id| t.device == id || t.name.eq_ignore_ascii_case(id)))?;
664 let info = ClientInfo {
665 id: target.device.clone(),
666 device: target.device.clone(),
667 name: target.name.clone(),
668 platform: target.platform.clone(),
669 username: target.username.clone(),
670 connected_at: 0,
671 state: LinkState::default(),
672 last_played_at: None,
673 state_at: 0,
674 reports: false,
675 notified: true,
676 };
677 let verb = match cmd {
678 LinkCommand::Play { .. } | LinkCommand::JumpTo { .. } | LinkCommand::Resume => Some("Play"),
679 _ => None,
680 };
681 let push = match verb {
682 Some(verb) => crate::push::Push::Notify {
683 title: format!("{verb} on {}", target.name),
684 body: outbox::describe(cmd).unwrap_or_else(|| "From your koan server".into()),
685 command: serde_json::to_value(cmd).ok()?,
686 image: cover_track(cmd)
687 .and_then(outbox::track_row)
688 .and_then(|t| pusher.cover_link(t)),
689 },
690 None => {
691 outbox::queue_for(&target.device, &target.username, cmd);
692 crate::push::Push::Wake
693 }
694 };
695 std::thread::spawn(move || deliver_push(pusher, &target, &push));
696 Some(Ok(info))
697}
698
699fn cover_track(cmd: &LinkCommand) -> Option<&str> {
702 match cmd {
703 LinkCommand::Play {
704 track_ids,
705 start_at,
706 ..
707 } => track_ids
708 .get(*start_at as usize)
709 .or(track_ids.first())
710 .map(String::as_str),
711 LinkCommand::JumpTo { track_id } => Some(track_id),
712 _ => None,
713 }
714}
715
716fn deliver_push(
718 pusher: &crate::push::Pusher,
719 target: &outbox::PushTarget,
720 push: &crate::push::Push,
721) {
722 use crate::push::Outcome;
723 match pusher.send(&target.token, target.sandbox, push) {
724 Outcome::Sent => log::info!("push: sent to {}", target.name),
725 Outcome::Gone => {
726 log::info!(
727 "push: {}'s token is no longer valid; forgotten",
728 target.name
729 );
730 outbox::forget_push(&target.username, &target.device);
731 }
732 Outcome::Failed(e) => log::warn!("push: to {} failed: {e}", target.name),
733 }
734}
735
736pub fn changed() {
741 let (_, queued) = registry().link_or_queue(None, &LinkCommand::Sync { full: false });
742 if queued.is_empty() || crate::push::pusher().is_none() {
743 return;
744 }
745 let mut wakes = WAKES.lock();
746 wakes.add(queued.iter().map(Absent::key), std::time::Instant::now());
747 if !wakes.timer {
748 wakes.timer = true;
749 std::thread::spawn(send_wakes);
750 }
751}
752
753const QUIET: std::time::Duration = std::time::Duration::from_secs(30);
757
758const LONGEST_WAIT: std::time::Duration = std::time::Duration::from_secs(5 * 60);
760
761static WAKES: Mutex<Wakes> = Mutex::new(Wakes {
762 pending: Vec::new(),
763 timer: false,
764});
765
766struct Wakes {
769 pending: Vec<((String, String), std::time::Instant, std::time::Instant)>,
771 timer: bool,
773}
774
775impl Wakes {
776 fn add(
777 &mut self,
778 devices: impl IntoIterator<Item = (String, String)>,
779 now: std::time::Instant,
780 ) {
781 for key in devices {
782 match self.pending.iter_mut().find(|(k, _, _)| *k == key) {
783 Some((_, _, last)) => *last = now,
784 None => self.pending.push((key, now, now)),
785 }
786 }
787 }
788
789 fn due_at(first: std::time::Instant, last: std::time::Instant) -> std::time::Instant {
790 (last + QUIET).min(first + LONGEST_WAIT)
791 }
792
793 fn next(&self) -> Option<std::time::Instant> {
795 self.pending
796 .iter()
797 .map(|(_, first, last)| Self::due_at(*first, *last))
798 .min()
799 }
800
801 fn take_due(&mut self, now: std::time::Instant) -> Vec<(String, String)> {
803 let (due, waiting) = std::mem::take(&mut self.pending)
804 .into_iter()
805 .partition(|(_, first, last)| Self::due_at(*first, *last) <= now);
806 self.pending = waiting;
807 due.into_iter().map(|(key, _, _)| key).collect()
808 }
809}
810
811fn send_wakes() {
814 loop {
815 let (due, next) = {
816 let mut wakes = WAKES.lock();
817 let due = wakes.take_due(std::time::Instant::now());
818 let next = wakes.next();
819 if due.is_empty() && next.is_none() {
820 wakes.timer = false;
821 return;
822 }
823 (due, next)
824 };
825 let live = registry().live();
826 let absent: Vec<(String, String)> = due.into_iter().filter(|d| !live.contains(d)).collect();
827 if !absent.is_empty() {
828 wake(&absent);
829 }
830 if let Some(next) = next {
831 std::thread::sleep(next.saturating_duration_since(std::time::Instant::now()));
832 }
833 }
834}
835
836pub fn changed_if_library_moved(conn: &rusqlite::Connection) {
840 static LAST: parking_lot::Mutex<Option<(i64, i64, i64)>> = parking_lot::Mutex::new(None);
841 let Ok(now) = conn.query_row(
842 "SELECT (SELECT COUNT(*) FROM tracks), (SELECT COALESCE(MAX(id), 0) FROM tracks),
843 (SELECT COUNT(*) FROM albums)",
844 [],
845 |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
846 ) else {
847 return;
848 };
849 let before = LAST.lock().replace(now);
850 if before.is_some_and(|b| b != now) {
851 changed();
852 }
853}
854
855mod outbox {
858 use koan_core::db::queries;
859 use koan_core::remote::link::LinkCommand;
860
861 const KEEP_SECS: i64 = 30 * 24 * 60 * 60;
866
867 fn db() -> Option<koan_core::db::pool::Handle<'static>> {
872 if cfg!(test) {
873 return None;
874 }
875 koan_core::db::pool::shared().get().ok()
876 }
877
878 pub fn load_orders() -> Vec<super::Order> {
879 let Some(db) = db() else { return Vec::new() };
880 db.conn
881 .prepare("SELECT body FROM link_orders ORDER BY created_at")
882 .and_then(|mut s| {
883 s.query_map([], |r| r.get::<_, String>(0))?
884 .collect::<Result<Vec<_>, _>>()
885 })
886 .unwrap_or_default()
887 .into_iter()
888 .filter_map(|b| serde_json::from_str(&b).ok())
889 .collect()
890 }
891
892 pub fn save_order(order: &super::Order) {
893 let (Some(db), Ok(body)) = (db(), serde_json::to_string(order)) else {
894 return;
895 };
896 let _ = db.conn.execute(
897 "INSERT OR REPLACE INTO link_orders (id, body, created_at) VALUES (?1, ?2, ?3)",
898 rusqlite::params![order.id, body, order.created_at],
899 );
900 }
901
902 pub fn drop_order(id: &str) {
903 if let Some(db) = db() {
904 let _ = db
905 .conn
906 .execute("DELETE FROM link_orders WHERE id = ?1", [id]);
907 }
908 }
909
910 pub fn take_and_remember(
911 username: &str,
912 device: &str,
913 name: &str,
914 platform: &str,
915 live: &[(String, String)],
916 ) -> Vec<LinkCommand> {
917 let Some(db) = db() else { return Vec::new() };
918 let now = chrono::Utc::now().timestamp();
919 let _ = db.conn.execute(
920 "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES (?1, ?2, ?3, ?4, ?5)
921 ON CONFLICT (device, username) DO UPDATE SET name = ?3, platform = ?4, last_seen = ?5",
922 rusqlite::params![device, username, name, platform, now],
923 );
924 forget_stale(&db.conn, live, now);
925 let waiting: Vec<(i64, String)> = db
926 .conn
927 .prepare("SELECT id, command FROM link_outbox WHERE device = ?1 AND username = ?2 ORDER BY id")
928 .and_then(|mut s| {
929 s.query_map([device, username], |r| Ok((r.get(0)?, r.get(1)?)))?
930 .collect()
931 })
932 .unwrap_or_default();
933 let _ = db.conn.execute(
934 "DELETE FROM link_outbox WHERE device = ?1 AND username = ?2",
935 [device, username],
936 );
937 if !waiting.is_empty() {
938 log::info!("link: {} waiting commands for {name}", waiting.len());
939 }
940 waiting
941 .into_iter()
942 .filter_map(|(_, c)| serde_json::from_str(&c).ok())
943 .collect()
944 }
945
946 pub(super) fn forget_stale(conn: &rusqlite::Connection, live: &[(String, String)], now: i64) {
949 touch_with(conn, live, now);
950 let _ = conn.execute(
951 "DELETE FROM link_outbox WHERE created_at < ?1",
952 [now - KEEP_SECS],
953 );
954 let _ = conn.execute(
955 "DELETE FROM link_push WHERE (device, username) IN
956 (SELECT device, username FROM link_devices WHERE last_seen < ?1)",
957 [now - KEEP_SECS],
958 );
959 let _ = conn.execute(
960 "DELETE FROM link_devices WHERE last_seen < ?1",
961 [now - KEEP_SECS],
962 );
963 }
964
965 pub fn touch(devices: &[(String, String)]) {
967 if let Some(db) = db() {
968 touch_with(&db.conn, devices, chrono::Utc::now().timestamp());
969 }
970 }
971
972 fn touch_with(conn: &rusqlite::Connection, devices: &[(String, String)], now: i64) {
973 for (device, username) in devices {
974 let _ = conn.execute(
975 "UPDATE link_devices SET last_seen = ?1 WHERE device = ?2 AND username = ?3",
976 rusqlite::params![now, device, username],
977 );
978 }
979 }
980
981 pub struct Absent {
983 pub device: String,
984 pub username: String,
985 pub name: String,
986 }
987
988 pub fn queue_for_absent(
990 username: Option<&str>,
991 live: &[(String, String)],
992 cmd: &LinkCommand,
993 ) -> Vec<Absent> {
994 let Some(db) = db() else { return Vec::new() };
995 let known: Vec<(String, String, String)> = db
996 .conn
997 .prepare("SELECT device, username, name FROM link_devices")
998 .and_then(|mut s| {
999 s.query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))?
1000 .collect()
1001 })
1002 .unwrap_or_default();
1003 let Ok(text) = serde_json::to_string(cmd) else {
1004 return Vec::new();
1005 };
1006 let absent: Vec<_> = known
1007 .into_iter()
1008 .filter(|(device, user, _)| {
1009 !username.is_some_and(|u| u != user)
1010 && !live.iter().any(|(d, u)| d == device && u == user)
1011 })
1012 .collect();
1013 if absent.is_empty() {
1014 return Vec::new();
1015 }
1016 let is_sync = matches!(cmd, LinkCommand::Sync { .. });
1017 let now = chrono::Utc::now().timestamp();
1018 koan_core::db::queries::atomically(&db.conn, || {
1021 let mut queued = Vec::new();
1022 for (device, user, name) in absent {
1023 if is_sync {
1024 let _ = db.conn.execute(
1026 "DELETE FROM link_outbox WHERE device = ?1 AND username = ?2 AND command LIKE '{\"type\":\"sync\"%'",
1027 [&device, &user],
1028 );
1029 }
1030 if db
1031 .conn
1032 .execute(
1033 "INSERT INTO link_outbox (device, username, command, created_at) VALUES (?1, ?2, ?3, ?4)",
1034 rusqlite::params![device, user, text, now],
1035 )
1036 .is_ok()
1037 {
1038 queued.push(Absent {
1039 device,
1040 username: user,
1041 name,
1042 });
1043 }
1044 }
1045 Ok::<_, rusqlite::Error>(queued)
1046 })
1047 .unwrap_or_default()
1048 }
1049
1050 #[derive(Clone)]
1052 pub struct PushTarget {
1053 pub device: String,
1054 pub username: String,
1055 pub name: String,
1056 pub platform: String,
1057 pub token: String,
1058 pub sandbox: bool,
1059 }
1060
1061 pub fn queue_for(device: &str, username: &str, cmd: &LinkCommand) {
1063 let (Some(db), Ok(text)) = (db(), serde_json::to_string(cmd)) else {
1064 return;
1065 };
1066 let _ = db.conn.execute(
1067 "INSERT INTO link_outbox (device, username, command, created_at) VALUES (?1, ?2, ?3, ?4)",
1068 rusqlite::params![device, username, text, chrono::Utc::now().timestamp()],
1069 );
1070 }
1071
1072 pub fn save_push(username: &str, device: &str, token: &str, sandbox: bool) {
1073 let Some(db) = db() else { return };
1074 let _ = db.conn.execute(
1075 "INSERT INTO link_push (device, username, token, sandbox, updated_at) VALUES (?1, ?2, ?3, ?4, ?5)
1076 ON CONFLICT (device, username) DO UPDATE SET token = ?3, sandbox = ?4, updated_at = ?5",
1077 rusqlite::params![device, username, token, sandbox, chrono::Utc::now().timestamp()],
1078 );
1079 }
1080
1081 pub fn forget_push(username: &str, device: &str) {
1082 if let Some(db) = db() {
1083 let _ = db.conn.execute(
1084 "DELETE FROM link_push WHERE device = ?1 AND username = ?2",
1085 [device, username],
1086 );
1087 }
1088 }
1089
1090 pub fn push_targets(username: Option<&str>) -> Vec<PushTarget> {
1092 let Some(db) = db() else { return Vec::new() };
1093 db.conn
1094 .prepare(
1095 "SELECT p.device, p.username, d.name, d.platform, p.token, p.sandbox
1096 FROM link_push p JOIN link_devices d ON d.device = p.device AND d.username = p.username
1097 WHERE ?1 IS NULL OR p.username = ?1
1098 ORDER BY d.last_seen DESC",
1099 )
1100 .and_then(|mut s| {
1101 s.query_map([username], |r| {
1102 Ok(PushTarget {
1103 device: r.get(0)?,
1104 username: r.get(1)?,
1105 name: r.get(2)?,
1106 platform: r.get(3)?,
1107 token: r.get(4)?,
1108 sandbox: r.get(5)?,
1109 })
1110 })?
1111 .collect()
1112 })
1113 .unwrap_or_default()
1114 }
1115
1116 pub fn track_row(id: &str) -> Option<i64> {
1118 let db = db()?;
1119 queries::resolve_id(&db.conn, queries::UidKind::Track, id)
1120 .ok()
1121 .flatten()
1122 }
1123
1124 pub fn describe(cmd: &LinkCommand) -> Option<String> {
1127 let ids: Vec<&String> = match cmd {
1128 LinkCommand::Play { track_ids, .. }
1129 | LinkCommand::Enqueue { track_ids }
1130 | LinkCommand::PlayNext { track_ids } => track_ids.iter().collect(),
1131 LinkCommand::JumpTo { track_id } => vec![track_id],
1132 _ => return None,
1133 };
1134 let db = db()?;
1135 let ids: Vec<i64> = ids
1136 .into_iter()
1137 .filter_map(|t| {
1138 queries::resolve_id(&db.conn, queries::UidKind::Track, t)
1139 .ok()
1140 .flatten()
1141 })
1142 .collect();
1143 let row = |id: i64| {
1144 db.conn
1145 .query_row(
1146 "SELECT t.title, COALESCE(a.name, ''), COALESCE(al.title, ''), t.album_id
1147 FROM tracks t LEFT JOIN artists a ON a.id = t.artist_id
1148 LEFT JOIN albums al ON al.id = t.album_id WHERE t.id = ?1",
1149 [id],
1150 |r| {
1151 Ok((
1152 r.get::<_, String>(0)?,
1153 r.get::<_, String>(1)?,
1154 r.get::<_, String>(2)?,
1155 r.get::<_, Option<i64>>(3)?,
1156 ))
1157 },
1158 )
1159 .ok()
1160 };
1161 let (title, artist, album, album_id) = row(*ids.first()?)?;
1162 let one_album = ids.len() > 1
1163 && album_id.is_some()
1164 && ids
1165 .iter()
1166 .all(|id| row(*id).is_some_and(|r| r.3 == album_id));
1167 Some(match (one_album, ids.len()) {
1168 (true, _) => format!("{album} — {artist}"),
1169 (false, 1) => format!("{title} — {artist}"),
1170 (false, n) => format!("{title} — {artist}, and {} more", n - 1),
1171 })
1172 }
1173}
1174
1175const RECENT: i64 = 6 * 60 * 60;
1178
1179fn pick(clients: &[ClientInfo], now: i64) -> Result<&ClientInfo, String> {
1180 if let Some(c) = clients.iter().find(|c| c.state.playing) {
1181 return Ok(c);
1182 }
1183 if let Some(c) = clients
1184 .iter()
1185 .filter(|c| c.last_played_at.is_some_and(|t| now - t < RECENT))
1186 .max_by_key(|c| c.last_played_at)
1187 {
1188 return Ok(c);
1189 }
1190 match clients {
1191 [] => Err("no koan app is linked to this server; open koan on the device".into()),
1192 [only] => Ok(only),
1193 several => Err(format!(
1194 "several koan apps are linked and none has played recently: {}. Ask which, then pass `client`",
1195 several
1196 .iter()
1197 .map(|c| format!("{} ({}, id {})", c.name, c.platform, c.device))
1198 .collect::<Vec<_>>()
1199 .join(", ")
1200 )),
1201 }
1202}
1203
1204#[cfg(test)]
1205mod tests {
1206
1207 #[test]
1208 fn a_playlist_order_adds_the_named_tracks_once_they_arrive() {
1209 let dir = tempfile::tempdir().unwrap();
1210 let path = dir.path().join("koan.db");
1211 let db = koan_core::db::connection::Database::open(&path).unwrap();
1212 let playlist = koan_core::db::queries::create_playlist(
1213 &db.conn,
1214 koan_core::db::queries::LOCAL_USER,
1215 "cyberpunk",
1216 None,
1217 )
1218 .unwrap();
1219 let order = Order {
1220 id: "o1".into(),
1221 username: None,
1222 client: None,
1223 artist: "Perturbator".into(),
1224 album: "Dangerous Days".into(),
1225 play_next: false,
1226 playlist: Some(playlist),
1227 titles: vec!["Future Club".into()],
1228 created_at: chrono::Utc::now().timestamp(),
1229 };
1230 registry().add_order(order);
1231
1232 fulfil_from(&path);
1234 assert!(registry().orders(None).iter().any(|o| o.id == "o1"));
1235
1236 db.conn
1237 .execute_batch(
1238 "INSERT INTO artists (id, name) VALUES (1, 'Perturbator');
1239 INSERT INTO albums (id, title, artist_id) VALUES (1, 'Dangerous Days', 1);
1240 INSERT INTO tracks (id, title, album_id, artist_id, track_number, path) VALUES
1241 (1, 'Welcome Back', 1, 1, 1, '/1.flac'), (2, 'Future Club', 1, 1, 2, '/2.flac');",
1242 )
1243 .unwrap();
1244 fulfil_from(&path);
1245 let held: Vec<i64> = db
1246 .conn
1247 .prepare("SELECT track_id FROM playlist_tracks WHERE playlist_id = ?1")
1248 .unwrap()
1249 .query_map([playlist], |r| r.get(0))
1250 .unwrap()
1251 .collect::<Result<_, _>>()
1252 .unwrap();
1253 assert_eq!(held, [2]);
1254 assert!(!registry().orders(None).iter().any(|o| o.id == "o1"));
1255 }
1256
1257 use super::*;
1258
1259 #[test]
1260 fn a_notification_cover_is_named_by_the_uid_the_command_carries() {
1261 let uid = "0190a5b2-7c3d-7e4f-8a1b-2c3d4e5f6a7b".to_string();
1263 let play = LinkCommand::Play {
1264 track_ids: vec!["x".into(), uid.clone()],
1265 start_at: 1,
1266 position_ms: 0,
1267 paused: false,
1268 handoff: false,
1269 };
1270 assert_eq!(cover_track(&play), Some(uid.as_str()));
1271 }
1272
1273 #[test]
1274 fn a_reconnect_replaces_the_device_and_commands_reach_it() {
1275 let reg = Registry::default();
1276 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
1277 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
1278 let (tx3, _rx3) = tokio::sync::mpsc::unbounded_channel();
1279 reg.register("j", "phone", "ios", "dev-1", tx1, false);
1280 let id = reg.register("j", "phone", "ios", "dev-1", tx2, false);
1281 reg.register("someone", "laptop", "macos", "dev-2", tx3, false);
1282
1283 assert_eq!(reg.list(Some("j")).len(), 1);
1284 assert_eq!(reg.list(None).len(), 2);
1285
1286 let sent = reg.send(Some("j"), None, LinkCommand::Pause).unwrap();
1287 assert_eq!(sent.id, id);
1288 assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
1289
1290 assert!(
1292 reg.send(Some("j"), Some("laptop"), LinkCommand::Pause)
1293 .is_err()
1294 );
1295
1296 reg.unregister(&id);
1297 assert!(reg.send(Some("j"), None, LinkCommand::Pause).is_err());
1298 }
1299
1300 #[test]
1301 fn disconnecting_an_account_closes_only_its_links() {
1302 let reg = Registry::default();
1303 let (tx1, mut rx1) = tokio::sync::mpsc::unbounded_channel();
1304 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
1305 reg.register("j", "phone", "ios", "dev-1", tx1, false);
1306 reg.register("someone", "laptop", "macos", "dev-2", tx2, false);
1307
1308 reg.disconnect("j");
1309 assert!(reg.list(Some("j")).is_empty());
1310 assert!(matches!(
1312 rx1.try_recv(),
1313 Err(tokio::sync::mpsc::error::TryRecvError::Disconnected)
1314 ));
1315 assert_eq!(reg.list(Some("someone")).len(), 1);
1316 assert!(matches!(
1317 rx2.try_recv(),
1318 Err(tokio::sync::mpsc::error::TryRecvError::Empty)
1319 ));
1320 }
1321
1322 #[test]
1323 fn a_burst_of_changes_wakes_each_device_once_when_it_goes_quiet() {
1324 let t0 = std::time::Instant::now();
1325 let s = std::time::Duration::from_secs;
1326 let phone = || ("dev-1".to_string(), "j".to_string());
1327 let ipad = || ("dev-2".to_string(), "j".to_string());
1328 let mut wakes = Wakes {
1329 pending: Vec::new(),
1330 timer: false,
1331 };
1332
1333 wakes.add([phone()], t0);
1334 wakes.add([phone(), ipad()], t0 + s(10));
1335 wakes.add([phone()], t0 + s(20));
1336 assert_eq!(wakes.pending.len(), 2);
1337
1338 assert!(wakes.take_due(t0 + s(39)).is_empty());
1340 assert_eq!(wakes.take_due(t0 + s(40)), [ipad()]);
1341 assert_eq!(wakes.next(), Some(t0 + s(50)));
1342 assert_eq!(wakes.take_due(t0 + s(50)), [phone()]);
1343 assert_eq!(wakes.next(), None);
1344 }
1345
1346 #[test]
1347 fn a_library_that_never_goes_quiet_still_wakes_devices() {
1348 let t0 = std::time::Instant::now();
1349 let phone = || ("dev-1".to_string(), "j".to_string());
1350 let mut wakes = Wakes {
1351 pending: Vec::new(),
1352 timer: false,
1353 };
1354 let mut sent = 0;
1355 for i in 0..40 {
1356 let now = t0 + std::time::Duration::from_secs(i * 10);
1357 sent += wakes.take_due(now).len();
1358 wakes.add([phone()], now);
1359 }
1360 assert_eq!(sent, 1);
1363 assert_eq!(wakes.pending.len(), 1);
1364 }
1365
1366 #[test]
1367 fn the_device_playing_is_the_one_meant() {
1368 let reg = Registry::default();
1369 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
1370 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
1371 let mac = reg.register("j", "mac", "macos", "dev-1", tx1, false);
1372 let phone = reg.register("j", "phone", "ios", "dev-2", tx2, false);
1373
1374 let err = reg.send(Some("j"), None, LinkCommand::Pause).unwrap_err();
1376 assert!(err.contains("mac") && err.contains("phone"), "{err}");
1377
1378 reg.report(
1379 &phone,
1380 LinkState {
1381 playing: true,
1382 ..Default::default()
1383 },
1384 );
1385 assert_eq!(
1386 reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
1387 phone
1388 );
1389 assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
1390
1391 reg.report(&phone, LinkState::default());
1393 assert_eq!(
1394 reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
1395 phone
1396 );
1397 let _ = mac;
1398 }
1399
1400 #[test]
1401 fn a_device_hears_its_peers_and_commands_reach_them_by_device_id() {
1402 let reg = Registry::default();
1403 let (tx1, mut rx1) = tokio::sync::mpsc::unbounded_channel();
1404 let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
1405 let (tx3, mut rx3) = tokio::sync::mpsc::unbounded_channel();
1406 reg.register("j", "mac", "macos", "dev-mac", tx1, true);
1407 let phone = reg.register("j", "phone", "ios", "dev-phone", tx2, true);
1408 reg.register("someone", "laptop", "macos", "dev-other", tx3, true);
1409
1410 reg.report(
1411 &phone,
1412 LinkState {
1413 playing: true,
1414 title: Some("Roygbiv".into()),
1415 ..Default::default()
1416 },
1417 );
1418 let mut last = None;
1419 while let Ok(LinkCommand::Devices { devices }) = rx1.try_recv() {
1420 last = Some(devices);
1421 }
1422 let devices = last.expect("the Mac is told of the phone");
1423 assert_eq!(devices.len(), 1, "not itself, not another account's");
1424 assert_eq!(devices[0].id, "dev-phone");
1425 assert_eq!(
1426 devices[0].state.as_ref().and_then(|s| s.title.as_deref()),
1427 Some("Roygbiv")
1428 );
1429 while rx3.try_recv().is_ok() {}
1430 assert!(rx3.try_recv().is_err());
1431
1432 while rx2.try_recv().is_ok() {}
1433 reg.relay("j", "dev-phone", LinkCommand::Pause).unwrap();
1434 assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
1435 assert!(
1436 reg.relay("someone", "dev-phone", LinkCommand::Pause)
1437 .is_err(),
1438 "another account cannot reach it"
1439 );
1440 assert!(
1441 reg.relay("j", "dev-phone", LinkCommand::Devices { devices: vec![] })
1442 .is_err()
1443 );
1444 }
1445
1446 #[test]
1447 fn a_device_linked_for_a_month_is_not_forgotten() {
1448 let conn = rusqlite::Connection::open_in_memory().unwrap();
1449 koan_core::db::schema::create_tables(&conn).unwrap();
1450 let now = 100 * 24 * 60 * 60;
1451 let long_ago = now - 40 * 24 * 60 * 60;
1452 for device in ["mac", "old-phone"] {
1453 conn.execute(
1454 "INSERT INTO link_devices (device, username, name, platform, last_seen) VALUES (?1, 'j', ?1, 'ios', ?2)",
1455 rusqlite::params![device, long_ago],
1456 )
1457 .unwrap();
1458 conn.execute(
1459 "INSERT INTO link_push (device, username, token, sandbox, updated_at) VALUES (?1, 'j', 't', 0, ?2)",
1460 rusqlite::params![device, long_ago],
1461 )
1462 .unwrap();
1463 }
1464 outbox::forget_stale(&conn, &[("mac".into(), "j".into())], now);
1465 let left = |table: &str| -> Vec<String> {
1466 conn.prepare(&format!("SELECT device FROM {table} ORDER BY device"))
1467 .unwrap()
1468 .query_map([], |r| r.get(0))
1469 .unwrap()
1470 .collect::<Result<_, _>>()
1471 .unwrap()
1472 };
1473 assert_eq!(left("link_devices"), ["mac"]);
1474 assert_eq!(left("link_push"), ["mac"]);
1475 }
1476}