Skip to main content

koan_server/
clients.rs

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