use std::sync::LazyLock;
use koan_core::remote::link::{LinkCommand, LinkState};
use parking_lot::Mutex;
use tokio::sync::mpsc::UnboundedSender;
#[derive(Debug, Clone)]
pub struct ClientInfo {
pub id: String,
pub name: String,
pub platform: String,
pub username: String,
pub connected_at: i64,
pub state: LinkState,
pub last_played_at: Option<i64>,
pub state_at: i64,
pub reports: bool,
}
impl ClientInfo {
pub fn position_ms(&self) -> u64 {
let pos = self.state.position_ms;
if !self.state.playing {
return pos;
}
let run = (chrono::Utc::now().timestamp_millis() - self.state_at).max(0) as u64;
let pos = pos + run;
if self.state.duration_ms > 0 {
pos.min(self.state.duration_ms)
} else {
pos
}
}
}
struct Entry {
info: ClientInfo,
device: String,
tx: UnboundedSender<LinkCommand>,
}
#[derive(Debug, Clone)]
pub struct Order {
pub id: String,
pub username: Option<String>,
pub client: Option<String>,
pub artist: String,
pub album: String,
pub play_next: bool,
pub created_at: i64,
}
const ORDER_TTL: i64 = 24 * 60 * 60;
#[derive(Default)]
pub struct Registry {
entries: Mutex<Vec<Entry>>,
orders: Mutex<Vec<Order>>,
}
pub fn registry() -> &'static Registry {
static REGISTRY: LazyLock<Registry> = LazyLock::new(Registry::default);
®ISTRY
}
impl Registry {
pub fn register(
&self,
username: &str,
name: &str,
platform: &str,
device: &str,
tx: UnboundedSender<LinkCommand>,
) -> String {
let id = uuid::Uuid::now_v7().to_string();
let mut entries = self.entries.lock();
entries.retain(|e| !(e.device == device && e.info.username == username));
entries.push(Entry {
info: ClientInfo {
id: id.clone(),
name: name.to_string(),
platform: platform.to_string(),
username: username.to_string(),
connected_at: chrono::Utc::now().timestamp(),
state: LinkState::default(),
last_played_at: None,
state_at: chrono::Utc::now().timestamp_millis(),
reports: false,
},
device: device.to_string(),
tx,
});
id
}
pub fn report(&self, id: &str, state: LinkState) {
let mut entries = self.entries.lock();
if let Some(e) = entries.iter_mut().find(|e| e.info.id == id) {
if state.playing || e.info.state.playing {
e.info.last_played_at = Some(chrono::Utc::now().timestamp());
}
e.info.state = state;
e.info.state_at = chrono::Utc::now().timestamp_millis();
e.info.reports = true;
}
}
pub fn unregister(&self, id: &str) {
self.entries.lock().retain(|e| e.info.id != id);
}
pub fn list(&self, username: Option<&str>) -> Vec<ClientInfo> {
let mut out: Vec<ClientInfo> = self
.entries
.lock()
.iter()
.filter(|e| username.is_none_or(|u| e.info.username == u))
.map(|e| e.info.clone())
.collect();
out.sort_by_key(|c| std::cmp::Reverse(c.connected_at));
out
}
pub fn send(
&self,
username: Option<&str>,
id: Option<&str>,
cmd: LinkCommand,
) -> Result<ClientInfo, String> {
let clients = self.list(username);
let target = match id {
Some(id) => clients
.iter()
.find(|c| c.id == id || c.name.eq_ignore_ascii_case(id))
.ok_or_else(|| format!("no linked client {id}; see `clients`"))?,
None => pick(&clients, chrono::Utc::now().timestamp())?,
};
let entries = self.entries.lock();
let entry = entries
.iter()
.find(|e| e.info.id == target.id)
.ok_or("that client has just gone")?;
entry
.tx
.send(cmd)
.map_err(|_| "that client has just gone".to_string())?;
Ok(target.clone())
}
}
impl Registry {
pub fn add_order(&self, order: Order) {
self.orders.lock().push(order);
}
pub fn orders(&self, username: Option<&str>) -> Vec<Order> {
self.orders
.lock()
.iter()
.filter(|o| username.is_none() || o.username.as_deref() == username)
.cloned()
.collect()
}
pub fn cancel_order(&self, username: Option<&str>, id: &str) -> bool {
let mut orders = self.orders.lock();
let before = orders.len();
orders
.retain(|o| !(o.id == id && (username.is_none() || o.username.as_deref() == username)));
orders.len() != before
}
pub fn fulfil_orders(&self, find: impl Fn(&Order) -> Option<Vec<i64>>) {
let now = chrono::Utc::now().timestamp();
let pending: Vec<Order> = {
let mut orders = self.orders.lock();
orders.retain(|o| now - o.created_at < ORDER_TTL);
orders.clone()
};
for order in pending {
let Some(ids) = find(&order).filter(|ids| !ids.is_empty()) else {
continue;
};
let track_ids = ids.iter().map(i64::to_string).collect();
let cmd = if order.play_next {
LinkCommand::PlayNext { track_ids }
} else {
LinkCommand::Enqueue { track_ids }
};
match self.send(order.username.as_deref(), order.client.as_deref(), cmd) {
Ok(c) => {
log::info!(
"link: {} — {} arrived; queued on {}",
order.artist,
order.album,
c.name
);
self.orders.lock().retain(|o| o.id != order.id);
}
Err(e) => log::info!("link: {} — {} arrived but {e}", order.artist, order.album),
}
}
}
}
pub fn fulfil_from(db_path: &std::path::Path) {
let registry = registry();
if registry.orders.lock().is_empty() {
return;
}
let Ok(db) = koan_core::db::connection::Database::open(db_path) else {
return;
};
registry.fulfil_orders(|order| album_tracks(&db.conn, &order.artist, &order.album));
}
pub fn album_tracks(conn: &rusqlite::Connection, artist: &str, album: &str) -> Option<Vec<i64>> {
let like = |s: &str| format!("%{}%", s.replace(['%', '_'], ""));
let album_id: i64 = conn
.query_row(
"SELECT al.id FROM albums al JOIN artists a ON a.id = al.artist_id
WHERE a.name LIKE ?1 COLLATE NOCASE AND al.title LIKE ?2 COLLATE NOCASE
ORDER BY al.id DESC LIMIT 1",
[like(artist), like(album)],
|r| r.get(0),
)
.ok()?;
let mut stmt = conn
.prepare("SELECT id FROM tracks WHERE album_id = ?1 ORDER BY disc, track_number, id")
.ok()?;
let ids = stmt
.query_map([album_id], |r| r.get(0))
.ok()?
.filter_map(Result::ok)
.collect();
Some(ids)
}
impl Registry {
pub fn broadcast(&self, username: Option<&str>, cmd: LinkCommand) -> Vec<String> {
let ids: Vec<String> = self.list(username).into_iter().map(|c| c.id).collect();
let entries = self.entries.lock();
entries
.iter()
.filter(|e| ids.contains(&e.info.id) && e.tx.send(cmd.clone()).is_ok())
.map(|e| e.info.name.clone())
.collect()
}
}
const RECENT: i64 = 6 * 60 * 60;
fn pick(clients: &[ClientInfo], now: i64) -> Result<&ClientInfo, String> {
if let Some(c) = clients.iter().find(|c| c.state.playing) {
return Ok(c);
}
if let Some(c) = clients
.iter()
.filter(|c| c.last_played_at.is_some_and(|t| now - t < RECENT))
.max_by_key(|c| c.last_played_at)
{
return Ok(c);
}
match clients {
[] => Err("no koan app is linked to this server; open koan on the device".into()),
[only] => Ok(only),
several => Err(format!(
"several koan apps are linked and none has played recently: {}. Ask which, then pass `client`",
several
.iter()
.map(|c| c.name.as_str())
.collect::<Vec<_>>()
.join(", ")
)),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_reconnect_replaces_the_device_and_commands_reach_it() {
let reg = Registry::default();
let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
let (tx3, _rx3) = tokio::sync::mpsc::unbounded_channel();
reg.register("j", "phone", "ios", "dev-1", tx1);
let id = reg.register("j", "phone", "ios", "dev-1", tx2);
reg.register("someone", "laptop", "macos", "dev-2", tx3);
assert_eq!(reg.list(Some("j")).len(), 1);
assert_eq!(reg.list(None).len(), 2);
let sent = reg.send(Some("j"), None, LinkCommand::Pause).unwrap();
assert_eq!(sent.id, id);
assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
assert!(
reg.send(Some("j"), Some("laptop"), LinkCommand::Pause)
.is_err()
);
reg.unregister(&id);
assert!(reg.send(Some("j"), None, LinkCommand::Pause).is_err());
}
#[test]
fn the_device_playing_is_the_one_meant() {
let reg = Registry::default();
let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
let mac = reg.register("j", "mac", "macos", "dev-1", tx1);
let phone = reg.register("j", "phone", "ios", "dev-2", tx2);
let err = reg.send(Some("j"), None, LinkCommand::Pause).unwrap_err();
assert!(err.contains("mac") && err.contains("phone"), "{err}");
reg.report(
&phone,
LinkState {
playing: true,
..Default::default()
},
);
assert_eq!(
reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
phone
);
assert_eq!(rx2.try_recv().unwrap(), LinkCommand::Pause);
reg.report(&phone, LinkState::default());
assert_eq!(
reg.send(Some("j"), None, LinkCommand::Pause).unwrap().id,
phone
);
let _ = mac;
}
}