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