koan-server 0.43.0

GraphQL, Subsonic REST, and MCP server for koan music player.
Documentation
//! Apple's push service, for reaching a koan iOS app that iOS has suspended.
//!
//! A linked app is reached over its WebSocket. iOS suspends an app in the
//! background that is not playing, and the socket dies with it; a push is the
//! supported way back. Two kinds: a background push that wakes the app for
//! half a minute to link and take what waits for it in the outbox, and a
//! notification for playback, which iOS will not let a suspended app start on
//! its own: the person taps it, the app opens and plays.
//!
//! Token auth: a JWT signed with the team's `.p8` key, reused for under an
//! hour as Apple asks. HTTP/2, which is all the gateway speaks.
//!
//! A notification carries a link to its album's cover, which the app's
//! notification service extension fetches and attaches. The extension holds no
//! credentials, so the link authorises itself: an HMAC over one track id and
//! an expiry, keyed by a secret that lives only in this process. It opens that
//! one cover for a few minutes and nothing else.

use std::path::Path;
use std::time::{Duration, SystemTime, UNIX_EPOCH};

use axum::extract::{Path as UrlPath, State};
use axum::response::Response;
use axum::routing::get;
use hmac::{Hmac, Mac};
use jsonwebtoken::{Algorithm, EncodingKey, Header};
use koan_core::config::PushConfig;
use koan_core::db::pool::Pool;
use parking_lot::Mutex;
use serde_json::{Value, json};

/// Apple rejects a token older than an hour and throttles one refreshed more
/// often than every twenty minutes.
const TOKEN_LIFE: Duration = Duration::from_secs(50 * 60);

/// How long a notification's cover link opens its cover. The extension fetches
/// it as the notification arrives; a notification delivered later than this
/// shows without art.
const COVER_LIFE: u64 = 10 * 60;

pub struct Pusher {
    key: EncodingKey,
    key_id: String,
    team_id: String,
    topic: String,
    bearer: Mutex<Option<(String, SystemTime)>>,
    /// Built on first send, which is always on a thread of its own: a blocking
    /// client runs a runtime, and building one inside the server's async
    /// runtime panics.
    http: std::sync::OnceLock<reqwest::blocking::Client>,
    /// `sharing.public_url`, which cover links are built on. Without it
    /// notifications carry no art.
    public_url: Option<String>,
}

/// What became of a push.
#[derive(Debug, PartialEq, Eq)]
pub enum Outcome {
    Sent,
    /// The token is no longer valid for this app: the app was deleted, or the
    /// token belongs to the other gateway. Forget it.
    Gone,
    Failed(String),
}

/// What a push asks of the device.
pub enum Push {
    /// Wake the app to link; whatever waits in the outbox follows.
    Wake,
    /// Show this, and carry `command` for the app to run when it is tapped.
    Notify {
        title: String,
        body: String,
        command: Value,
        /// A cover link from `Pusher::cover_link`.
        image: Option<String>,
    },
}

impl Pusher {
    /// The pusher `cfg` describes, if it names a key that can be read.
    pub fn from_config(cfg: &PushConfig) -> Option<Self> {
        if cfg.key_id.is_empty() || cfg.team_id.is_empty() {
            return None;
        }
        let pem = match (&cfg.key, &cfg.key_path) {
            (Some(pem), _) if !pem.trim().is_empty() => pem.clone(),
            (_, Some(path)) => read_key(path)?,
            _ => return None,
        };
        let key = match EncodingKey::from_ec_pem(pem.as_bytes()) {
            Ok(key) => key,
            Err(e) => {
                log::warn!("push: the APNs key is not an EC private key: {e}");
                return None;
            }
        };
        Some(Self {
            key,
            key_id: cfg.key_id.clone(),
            team_id: cfg.team_id.clone(),
            topic: cfg.topic.clone(),
            bearer: Mutex::new(None),
            http: std::sync::OnceLock::new(),
            public_url: None,
        })
    }

    /// A link that opens the cover of `track_id`'s album for `COVER_LIFE`.
    pub fn cover_link(&self, track_id: i64) -> Option<String> {
        let base = self.public_url.as_deref()?.trim_end_matches('/');
        let expires = unix_now() + COVER_LIFE;
        let sig = cover_sig(track_id, expires);
        Some(format!("{base}/push/cover/{track_id}/{expires}/{sig}"))
    }

    fn bearer(&self) -> Result<String, String> {
        let mut cached = self.bearer.lock();
        if let Some((token, at)) = cached.as_ref()
            && at.elapsed().is_ok_and(|age| age < TOKEN_LIFE)
        {
            return Ok(token.clone());
        }
        let mut header = Header::new(Algorithm::ES256);
        header.kid = Some(self.key_id.clone());
        let iat = SystemTime::now()
            .duration_since(UNIX_EPOCH)
            .map_err(|e| e.to_string())?
            .as_secs();
        let token = jsonwebtoken::encode(
            &header,
            &json!({ "iss": self.team_id, "iat": iat }),
            &self.key,
        )
        .map_err(|e| e.to_string())?;
        *cached = Some((token.clone(), SystemTime::now()));
        Ok(token)
    }

