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