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