    /// Send `push` to the device holding `token`. Blocks for the round trip,
    /// so never call it from async code.
    pub fn send(&self, token: &str, sandbox: bool, push: &Push) -> Outcome {
        let http = self.http.get_or_init(|| {
            reqwest::blocking::Client::builder()
                .connect_timeout(Duration::from_secs(10))
                .timeout(Duration::from_secs(20))
                .build()
                .unwrap_or_default()
        });
        let bearer = match self.bearer() {
            Ok(b) => b,
            Err(e) => return Outcome::Failed(format!("signing: {e}")),
        };
        let host = if sandbox {
            "api.sandbox.push.apple.com"
        } else {
            "api.push.apple.com"
        };
        let (kind, priority, expires_in) = match push {
            // Low priority is the only priority a background push may have.
            Push::Wake => ("background", "5", 60 * 60),
            // "Play this" an hour late is not what anyone asked for.
            Push::Notify { .. } => ("alert", "10", 10 * 60),
        };
        let expiration = unix_now() + expires_in;
        let response = http
            .post(format!("https://{host}/3/device/{token}"))
            .bearer_auth(bearer)
            .header("apns-topic", &self.topic)
            .header("apns-push-type", kind)
            .header("apns-priority", priority)
            .header("apns-expiration", expiration.to_string())
            .json(&payload(push))
            .send();
        match response {
            Ok(r) if r.status().is_success() => Outcome::Sent,
            Ok(r) => {
                let status = r.status();
                let reason = r
                    .json::<Value>()
                    .ok()
                    .and_then(|v| v["reason"].as_str().map(str::to_owned))
                    .unwrap_or_default();
                if status == reqwest::StatusCode::GONE
                    || matches!(
                        reason.as_str(),
                        "BadDeviceToken" | "Unregistered" | "DeviceTokenNotForTopic"
                    )
                {
                    Outcome::Gone
                } else {
                    Outcome::Failed(format!("{status} {reason}"))
                }
            }
            Err(e) => Outcome::Failed(e.to_string()),
        }
    }
}

/// The JSON a push carries. The command rides under `koan`, beside Apple's
/// `aps`, so the app can act on it without asking the server first.
pub fn payload(push: &Push) -> Value {
    match push {
        Push::Wake => json!({ "aps": { "content-available": 1 } }),
        Push::Notify {
            title,
            body,
            command,
            image,
        } => {
            let mut body = json!({
                "aps": {
                    "alert": { "title": title, "body": body },
                    "sound": "default",
                    // Asked for this moment, by the person it is for: through Focus.
                    "interruption-level": "time-sensitive",
                },
                "koan": command,
            });
            if let Some(image) = image {
                // Hands the notification to the service extension, which
                // attaches the cover before it is shown.
                body["aps"]["mutable-content"] = json!(1);
                body["image"] = json!(image);
            }
            body
        }
    }
}

fn read_key(path: &Path) -> Option<String> {
    match std::fs::read_to_string(path) {
        Ok(pem) => Some(pem),
        Err(e) => {
            log::warn!("push: cannot read the APNs key at {}: {e}", path.display());
            None
        }
    }
}

/// The server's pusher, from the config it started with; `None` when no key
/// is configured, and pushes are simply not sent.
pub fn pusher() -> Option<&'static Pusher> {
    static PUSHER: std::sync::LazyLock<Option<Pusher>> = std::sync::LazyLock::new(|| {
        let cfg = koan_core::config::Config::load().unwrap_or_default();
        let pusher = Pusher::from_config(&cfg.push).map(|p| Pusher {
            public_url: cfg.sharing.public_url.filter(|u| !u.trim().is_empty()),
            ..p
        });
        if pusher.is_some() {
            log::info!("push: APNs key {} for {}", cfg.push.key_id, cfg.push.topic);
        }
        pusher
    });
    PUSHER.as_ref()
}

fn unix_now() -> u64 {
    SystemTime::now()
        .duration_since(UNIX_EPOCH)
        .map_or(0, |d| d.as_secs())
}

/// Minted at start-up and never stored: a restart voids outstanding cover
/// links, which live minutes anyway.
fn cover_mac() -> Hmac<sha2::Sha256> {
    static KEY: std::sync::LazyLock<[u8; 32]> = std::sync::LazyLock::new(|| {
        let mut key = [0; 32];
        getrandom::fill(&mut key).expect("system randomness");
        key
    });
    let mut mac = Hmac::<sha2::Sha256>::new_from_slice(&*KEY).expect("any key length");
    mac.update(b"koan notification cover\0");
    mac
}

fn cover_sig(track_id: i64, expires: u64) -> String {
    use base64::Engine;
    let mut mac = cover_mac();
    mac.update(format!("{track_id}.{expires}").as_bytes());
    base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(mac.finalize().into_bytes())
}

fn cover_sig_valid(track_id: i64, expires: u64, sig: &str) -> bool {
    use base64::Engine;
    let Ok(sig) = base64::engine::general_purpose::URL_SAFE_NO_PAD.decode(sig) else {
        return false;
    };
    let mut mac = cover_mac();
    mac.update(format!("{track_id}.{expires}").as_bytes());
    expires >= unix_now() && mac.verify_slice(&sig).is_ok()
}

#[derive(Clone)]
struct CoverState {
    pool: std::sync::Arc<Pool>,
    covers: std::sync::Arc<crate::covers::Covers>,
}

/// `/push/cover/{track}/{expires}/{sig}`: public, since the extension that
/// fetches it has no login. Anything but a valid, unexpired link is a 404.
pub fn router(
    pool: std::sync::Arc<Pool>,
    covers: std::sync::Arc<crate::covers::Covers>,
) -> axum::Router {
    axum::Router::new()
        .route("/push/cover/{track}/{expires}/{sig}", get(cover))
        .with_state(CoverState { pool, covers })
}

async fn cover(
    State(s): State<CoverState>,
    UrlPath((track, expires, sig)): UrlPath<(i64, u64, String)>,
) -> Response {
    if !cover_sig_valid(track, expires, &sig) {
        return crate::share::not_found();
    }
    let art = crate::share::blocking(move || {
        let db = s.pool.get().ok()?;
        let row = koan_core::db::queries::get_track_row(&db.conn, track).ok()??;
        s.covers
            .cover(std::slice::from_ref(&row), crate::covers::LARGE)
    })
    .await;
    crate::share::jpeg(art, false)
}

#[cfg(test)]
mod tests {
    use super::*;

    // A throwaway P-256 key, generated for these tests and used nowhere else.
    const TEST_KEY: &str = "-----BEGIN PRIVATE KEY-----
MIGHAgEAMBMGByqGSM49AgEGCCqGSM49AwEHBG0wawIBAQQgMavrDWJ3FFXFskYn
rsmSSRycut7lJn10pzM1NivXehuhRANCAARKklUe/y7J3SZEK36mnyt5ejhmbNKT
AtQSJr6Wg9OtOkzZdoOhdRVcNFW8q9peFQ+S7qIcWNbXlhi+cAlpf0ce
-----END PRIVATE KEY-----";

    fn config() -> PushConfig {
        PushConfig {
            key: Some(TEST_KEY.into()),
            key_id: "ABC123DEFG".into(),
            team_id: "TEAM123456".into(),
            ..PushConfig::default()
        }
    }

    #[test]
    fn no_key_no_pusher() {
        assert!(Pusher::from_config(&PushConfig::default()).is_none());
    }

    #[test]
    fn the_token_is_es256_with_the_key_id_and_reused() {
        let pusher = Pusher::from_config(&config()).expect("pusher");
        let token = pusher.bearer().unwrap();
        let header = jsonwebtoken::decode_header(&token).unwrap();
        assert_eq!(header.alg, Algorithm::ES256);
        assert_eq!(header.kid.as_deref(), Some("ABC123DEFG"));
        assert_eq!(pusher.bearer().unwrap(), token);
    }

    #[test]
    fn a_notification_carries_its_command() {
        let command = json!({ "type": "play", "trackIds": ["1"], "startAt": 0 });
        let body = payload(&Push::Notify {
            title: "Play on this iPhone".into(),
            body: "Golden Standard".into(),
            command: command.clone(),
            image: None,
        });
        assert_eq!(body["koan"], command);
        assert_eq!(body["aps"]["alert"]["body"], "Golden Standard");
        assert_eq!(body["aps"]["interruption-level"], "time-sensitive");
        assert!(body["aps"].get("mutable-content").is_none());
        assert_eq!(payload(&Push::Wake)["aps"]["content-available"], 1);

        let body = payload(&Push::Notify {
            title: "Play on this iPhone".into(),
            body: "Golden Standard".into(),
            command,
            image: Some("https://koan.example/push/cover/1/2/x".into()),
        });
        assert_eq!(body["aps"]["mutable-content"], 1);
        assert_eq!(body["image"], "https://koan.example/push/cover/1/2/x");
    }

    #[test]
    fn cover_links_open_one_cover_until_they_expire() {
        let later = unix_now() + 60;
        let sig = cover_sig(7, later);
        assert!(cover_sig_valid(7, later, &sig));
        assert!(!cover_sig_valid(8, later, &sig));
        assert!(!cover_sig_valid(7, later + 1, &sig));
        assert!(!cover_sig_valid(7, later, "not-a-signature"));
        let past = unix_now() - 1;
        assert!(!cover_sig_valid(7, past, &cover_sig(7, past)));
    }
}