diff --git a/src/app/cmdline.rs b/src/app/cmdline.rs index f9c2920..a0dffb5 100644 --- a/src/app/cmdline.rs +++ b/src/app/cmdline.rs @@ -50,6 +50,7 @@ fn apply_live(state: &mut AppState, runtime: &Runtime, command: Command) { Command::Quit | Command::Import(_) | Command::Open(_) + | Command::ConnectInvite(_) | Command::Volume(_) | Command::Seek(_) | Command::SeekTo(_) @@ -167,6 +168,7 @@ fn execute(state: &mut AppState, runtime: &mut Runtime, command: Command) { Command::Quit => state.should_quit = true, Command::Import(path) => super::spawn_import(state, runtime, &path), Command::Open(link) => open_frid_link(state, runtime, link), + Command::ConnectInvite(invite) => super::device_connect(runtime, invite), Command::Volume(value) => { state.player.volume = value; super::perform_effect(state, runtime, Effect::SetVolume(value)); diff --git a/src/app/command.rs b/src/app/command.rs index ca2eca9..d9f3f0a 100644 --- a/src/app/command.rs +++ b/src/app/command.rs @@ -27,6 +27,8 @@ pub enum Command { Import(String), /// `:open frid://...` — open a shared federation content link. Open(String), + /// `:connect frid://i/...` — pair this client with a trusted device. + ConnectInvite(String), /// `:volume 40` (also `:vol`) — set the volume precisely. Volume(u8), /// `:seek +30` / `:seek -10` — relative seek in seconds. @@ -91,6 +93,15 @@ pub fn parse(input: &str) -> Parsed { _ => Parsed::Invalid("usage: :open frid://".to_string()), } } + "connect" => { + let value = input.trim_start().split_once(char::is_whitespace); + match value.map(|(_, rest)| rest.trim()) { + Some(value) if !value.is_empty() => { + Parsed::Command(Command::ConnectInvite(value.into())) + } + _ => Parsed::Invalid("usage: :connect frid://i/".to_string()), + } + } "volume" | "vol" => match arg.and_then(|a| a.parse::().ok()) { Some(value) if value <= 100 => Parsed::Command(Command::Volume(value)), _ => Parsed::Invalid("usage: :volume 0-100".to_string()), @@ -182,6 +193,10 @@ mod tests { ); assert!(matches!(parse("import"), Parsed::Invalid(_))); assert!(matches!(parse("open"), Parsed::Invalid(_))); + assert_eq!( + parse("connect frid://i/abcd"), + Parsed::Command(Command::ConnectInvite("frid://i/abcd".to_string())) + ); assert_eq!(parse("volume 40"), Parsed::Command(Command::Volume(40))); assert_eq!(parse("vol 0"), Parsed::Command(Command::Volume(0))); assert_eq!(parse("shuffle"), Parsed::Command(Command::Shuffle)); diff --git a/src/app/event.rs b/src/app/event.rs index 598e9ec..7495dd4 100644 --- a/src/app/event.rs +++ b/src/app/event.rs @@ -128,4 +128,12 @@ pub enum AppEvent { }, /// This peer's connection ticket, requested from the Federation tab. FedTicket(Result), + /// Fresh personal-device sync status snapshot for Settings. + DeviceSyncStatus(crate::devices::DeviceSyncStatus), + /// Invite link for pairing another device. + DeviceInvite(Result), + /// Result of `:connect frid://i/...`. + DeviceConnectResult(Result), + /// Incoming pairing request that passed the invite-secret check. + DevicePairingRequest(crate::devices::PendingPairing), } diff --git a/src/app/mod.rs b/src/app/mod.rs index db3a984..d7c85e7 100644 --- a/src/app/mod.rs +++ b/src/app/mod.rs @@ -7,9 +7,9 @@ mod popup; pub mod state; pub mod update; -use std::io; +use std::io::{self, Write as _}; use std::path::{Path, PathBuf}; -use std::process::Command; +use std::process::{Command, Stdio}; use std::sync::Arc; use std::time::Duration; @@ -36,6 +36,7 @@ const VISUALIZER_TICK_INTERVAL: Duration = Duration::from_millis(50); pub struct Runtime { pub event_tx: mpsc::UnboundedSender, pub library: Arc, + pub devices: Arc, pub federation: Arc, /// When the last Federation-tab status snapshot was requested. pub fed_status_at: Option, @@ -92,12 +93,16 @@ pub async fn run( state.status_message = Some(format!("visualizations disabled: {err:#}")); } - let federation = crate::federation::Federation::new(Arc::clone(&library)); + let devices = crate::devices::DeviceSync::new(Arc::clone(&library))?; + devices.set_event_tx(event_tx.clone()); + let federation = crate::federation::Federation::new(Arc::clone(&library), Arc::clone(&devices)); state.federation.settings = federation.settings(); + state.federation.devices = Some(devices.status()); let player_events = event_tx.clone(); let mut runtime = Runtime { event_tx, library, + devices, federation, fed_status_at: None, fed_resolving: std::sync::Mutex::new(std::collections::HashSet::new()), @@ -114,11 +119,13 @@ pub async fn run( { let fed = Arc::clone(&runtime.federation); + let devices = Arc::clone(&runtime.devices); let tx = runtime.event_tx.clone(); tokio::spawn(async move { fed.start_if_enabled().await; let status = fed.status().await; let _ = tx.send(AppEvent::FederationStatus(status)); + let _ = tx.send(AppEvent::DeviceSyncStatus(devices.status())); }); } @@ -430,11 +437,15 @@ fn perform_effect(state: &mut AppState, runtime: &mut Runtime, effect: Effect) { return; } let library = Arc::clone(&runtime.library); + let devices = Arc::clone(&runtime.devices); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || { for track_id in track_ids { match library.toggle_like(track_id) { Ok(liked) => { + if let Err(err) = devices.record_track_like(track_id, liked) { + tracing::warn!(%err, track_id, "recording synced like failed"); + } let _ = tx.send(AppEvent::LikeToggled { track_id, liked }); } Err(err) => { @@ -448,6 +459,11 @@ fn perform_effect(state: &mut AppState, runtime: &mut Runtime, effect: Effect) { for fed in fed_tracks { match library.toggle_fed_like(&fed) { Ok(liked) => { + if let Err(err) = + devices.record_fed_like(fed.content_id.as_deref(), liked) + { + tracing::warn!(%err, title = %fed.title, "recording synced federated like failed"); + } let _ = tx.send(AppEvent::FedLikeToggled { item_id: fed.item_id.clone(), liked, @@ -468,12 +484,20 @@ fn perform_effect(state: &mut AppState, runtime: &mut Runtime, effect: Effect) { track_ids, } => { let library = Arc::clone(&runtime.library); + let devices = Arc::clone(&runtime.devices); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || { let event = match library.remove_tracks_from_playlist(playlist_id, &track_ids) { - Ok(()) => AppEvent::LibraryChanged { - message: Some(format!("removed {} track(s)", track_ids.len())), - }, + Ok(()) => { + if let Err(err) = + devices.record_playlist_tracks_removed(playlist_id, &track_ids) + { + tracing::warn!(%err, playlist_id, "recording synced playlist removal failed"); + } + AppEvent::LibraryChanged { + message: Some(format!("removed {} track(s)", track_ids.len())), + } + } Err(err) => AppEvent::StatusMessage(format!("remove failed: {err:#}")), }; let _ = tx.send(event); @@ -500,6 +524,57 @@ fn perform_effect(state: &mut AppState, runtime: &mut Runtime, effect: Effect) { let _ = tx.send(AppEvent::FedTicket(result)); }); } + Effect::DeviceShowInvite => { + let fed = Arc::clone(&runtime.federation); + let devices = Arc::clone(&runtime.devices); + let tx = runtime.event_tx.clone(); + tokio::spawn(async move { + let result = fed.device_invite().await.map_err(|err| format!("{err:#}")); + let _ = tx.send(AppEvent::DeviceInvite(result)); + let _ = tx.send(AppEvent::FederationStatus(fed.status().await)); + let _ = tx.send(AppEvent::DeviceSyncStatus(devices.status())); + }); + } + Effect::DeviceConnectInvite(invite) => device_connect(runtime, invite), + Effect::DeviceSyncNow => { + let fed = Arc::clone(&runtime.federation); + let devices = Arc::clone(&runtime.devices); + let tx = runtime.event_tx.clone(); + tokio::spawn(async move { + let message = match fed.device_sync_now().await { + Ok(()) => "devices: sync complete".to_string(), + Err(err) => format!("devices: {err:#}"), + }; + let _ = tx.send(AppEvent::DeviceSyncStatus(devices.status())); + let _ = tx.send(AppEvent::StatusMessage(message)); + }); + } + Effect::DeviceSetName(name) => { + let fed = Arc::clone(&runtime.federation); + let devices = Arc::clone(&runtime.devices); + let tx = runtime.event_tx.clone(); + tokio::spawn(async move { + let ticket = fed.ticket().await.ok(); + let message = match devices.set_device_name(&name, ticket.as_deref()) { + Ok(()) => "device name saved".to_string(), + Err(err) => format!("device name: {err:#}"), + }; + let _ = tx.send(AppEvent::DeviceSyncStatus(devices.status())); + let _ = tx.send(AppEvent::StatusMessage(message)); + }); + } + Effect::DeviceRevoke(device_id) => { + let devices = Arc::clone(&runtime.devices); + let tx = runtime.event_tx.clone(); + tokio::task::spawn_blocking(move || { + let message = match devices.revoke_device(&device_id) { + Ok(()) => format!("device {} revoked", &device_id[..device_id.len().min(10)]), + Err(err) => format!("revoke failed: {err:#}"), + }; + let _ = tx.send(AppEvent::DeviceSyncStatus(devices.status())); + let _ = tx.send(AppEvent::StatusMessage(message)); + }); + } Effect::FedOpenArtist(name) => { let fed = Arc::clone(&runtime.federation); let tx = runtime.event_tx.clone(); @@ -646,6 +721,45 @@ fn shell_quote(path: &Path) -> String { format!("'{}'", value.replace('\'', "'\\''")) } +fn copy_text_to_clipboard(text: &str) -> bool { + let command: &[&str] = if cfg!(target_os = "macos") { + &["pbcopy"] + } else if cfg!(target_os = "windows") { + &["clip"] + } else { + &["wl-copy"] + }; + let Some((program, args)) = command.split_first() else { + return false; + }; + let mut child = match Command::new(program) + .args(args) + .stdin(Stdio::piped()) + .spawn() + { + Ok(child) => child, + Err(_) if !cfg!(target_os = "macos") && !cfg!(target_os = "windows") => { + match Command::new("xclip") + .args(["-selection", "clipboard"]) + .stdin(Stdio::piped()) + .spawn() + { + Ok(child) => child, + Err(_) => return false, + } + } + Err(_) => return false, + }; + let Some(mut stdin) = child.stdin.take() else { + return false; + }; + if stdin.write_all(text.as_bytes()).is_err() { + return false; + } + drop(stdin); + child.wait().is_ok_and(|status| status.success()) +} + /// Start playing `queue[queue_pos]`: open the local file in a background /// task and hand the reader to the audio thread. fn play_current(state: &mut AppState, runtime: &mut Runtime) { @@ -819,6 +933,22 @@ pub(crate) fn fed_connect(runtime: &Runtime, ticket: String) { }); } +/// Pair with another trusted client by an opaque `frid://i/...` invite. +pub(crate) fn device_connect(runtime: &Runtime, invite: String) { + let fed = Arc::clone(&runtime.federation); + let devices = Arc::clone(&runtime.devices); + let tx = runtime.event_tx.clone(); + tokio::spawn(async move { + let result = fed + .device_connect(&invite) + .await + .map_err(|err| format!("{err:#}")); + let _ = tx.send(AppEvent::FederationStatus(fed.status().await)); + let _ = tx.send(AppEvent::DeviceSyncStatus(devices.status())); + let _ = tx.send(AppEvent::DeviceConnectResult(result)); + }); +} + /// Downloads one pending federated track (into the cache, or the library /// when save-on-listen is enabled) and reports back with the placeholder id /// so the queue can swap the resolved track in. @@ -863,6 +993,7 @@ pub(crate) fn fed_download_spawn( } let fed = Arc::clone(&runtime.federation); let library = Arc::clone(&runtime.library); + let devices = Arc::clone(&runtime.devices); let tx = runtime.event_tx.clone(); tokio::spawn(async move { let total = tracks.len(); @@ -890,12 +1021,19 @@ pub(crate) fn fed_download_spawn( && !imported_ids.is_empty() { let library = Arc::clone(&library); + let devices = Arc::clone(&devices); let tx_add = tx.clone(); let title = playlist_title.clone(); tokio::task::spawn_blocking(move || { let result = library .add_tracks_to_playlist(playlist_id, &imported_ids) .map_err(|err| format!("{err:#}")); + if result.is_ok() + && let Err(err) = + devices.record_playlist_tracks_added(playlist_id, &imported_ids) + { + tracing::warn!(%err, playlist_id, "recording synced playlist add failed"); + } let _ = tx_add.send(AppEvent::PlaylistTracksAdded { playlist_id, playlist_title: title, @@ -912,9 +1050,11 @@ pub(crate) fn fed_download_spawn( /// Request a fresh status snapshot for the Federation tab. fn fed_spawn_status(runtime: &Runtime) { let fed = Arc::clone(&runtime.federation); + let devices = Arc::clone(&runtime.devices); let tx = runtime.event_tx.clone(); tokio::spawn(async move { let _ = tx.send(AppEvent::FederationStatus(fed.status().await)); + let _ = tx.send(AppEvent::DeviceSyncStatus(devices.status())); }); } @@ -1221,6 +1361,38 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent AppEvent::FederationStatus(status) => { state.federation.status = Some(status); } + AppEvent::DeviceSyncStatus(status) => { + state.federation.devices = Some(status); + clamp_settings_cursor(state); + } + AppEvent::DeviceInvite(result) => match result { + Ok(invite) => { + let copied = copy_text_to_clipboard(&invite); + state.popup = Some(state::Popup::FedText { + title: "Device invite".to_string(), + text: invite, + }); + state.status_message = Some(if copied { + "device invite copied to clipboard".to_string() + } else { + "device invite generated".to_string() + }); + } + Err(message) => state.status_message = Some(format!("device invite: {message}")), + }, + AppEvent::DeviceConnectResult(result) => match result { + Ok(message) => state.status_message = Some(message), + Err(message) => state.status_message = Some(format!("connect failed: {message}")), + }, + AppEvent::DevicePairingRequest(request) => { + state.popup = Some(state::Popup::DevicePairing { + request_id: request.request_id, + device_id: request.device_id, + name: request.name, + client_version: request.client_version, + }); + state.federation.devices = Some(runtime.devices.status()); + } AppEvent::FedSearchLoaded { seq, result } => { if runtime.search_seq.load(std::sync::atomic::Ordering::SeqCst) != seq { return; diff --git a/src/app/popup.rs b/src/app/popup.rs index 04dda3c..88cef61 100644 --- a/src/app/popup.rs +++ b/src/app/popup.rs @@ -61,6 +61,23 @@ pub fn handle_key(state: &mut AppState, runtime: &mut Runtime, key: KeyEvent) { KeyCode::Esc | KeyCode::Enter | KeyCode::Char('q') => {} _ => state.popup = Some(Popup::FedText { title, text }), }, + Popup::DevicePairing { + request_id, + device_id, + name, + client_version, + } => handle_device_pairing( + state, + runtime, + request_id, + device_id, + name, + client_version, + key, + ), + Popup::ConfirmDeviceRevoke { device_id, name } => { + handle_device_revoke(state, runtime, device_id, name, key); + } } } @@ -84,7 +101,7 @@ fn handle_library_filters(state: &mut AppState, runtime: &Runtime, cursor: usize /// One-line text entry on the Federation tab (network id / peer ticket). fn handle_fed_input( state: &mut AppState, - runtime: &Runtime, + runtime: &mut Runtime, field: FedInputField, mut input: crate::app::input::LineEdit, key: KeyEvent, @@ -110,6 +127,28 @@ fn handle_fed_input( super::fed_connect(runtime, value); } } + FedInputField::DeviceName => { + if value.is_empty() { + state.status_message = Some("device name is empty".into()); + } else { + super::perform_effect( + state, + runtime, + crate::app::update::Effect::DeviceSetName(value), + ); + } + } + FedInputField::ConnectInvite => { + if value.is_empty() { + state.status_message = Some("invite is empty".into()); + } else { + super::perform_effect( + state, + runtime, + crate::app::update::Effect::DeviceConnectInvite(value), + ); + } + } } } _ => { @@ -119,6 +158,63 @@ fn handle_fed_input( } } +fn handle_device_pairing( + state: &mut AppState, + runtime: &Runtime, + request_id: String, + device_id: String, + name: String, + client_version: String, + key: KeyEvent, +) { + match key.code { + KeyCode::Esc | KeyCode::Char('n') | KeyCode::Char('q') => { + if let Err(err) = runtime.devices.answer_pairing(&request_id, false) { + state.status_message = Some(format!("pairing: {err:#}")); + } else { + state.status_message = Some("device pairing denied".to_string()); + } + state.federation.devices = Some(runtime.devices.status()); + } + KeyCode::Enter | KeyCode::Char('y') => { + if let Err(err) = runtime.devices.answer_pairing(&request_id, true) { + state.status_message = Some(format!("pairing: {err:#}")); + } else { + state.status_message = Some(format!("device \"{name}\" accepted")); + } + state.federation.devices = Some(runtime.devices.status()); + } + _ => { + state.popup = Some(Popup::DevicePairing { + request_id, + device_id, + name, + client_version, + }); + } + } +} + +fn handle_device_revoke( + state: &mut AppState, + runtime: &mut Runtime, + device_id: String, + name: String, + key: KeyEvent, +) { + match key.code { + KeyCode::Esc | KeyCode::Char('n') | KeyCode::Char('q') => {} + KeyCode::Enter | KeyCode::Char('y') => { + super::perform_effect( + state, + runtime, + crate::app::update::Effect::DeviceRevoke(device_id), + ); + } + _ => state.popup = Some(Popup::ConfirmDeviceRevoke { device_id, name }), + } +} + /// Pasted text goes into the focused text field when one is open. pub fn handle_paste(state: &mut AppState, pasted: &str) { let cleaned: String = pasted.chars().filter(|c| !c.is_control()).collect(); @@ -278,7 +374,13 @@ fn save_edit( if title.is_empty() { return Err("title is empty".to_string()); } - library.update_playlist(id, &title, None) + let result = library.update_playlist(id, &title, None); + if result.is_ok() + && let Err(err) = runtime.devices.record_playlist_renamed(id, &title) + { + tracing::warn!(%err, playlist = id, "recording synced playlist rename failed"); + } + result } }; match result { @@ -311,7 +413,13 @@ fn handle_confirm_delete( DeleteTarget::Track(id) => library.delete_track(id), DeleteTarget::Release(id) => library.delete_release(id), DeleteTarget::Artist(id) => library.delete_artist(id), - DeleteTarget::Playlist(id) => library.delete_playlist(id), + DeleteTarget::Playlist(id) => match runtime.devices.record_playlist_deleted(id) { + Ok(()) => library.delete_playlist(id), + Err(err) => { + tracing::warn!(%err, playlist = id, "recording synced playlist deletion failed"); + library.delete_playlist(id) + } + }, }; match result { Ok(()) => { @@ -656,12 +764,18 @@ pub(crate) fn spawn_add_target( match target { crate::app::state::PlaylistAddTarget::Local(tracks) => { let library = Arc::clone(&runtime.library); + let devices = Arc::clone(&runtime.devices); let tx = runtime.event_tx.clone(); let ids: Vec = tracks.iter().map(|t| t.id).filter(|id| *id >= 0).collect(); tokio::task::spawn_blocking(move || { let result = library .add_tracks_to_playlist(playlist_id, &ids) .map_err(|err| format!("{err:#}")); + if result.is_ok() + && let Err(err) = devices.record_playlist_tracks_added(playlist_id, &ids) + { + tracing::warn!(%err, playlist_id, "recording synced playlist add failed"); + } let _ = tx.send(AppEvent::PlaylistTracksAdded { playlist_id, playlist_title, @@ -681,11 +795,17 @@ fn spawn_create_playlist( add_target: Option, ) { let library = Arc::clone(&runtime.library); + let devices = Arc::clone(&runtime.devices); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || { let result = library .create_playlist(&title) .map_err(|err| format!("{err:#}")); + if let Ok(playlist) = &result + && let Err(err) = devices.record_playlist_created(playlist.id, &playlist.title) + { + tracing::warn!(%err, playlist = playlist.id, "recording synced playlist creation failed"); + } let _ = tx.send(AppEvent::PlaylistCreated { result, add_target }); }); } diff --git a/src/app/state.rs b/src/app/state.rs index d09ad1e..4493fd8 100644 --- a/src/app/state.rs +++ b/src/app/state.rs @@ -530,12 +530,23 @@ pub enum Popup { }, /// Wrapped read-only text (this peer's connection ticket). FedText { title: String, text: String }, + /// Incoming trusted-device pairing request. + DevicePairing { + request_id: String, + device_id: String, + name: String, + client_version: String, + }, + /// Confirmation before revoking a trusted device. + ConfirmDeviceRevoke { device_id: String, name: String }, } #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum FedInputField { NetworkId, ConnectTicket, + DeviceName, + ConnectInvite, } impl FedInputField { @@ -543,6 +554,8 @@ impl FedInputField { match self { FedInputField::NetworkId => "Network ID", FedInputField::ConnectTicket => "Connect to peer (paste ticket)", + FedInputField::DeviceName => "Device name", + FedInputField::ConnectInvite => "Connect device (paste frid://i invite)", } } } @@ -575,6 +588,11 @@ impl FedRow { #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum SettingsRow { Federation(FedRow), + DeviceName, + DeviceInvite, + DeviceConnect, + DeviceSyncNow, + Device(usize), VisualizationClock, VisualizationScript(usize), VisualizationNew, @@ -584,6 +602,13 @@ pub enum SettingsRow { pub fn settings_rows(state: &AppState) -> Vec { let mut rows = Vec::new(); rows.extend(FedRow::ALL.into_iter().map(SettingsRow::Federation)); + rows.push(SettingsRow::DeviceName); + rows.push(SettingsRow::DeviceInvite); + rows.push(SettingsRow::DeviceConnect); + rows.push(SettingsRow::DeviceSyncNow); + if let Some(status) = &state.federation.devices { + rows.extend((0..status.devices.len()).map(SettingsRow::Device)); + } rows.push(SettingsRow::VisualizationClock); rows.extend( state @@ -605,6 +630,7 @@ pub fn settings_rows(state: &AppState) -> Vec { pub struct FederationTab { pub settings: crate::federation::FedSettings, pub status: Option, + pub devices: Option, } /// Playlists eligible as add-targets (the virtual Likes playlist is managed diff --git a/src/app/update.rs b/src/app/update.rs index c136e4c..9ddc4bb 100644 --- a/src/app/update.rs +++ b/src/app/update.rs @@ -49,6 +49,16 @@ pub enum Effect { FedSyncNow, /// Fetch this peer's ticket and show it in a popup. FedShowTicket, + /// Generate a personal-device invite and show it in a popup. + DeviceShowInvite, + /// Connect this client to a trusted device by opaque frid invite. + DeviceConnectInvite(String), + /// Force an immediate personal-device sync. + DeviceSyncNow, + /// Persist and publish this device's display name. + DeviceSetName(String), + /// Revoke a trusted device. + DeviceRevoke(String), /// Assemble the federated artist card (fan-out to the owning peers). FedOpenArtist(String), /// Download federated tracks into the local library, one by one. @@ -2333,6 +2343,48 @@ fn federation_select(state: &mut AppState) -> Option { }); None } + SettingsRow::DeviceName => { + let name = state + .federation + .devices + .as_ref() + .map(|status| status.this_device_name.clone()) + .unwrap_or_default(); + state.popup = Some(Popup::FedInput { + field: FedInputField::DeviceName, + input: crate::app::input::LineEdit::new(name), + }); + None + } + SettingsRow::DeviceInvite => Some(Effect::DeviceShowInvite), + SettingsRow::DeviceConnect => { + state.popup = Some(Popup::FedInput { + field: FedInputField::ConnectInvite, + input: crate::app::input::LineEdit::default(), + }); + None + } + SettingsRow::DeviceSyncNow => Some(Effect::DeviceSyncNow), + SettingsRow::Device(index) => { + let Some(device) = state + .federation + .devices + .as_ref() + .and_then(|status| status.devices.get(index)) + .cloned() + else { + return None; + }; + if device.is_self || device.revoked { + state.status_message = Some("this device cannot be revoked here".to_string()); + return None; + } + state.popup = Some(Popup::ConfirmDeviceRevoke { + device_id: device.device_id, + name: device.name, + }); + None + } SettingsRow::VisualizationClock => { match state.visualizer.toggle_clock() { Ok(()) => { diff --git a/src/devices.rs b/src/devices.rs new file mode 100644 index 0000000..649aab0 --- /dev/null +++ b/src/devices.rs @@ -0,0 +1,2212 @@ +//! Personal-device sync: trusted furumi clients that belong to one listener. +//! +//! The transport is a small JSON-lines protocol over a dedicated frid/iroh +//! byte-stream ALPN. Local state is stored as an append-only operation log plus +//! materialized tables, so offline clients can merge likes, playlists and +//! membership changes deterministically. + +use std::collections::BTreeMap; +use std::path::PathBuf; +use std::str::FromStr; +use std::sync::Arc; +use std::time::Duration; + +use anyhow::{Context as _, Result}; +use music_dht::{ByteStream, MusicDhtService, NetworkId, PeerTicket, SecretKey, StreamAcceptor}; +use rusqlite::{Connection, OptionalExtension, params}; +use serde::{Deserialize, Serialize}; +use tokio::io::{AsyncRead, AsyncReadExt}; + +use crate::app::event::AppEvent; +use crate::library::Library; + +pub const SYNC_ALPN: &[u8] = b"furumi/sync/1"; +const CLIENT_VERSION: &str = env!("CARGO_PKG_VERSION"); +const PROTOCOL_VERSION: u16 = 1; +const INVITE_TTL_MS: i64 = 10 * 60 * 1000; +const PAIRING_WAIT_MS: i64 = 120 * 1000; +const SYNC_INTERVAL: Duration = Duration::from_secs(30); +const MAX_LINE: usize = 8 * 1024 * 1024; +const MAX_OPS_PER_BATCH: usize = 1000; + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct PendingPairing { + pub request_id: String, + pub device_id: String, + pub name: String, + pub client_version: String, +} + +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct DeviceStatusRow { + pub device_id: String, + pub name: String, + pub client_version: String, + pub endpoint_id: String, + pub last_seen_ms: Option, + pub revoked: bool, + pub is_self: bool, +} + +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct DeviceSyncStatus { + pub this_device_id: String, + pub this_device_name: String, + pub group_id: String, + pub active_devices: usize, + pub revoked_devices: usize, + pub pending_requests: usize, + pub ops_total: usize, + pub tombstone_ops: usize, + pub compactable_tombstones: usize, + pub outbox_ops: usize, + pub snapshot_likes: usize, + pub snapshot_playlists: usize, + pub snapshot_items: usize, + pub unresolved_playlist_items: usize, + pub peer_ack_floor: String, + pub last_sync: Option, + pub last_error: Option, + pub devices: Vec, +} + +#[derive(Clone)] +pub struct DeviceSync { + conn: Arc>, + library: Arc, + event_tx: Arc>>>, +} + +#[derive(Debug, Clone)] +struct LocalIdentity { + device_id: String, + group_id: String, + name: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct InviteWire { + v: u16, + #[serde(rename = "t")] + ticket: String, + #[serde(rename = "d")] + device_id: String, + #[serde(rename = "i")] + invite_id: String, + #[serde(rename = "s")] + secret: String, + #[serde(rename = "e")] + expires_at_ms: i64, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct DeviceProfileWire { + device_id: String, + name: String, + client_version: String, + protocol_version: u16, + endpoint_id: String, + endpoint_ticket: String, + #[serde(default)] + revoked: bool, + #[serde(default)] + revoke_cutoff_seq: Option, + #[serde(default)] + updated_at_ms: i64, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct SyncOpWire { + op_id: String, + origin_device_id: String, + seq: i64, + hlc_ms: i64, + payload: SyncOpPayload, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(tag = "kind", rename_all = "snake_case")] +pub enum SyncOpPayload { + TrackLikeSet { + content_id: String, + liked: bool, + }, + PlaylistCreated { + playlist_id: String, + title: String, + }, + PlaylistRenamed { + playlist_id: String, + title: String, + }, + PlaylistDeleted { + playlist_id: String, + }, + PlaylistTrackAdded { + playlist_id: String, + content_id: String, + position: i64, + }, + PlaylistTrackRemoved { + playlist_id: String, + content_id: String, + }, + DeviceProfileSet { + name: String, + client_version: String, + endpoint_ticket: String, + endpoint_id: String, + }, + DeviceRevoked { + target_device_id: String, + target_max_seq_seen: i64, + }, +} + +impl SyncOpPayload { + fn is_tombstone(&self) -> bool { + matches!( + self, + SyncOpPayload::TrackLikeSet { liked: false, .. } + | SyncOpPayload::PlaylistDeleted { .. } + | SyncOpPayload::PlaylistTrackRemoved { .. } + | SyncOpPayload::DeviceRevoked { .. } + ) + } +} + +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +struct SyncSnapshot { + #[serde(default)] + likes: Vec, + #[serde(default)] + playlists: Vec, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct SnapshotLike { + content_id: String, + hlc_ms: i64, + op_id: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct SnapshotPlaylist { + playlist_id: String, + title: String, + hlc_ms: i64, + op_id: String, + #[serde(default)] + items: Vec, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct SnapshotPlaylistItem { + content_id: String, + position: i64, + hlc_ms: i64, + op_id: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(tag = "type", rename_all = "snake_case")] +enum WireMessage { + PairRequest { + invite_id: String, + secret: String, + profile: DeviceProfileWire, + vector: BTreeMap, + ops: Vec, + snapshot: SyncSnapshot, + }, + PairResponse { + accepted: bool, + #[serde(default)] + error: Option, + #[serde(default)] + group_id: Option, + #[serde(default)] + profile: Option, + #[serde(default)] + devices: Vec, + #[serde(default)] + vector: BTreeMap, + #[serde(default)] + ops: Vec, + #[serde(default)] + snapshot: SyncSnapshot, + }, + Hello { + group_id: String, + profile: DeviceProfileWire, + devices: Vec, + vector: BTreeMap, + ops: Vec, + snapshot: SyncSnapshot, + }, + SyncResponse { + accepted: bool, + #[serde(default)] + error: Option, + #[serde(default)] + devices: Vec, + #[serde(default)] + vector: BTreeMap, + #[serde(default)] + ops: Vec, + #[serde(default)] + snapshot: SyncSnapshot, + }, +} + +fn lock(mutex: &std::sync::Mutex) -> std::sync::MutexGuard<'_, T> { + mutex + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) +} + +fn now_ms() -> i64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|d| d.as_millis() as i64) + .unwrap_or(0) +} + +fn now_label() -> String { + let secs = (now_ms() / 1000).max(0); + format!( + "{:02}:{:02}:{:02} UTC", + secs / 3600 % 24, + secs / 60 % 60, + secs % 60 + ) +} + +fn default_db_path() -> PathBuf { + crate::config::project_dirs() + .map(|dirs| dirs.data_dir().join("devices").join("sync.sqlite3")) + .unwrap_or_else(|| PathBuf::from("devices").join("sync.sqlite3")) +} + +impl DeviceSync { + pub fn new(library: Arc) -> Result> { + let path = default_db_path(); + if let Some(parent) = path.parent() { + std::fs::create_dir_all(parent)?; + } + let conn = + Connection::open(&path).with_context(|| format!("opening {}", path.display()))?; + init_schema(&conn)?; + let sync = Arc::new(Self { + conn: Arc::new(std::sync::Mutex::new(conn)), + library, + event_tx: Arc::new(std::sync::Mutex::new(None)), + }); + sync.ensure_identity()?; + Ok(sync) + } + + pub fn set_event_tx(&self, tx: tokio::sync::mpsc::UnboundedSender) { + *lock(&self.event_tx) = Some(tx); + } + + pub fn set_device_name(&self, name: &str, endpoint_ticket: Option<&str>) -> Result<()> { + let name = if name.trim().is_empty() { + "furumi".to_string() + } else { + name.trim().to_string() + }; + { + let conn = lock(&self.conn); + set_meta(&conn, "device_name", &name)?; + if let Some(device_id) = get_meta(&conn, "device_id")? { + conn.execute( + "UPDATE sync_devices + SET name = ?2, client_version = ?3 + WHERE device_id = ?1", + params![device_id, name, CLIENT_VERSION], + )?; + } + } + if let Some(ticket) = endpoint_ticket { + let endpoint_id = ticket_endpoint_id(ticket).unwrap_or_default(); + self.record_local_op(SyncOpPayload::DeviceProfileSet { + name, + client_version: CLIENT_VERSION.to_string(), + endpoint_ticket: ticket.to_string(), + endpoint_id, + })?; + } + Ok(()) + } + + pub fn status(&self) -> DeviceSyncStatus { + match self.status_inner() { + Ok(status) => status, + Err(err) => DeviceSyncStatus { + last_error: Some(format!("{err:#}")), + ..DeviceSyncStatus::default() + }, + } + } + + pub async fn create_invite(&self, service: Arc) -> Result { + let identity = self.ensure_identity()?; + let ticket = service.ticket().await?.to_string(); + let secret = random_hex(16); + let invite_id = random_hex(8); + let expires_at_ms = now_ms() + INVITE_TTL_MS; + { + let conn = lock(&self.conn); + conn.execute( + "INSERT INTO sync_invites (invite_id, secret_hash, expires_at_ms, created_at_ms) + VALUES (?1, ?2, ?3, ?4)", + params![invite_id, hash_secret(&secret), expires_at_ms, now_ms()], + )?; + } + let payload = InviteWire { + v: 1, + ticket, + device_id: identity.device_id, + invite_id, + secret, + expires_at_ms, + }; + let bytes = serde_json::to_vec(&payload)?; + Ok(format!("frid://i/{}", base64url_encode(&bytes))) + } + + pub async fn connect_invite( + &self, + service: Arc, + invite_link: &str, + ) -> Result { + let invite = parse_invite(invite_link)?; + anyhow::ensure!(invite.expires_at_ms >= now_ms(), "invite expired"); + let ticket: PeerTicket = invite + .ticket + .parse() + .map_err(|err| anyhow::anyhow!("malformed invite ticket: {err}"))?; + let peer = service.connect(ticket).await?; + let own_ticket = service.ticket().await?.to_string(); + let profile = self.own_profile(&own_ticket)?; + let vector = self.vector()?; + let ops = self.ops_for_peer(&invite.device_id)?; + let snapshot = self.snapshot()?; + let mut stream = service.open_stream(peer, SYNC_ALPN).await?; + write_msg( + &mut stream, + &WireMessage::PairRequest { + invite_id: invite.invite_id, + secret: invite.secret, + profile, + vector, + ops, + snapshot, + }, + ) + .await?; + let response = read_msg(&mut stream).await?; + match response { + WireMessage::PairResponse { + accepted: true, + group_id: Some(group_id), + profile, + devices, + vector, + ops, + snapshot, + .. + } => { + self.set_group_id(&group_id)?; + if let Some(profile) = profile { + self.apply_device_profile(&profile, false)?; + } + self.apply_device_profiles(&devices)?; + self.apply_snapshot(snapshot)?; + self.apply_ops(ops)?; + self.note_peer_vector(&invite.device_id, &vector)?; + self.set_last_sync(Some(format!("paired with {}", short_id(&invite.device_id))))?; + self.gc_tombstones()?; + Ok(format!("connected device {}", short_id(&invite.device_id))) + } + WireMessage::PairResponse { + accepted: false, + error, + .. + } => anyhow::bail!(error.unwrap_or_else(|| "pairing denied".to_string())), + _ => anyhow::bail!("unexpected pairing response"), + } + } + + pub fn answer_pairing(&self, request_id: &str, accept: bool) -> Result<()> { + let conn = lock(&self.conn); + conn.execute( + "UPDATE sync_pending_pairing + SET status = ?2, answered_at_ms = ?3 + WHERE request_id = ?1 AND status = 'pending'", + params![ + request_id, + if accept { "accepted" } else { "denied" }, + now_ms() + ], + )?; + Ok(()) + } + + pub fn revoke_device(&self, device_id: &str) -> Result<()> { + let own = self.ensure_identity()?.device_id; + anyhow::ensure!(device_id != own, "cannot revoke this device from itself"); + let cutoff = { + let conn = lock(&self.conn); + conn.query_row( + "SELECT COALESCE(MAX(seq), 0) FROM sync_ops WHERE origin_device_id = ?1", + [device_id], + |row| row.get::<_, i64>(0), + )? + }; + { + let conn = lock(&self.conn); + conn.execute( + "UPDATE sync_devices + SET revoked_at_ms = ?2, revoked_by = ?3, revoke_cutoff_seq = ?4 + WHERE device_id = ?1", + params![device_id, now_ms(), own, cutoff], + )?; + } + self.record_local_op(SyncOpPayload::DeviceRevoked { + target_device_id: device_id.to_string(), + target_max_seq_seen: cutoff, + })?; + self.gc_tombstones()?; + Ok(()) + } + + pub fn record_track_like(&self, track_id: i64, liked: bool) -> Result<()> { + if let Some(content_id) = self.library.track_content_id_by_id(track_id)? { + self.record_local_op(SyncOpPayload::TrackLikeSet { content_id, liked })?; + } + Ok(()) + } + + pub fn record_fed_like(&self, content_id: Option<&str>, liked: bool) -> Result<()> { + let Some(content_id) = content_id.and_then(music_dht::normalize_content_id) else { + return Ok(()); + }; + self.record_local_op(SyncOpPayload::TrackLikeSet { content_id, liked })?; + Ok(()) + } + + pub fn record_playlist_created(&self, playlist_id: i64, title: &str) -> Result<()> { + let playlist_id = self.library.ensure_playlist_sync_id(playlist_id)?; + self.record_local_op(SyncOpPayload::PlaylistCreated { + playlist_id, + title: title.to_string(), + })?; + Ok(()) + } + + pub fn record_playlist_renamed(&self, playlist_id: i64, title: &str) -> Result<()> { + let playlist_id = self.library.ensure_playlist_sync_id(playlist_id)?; + self.record_local_op(SyncOpPayload::PlaylistRenamed { + playlist_id, + title: title.to_string(), + })?; + Ok(()) + } + + pub fn record_playlist_deleted(&self, playlist_id: i64) -> Result<()> { + let Some(playlist_id) = self.library.playlist_sync_id(playlist_id)? else { + return Ok(()); + }; + self.record_local_op(SyncOpPayload::PlaylistDeleted { playlist_id })?; + Ok(()) + } + + pub fn record_playlist_tracks_added(&self, playlist_id: i64, track_ids: &[i64]) -> Result<()> { + let playlist_id = self.library.ensure_playlist_sync_id(playlist_id)?; + for (position, content_id) in self + .library + .track_content_ids(track_ids)? + .into_iter() + .enumerate() + { + self.record_local_op(SyncOpPayload::PlaylistTrackAdded { + playlist_id: playlist_id.clone(), + content_id, + position: position as i64, + })?; + } + Ok(()) + } + + pub fn record_playlist_tracks_removed( + &self, + playlist_id: i64, + track_ids: &[i64], + ) -> Result<()> { + let Some(playlist_id) = self.library.playlist_sync_id(playlist_id)? else { + return Ok(()); + }; + for content_id in self.library.track_content_ids(track_ids)? { + self.record_local_op(SyncOpPayload::PlaylistTrackRemoved { + playlist_id: playlist_id.clone(), + content_id, + })?; + } + Ok(()) + } + + pub async fn sync_once(&self, service: Arc) -> Result<()> { + let devices = self.active_remote_devices()?; + for device in devices { + if device.endpoint_ticket.trim().is_empty() { + continue; + } + if let Err(err) = self.sync_device(Arc::clone(&service), &device).await { + tracing::debug!(device = %device.device_id, "device sync failed: {err:#}"); + self.set_last_error(Some(format!("{}: {err:#}", short_id(&device.device_id))))?; + } + } + self.gc_tombstones()?; + Ok(()) + } + + async fn sync_device( + &self, + service: Arc, + device: &StoredDevice, + ) -> Result<()> { + let ticket: PeerTicket = device.endpoint_ticket.parse()?; + let peer = service.connect(ticket).await?; + let own_ticket = service.ticket().await?.to_string(); + let identity = self.ensure_identity()?; + let profile = self.own_profile(&own_ticket)?; + let devices = self.device_profiles()?; + let vector = self.vector()?; + let ops = self.ops_for_peer(&device.device_id)?; + let snapshot = self.snapshot()?; + let mut stream = service.open_stream(peer, SYNC_ALPN).await?; + write_msg( + &mut stream, + &WireMessage::Hello { + group_id: identity.group_id, + profile, + devices, + vector, + ops, + snapshot, + }, + ) + .await?; + match read_msg(&mut stream).await? { + WireMessage::SyncResponse { + accepted: true, + devices, + vector, + ops, + snapshot, + .. + } => { + self.apply_device_profiles(&devices)?; + self.apply_snapshot(snapshot)?; + self.apply_ops(ops)?; + self.note_peer_vector(&device.device_id, &vector)?; + self.mark_seen(&device.device_id, Some(peer.to_string()))?; + self.set_last_sync(Some(format!("synced {}", short_id(&device.device_id))))?; + self.set_last_error(None)?; + Ok(()) + } + WireMessage::SyncResponse { + accepted: false, + error, + .. + } => anyhow::bail!(error.unwrap_or_else(|| "sync refused".to_string())), + _ => anyhow::bail!("unexpected sync response"), + } + } + + fn ensure_identity(&self) -> Result { + let conn = lock(&self.conn); + if let (Some(device_id), Some(group_id), Some(name)) = ( + get_meta(&conn, "device_id")?, + get_meta(&conn, "group_id")?, + get_meta(&conn, "device_name")?, + ) { + return Ok(LocalIdentity { + device_id, + group_id, + name, + }); + } + + let key = SecretKey::generate(); + let secret_hex = hex_encode(&key.to_bytes()); + let device_id = format!( + "dev_{}", + &blake3::hash(secret_hex.as_bytes()).to_hex()[..24] + ); + let group_id = format!("grp_{}", &blake3::hash(device_id.as_bytes()).to_hex()[..24]); + let name = default_device_name(&device_id); + set_meta(&conn, "device_secret", &secret_hex)?; + set_meta(&conn, "device_id", &device_id)?; + set_meta(&conn, "group_id", &group_id)?; + set_meta(&conn, "device_name", &name)?; + set_meta(&conn, "local_seq", "0")?; + set_meta(&conn, "last_hlc_ms", "0")?; + conn.execute( + "INSERT OR IGNORE INTO sync_devices + (device_id, name, client_version, protocol_version, endpoint_id, + endpoint_ticket, trusted_at_ms, last_seen_ms) + VALUES (?1, ?2, ?3, ?4, '', '', ?5, ?5)", + params![device_id, name, CLIENT_VERSION, PROTOCOL_VERSION, now_ms()], + )?; + Ok(LocalIdentity { + device_id, + group_id, + name, + }) + } + + fn set_group_id(&self, group_id: &str) -> Result<()> { + let conn = lock(&self.conn); + set_meta(&conn, "group_id", group_id) + } + + fn own_profile(&self, endpoint_ticket: &str) -> Result { + let identity = self.ensure_identity()?; + let endpoint_id = ticket_endpoint_id(endpoint_ticket).unwrap_or_default(); + { + let conn = lock(&self.conn); + conn.execute( + "INSERT INTO sync_devices + (device_id, name, client_version, protocol_version, endpoint_id, + endpoint_ticket, trusted_at_ms, last_seen_ms) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?7) + ON CONFLICT(device_id) DO UPDATE SET + name = excluded.name, + client_version = excluded.client_version, + protocol_version = excluded.protocol_version, + endpoint_id = excluded.endpoint_id, + endpoint_ticket = excluded.endpoint_ticket, + last_seen_ms = excluded.last_seen_ms", + params![ + identity.device_id, + identity.name, + CLIENT_VERSION, + PROTOCOL_VERSION, + endpoint_id, + endpoint_ticket, + now_ms(), + ], + )?; + } + Ok(DeviceProfileWire { + device_id: identity.device_id, + name: identity.name, + client_version: CLIENT_VERSION.to_string(), + protocol_version: PROTOCOL_VERSION, + endpoint_id, + endpoint_ticket: endpoint_ticket.to_string(), + revoked: false, + revoke_cutoff_seq: None, + updated_at_ms: now_ms(), + }) + } + + fn record_local_op(&self, payload: SyncOpPayload) -> Result<()> { + let identity = self.ensure_identity()?; + let (op, payload_json, tombstone) = { + let conn = lock(&self.conn); + let seq = get_meta(&conn, "local_seq")? + .and_then(|v| v.parse::().ok()) + .unwrap_or(0) + + 1; + let last_hlc = get_meta(&conn, "last_hlc_ms")? + .and_then(|v| v.parse::().ok()) + .unwrap_or(0); + let hlc_ms = now_ms().max(last_hlc + 1); + let op_id = format!("{}:{seq}", identity.device_id); + let payload_json = serde_json::to_string(&payload)?; + let tombstone = payload.is_tombstone(); + conn.execute( + "INSERT OR IGNORE INTO sync_ops + (op_id, origin_device_id, seq, kind, payload_json, hlc_ms, + received_at_ms, tombstone) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)", + params![ + op_id, + identity.device_id, + seq, + payload_kind(&payload), + payload_json, + hlc_ms, + now_ms(), + i64::from(tombstone), + ], + )?; + set_meta(&conn, "local_seq", &seq.to_string())?; + set_meta(&conn, "last_hlc_ms", &hlc_ms.to_string())?; + conn.execute( + "INSERT INTO sync_vectors (device_id, max_seq) + VALUES (?1, ?2) + ON CONFLICT(device_id) DO UPDATE SET + max_seq = MAX(max_seq, excluded.max_seq)", + params![identity.device_id, seq], + )?; + ( + SyncOpWire { + op_id, + origin_device_id: identity.device_id, + seq, + hlc_ms, + payload, + }, + payload_json, + tombstone, + ) + }; + let _ = self.apply_op(&op)?; + tracing::debug!( + op_id = %op.op_id, + tombstone, + payload = %payload_json, + "recorded personal-sync op" + ); + let _ = self.gc_tombstones(); + Ok(()) + } + + fn apply_ops(&self, ops: Vec) -> Result<()> { + let mut changed = false; + for op in ops { + if !self.should_accept_op(&op)? { + continue; + } + let inserted = { + let conn = lock(&self.conn); + conn.execute( + "INSERT OR IGNORE INTO sync_ops + (op_id, origin_device_id, seq, kind, payload_json, hlc_ms, + received_at_ms, tombstone) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)", + params![ + op.op_id, + op.origin_device_id, + op.seq, + payload_kind(&op.payload), + serde_json::to_string(&op.payload)?, + op.hlc_ms, + now_ms(), + i64::from(op.payload.is_tombstone()), + ], + )? + }; + if inserted > 0 { + changed |= self.apply_op(&op)?; + } + let conn = lock(&self.conn); + conn.execute( + "INSERT INTO sync_vectors (device_id, max_seq) + VALUES (?1, ?2) + ON CONFLICT(device_id) DO UPDATE SET + max_seq = MAX(max_seq, excluded.max_seq)", + params![op.origin_device_id, op.seq], + )?; + } + if changed { + self.notify_library_changed(); + } + Ok(()) + } + + fn should_accept_op(&self, op: &SyncOpWire) -> Result { + let identity = self.ensure_identity()?; + if op.origin_device_id == identity.device_id { + return Ok(true); + } + let conn = lock(&self.conn); + let revoked: Option<(Option,)> = conn + .query_row( + "SELECT revoke_cutoff_seq + FROM sync_devices + WHERE device_id = ?1 AND revoked_at_ms IS NOT NULL", + [&op.origin_device_id], + |row| Ok((row.get(0)?,)), + ) + .optional()?; + if let Some((cutoff,)) = revoked { + return Ok(op.seq <= cutoff.unwrap_or(0)); + } + let known: Option = conn + .query_row( + "SELECT 1 FROM sync_devices + WHERE device_id = ?1 AND trusted_at_ms IS NOT NULL", + [&op.origin_device_id], + |row| row.get(0), + ) + .optional()?; + Ok(known.is_some()) + } + + fn apply_op(&self, op: &SyncOpWire) -> Result { + let changed = match &op.payload { + SyncOpPayload::TrackLikeSet { content_id, liked } => { + self.apply_like_state(content_id, *liked, op.hlc_ms, &op.op_id)? + } + SyncOpPayload::PlaylistCreated { playlist_id, title } => { + self.apply_playlist_state(playlist_id, title, false, op.hlc_ms, &op.op_id)? + } + SyncOpPayload::PlaylistRenamed { playlist_id, title } => { + self.apply_playlist_state(playlist_id, title, false, op.hlc_ms, &op.op_id)? + } + SyncOpPayload::PlaylistDeleted { playlist_id } => { + self.apply_playlist_state(playlist_id, "", true, op.hlc_ms, &op.op_id)? + } + SyncOpPayload::PlaylistTrackAdded { + playlist_id, + content_id, + position, + } => self.apply_playlist_item_state( + playlist_id, + content_id, + true, + *position, + op.hlc_ms, + &op.op_id, + )?, + SyncOpPayload::PlaylistTrackRemoved { + playlist_id, + content_id, + } => self.apply_playlist_item_state( + playlist_id, + content_id, + false, + 0, + op.hlc_ms, + &op.op_id, + )?, + SyncOpPayload::DeviceProfileSet { + name, + client_version, + endpoint_ticket, + endpoint_id, + } => { + let profile = DeviceProfileWire { + device_id: op.origin_device_id.clone(), + name: name.clone(), + client_version: client_version.clone(), + protocol_version: PROTOCOL_VERSION, + endpoint_id: endpoint_id.clone(), + endpoint_ticket: endpoint_ticket.clone(), + revoked: false, + revoke_cutoff_seq: None, + updated_at_ms: op.hlc_ms, + }; + self.apply_device_profile(&profile, false)?; + false + } + SyncOpPayload::DeviceRevoked { + target_device_id, + target_max_seq_seen, + } => { + let conn = lock(&self.conn); + conn.execute( + "UPDATE sync_devices + SET revoked_at_ms = COALESCE(revoked_at_ms, ?2), + revoked_by = ?3, + revoke_cutoff_seq = COALESCE(revoke_cutoff_seq, ?4) + WHERE device_id = ?1", + params![ + target_device_id, + op.hlc_ms, + op.origin_device_id, + target_max_seq_seen, + ], + )?; + true + } + }; + Ok(changed) + } + + fn apply_like_state( + &self, + content_id: &str, + liked: bool, + hlc_ms: i64, + op_id: &str, + ) -> Result { + let apply = { + let conn = lock(&self.conn); + let current: Option<(i64, String)> = conn + .query_row( + "SELECT hlc_ms, op_id FROM sync_state_likes WHERE content_id = ?1", + [content_id], + |row| Ok((row.get(0)?, row.get(1)?)), + ) + .optional()?; + current.as_ref().is_none_or(|(current_hlc, current_op)| { + (hlc_ms, op_id) > (*current_hlc, current_op.as_str()) + }) + }; + if !apply { + return Ok(false); + } + { + let conn = lock(&self.conn); + conn.execute( + "INSERT INTO sync_state_likes (content_id, liked, hlc_ms, op_id) + VALUES (?1, ?2, ?3, ?4) + ON CONFLICT(content_id) DO UPDATE SET + liked = excluded.liked, + hlc_ms = excluded.hlc_ms, + op_id = excluded.op_id", + params![content_id, i64::from(liked), hlc_ms, op_id], + )?; + } + if let Some(track_id) = self.library.track_id_by_content_id(content_id)? { + self.library.set_like(track_id, liked)?; + } + Ok(true) + } + + fn apply_playlist_state( + &self, + playlist_id: &str, + title: &str, + deleted: bool, + hlc_ms: i64, + op_id: &str, + ) -> Result { + let apply = { + let conn = lock(&self.conn); + let current: Option<(i64, String)> = conn + .query_row( + "SELECT hlc_ms, op_id FROM sync_state_playlists WHERE playlist_id = ?1", + [playlist_id], + |row| Ok((row.get(0)?, row.get(1)?)), + ) + .optional()?; + current.as_ref().is_none_or(|(current_hlc, current_op)| { + (hlc_ms, op_id) > (*current_hlc, current_op.as_str()) + }) + }; + if !apply { + return Ok(false); + } + { + let conn = lock(&self.conn); + conn.execute( + "INSERT INTO sync_state_playlists + (playlist_id, title, deleted, hlc_ms, op_id) + VALUES (?1, ?2, ?3, ?4, ?5) + ON CONFLICT(playlist_id) DO UPDATE SET + title = excluded.title, + deleted = excluded.deleted, + hlc_ms = excluded.hlc_ms, + op_id = excluded.op_id", + params![playlist_id, title, i64::from(deleted), hlc_ms, op_id], + )?; + } + if deleted { + self.library.delete_playlist_by_sync_id(playlist_id)?; + } else if !title.trim().is_empty() { + self.library.upsert_synced_playlist(playlist_id, title)?; + } + Ok(true) + } + + fn apply_playlist_item_state( + &self, + playlist_id: &str, + content_id: &str, + present: bool, + position: i64, + hlc_ms: i64, + op_id: &str, + ) -> Result { + let apply = { + let conn = lock(&self.conn); + let current: Option<(i64, String)> = conn + .query_row( + "SELECT hlc_ms, op_id + FROM sync_state_playlist_items + WHERE playlist_id = ?1 AND content_id = ?2", + params![playlist_id, content_id], + |row| Ok((row.get(0)?, row.get(1)?)), + ) + .optional()?; + current.as_ref().is_none_or(|(current_hlc, current_op)| { + (hlc_ms, op_id) > (*current_hlc, current_op.as_str()) + }) + }; + if !apply { + return Ok(false); + } + { + let conn = lock(&self.conn); + conn.execute( + "INSERT INTO sync_state_playlist_items + (playlist_id, content_id, present, position, hlc_ms, op_id) + VALUES (?1, ?2, ?3, ?4, ?5, ?6) + ON CONFLICT(playlist_id, content_id) DO UPDATE SET + present = excluded.present, + position = excluded.position, + hlc_ms = excluded.hlc_ms, + op_id = excluded.op_id", + params![ + playlist_id, + content_id, + i64::from(present), + position, + hlc_ms, + op_id + ], + )?; + } + if present { + self.library + .add_content_id_to_synced_playlist(playlist_id, content_id)?; + } else { + self.library + .remove_content_id_from_synced_playlist(playlist_id, content_id)?; + } + Ok(true) + } + + fn apply_snapshot(&self, snapshot: SyncSnapshot) -> Result<()> { + let mut changed = false; + for like in snapshot.likes { + changed |= self.apply_like_state(&like.content_id, true, like.hlc_ms, &like.op_id)?; + } + for playlist in snapshot.playlists { + changed |= self.apply_playlist_state( + &playlist.playlist_id, + &playlist.title, + false, + playlist.hlc_ms, + &playlist.op_id, + )?; + for item in playlist.items { + changed |= self.apply_playlist_item_state( + &playlist.playlist_id, + &item.content_id, + true, + item.position, + item.hlc_ms, + &item.op_id, + )?; + } + } + if changed { + self.notify_library_changed(); + } + Ok(()) + } + + fn apply_device_profiles(&self, devices: &[DeviceProfileWire]) -> Result<()> { + for profile in devices { + self.apply_device_profile(profile, false)?; + } + Ok(()) + } + + fn apply_device_profile(&self, profile: &DeviceProfileWire, trusted: bool) -> Result<()> { + let own = self.ensure_identity()?.device_id; + if profile.device_id == own { + return Ok(()); + } + let conn = lock(&self.conn); + conn.execute( + "INSERT INTO sync_devices + (device_id, name, client_version, protocol_version, endpoint_id, + endpoint_ticket, trusted_at_ms, last_seen_ms, revoked_at_ms, + revoke_cutoff_seq) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10) + ON CONFLICT(device_id) DO UPDATE SET + name = excluded.name, + client_version = excluded.client_version, + protocol_version = excluded.protocol_version, + endpoint_id = excluded.endpoint_id, + endpoint_ticket = CASE + WHEN excluded.endpoint_ticket != '' THEN excluded.endpoint_ticket + ELSE sync_devices.endpoint_ticket + END, + trusted_at_ms = COALESCE(sync_devices.trusted_at_ms, excluded.trusted_at_ms), + last_seen_ms = COALESCE(excluded.last_seen_ms, sync_devices.last_seen_ms), + revoked_at_ms = COALESCE(sync_devices.revoked_at_ms, excluded.revoked_at_ms), + revoke_cutoff_seq = COALESCE(sync_devices.revoke_cutoff_seq, excluded.revoke_cutoff_seq)", + params![ + profile.device_id, + profile.name, + profile.client_version, + profile.protocol_version, + profile.endpoint_id, + profile.endpoint_ticket, + if trusted { now_ms() } else { profile.updated_at_ms }, + profile.updated_at_ms, + if profile.revoked { Some(profile.updated_at_ms) } else { None }, + profile.revoke_cutoff_seq, + ], + )?; + Ok(()) + } + + fn mark_seen(&self, device_id: &str, endpoint_id: Option) -> Result<()> { + let conn = lock(&self.conn); + conn.execute( + "UPDATE sync_devices + SET last_seen_ms = ?2, + endpoint_id = COALESCE(?3, endpoint_id) + WHERE device_id = ?1", + params![device_id, now_ms(), endpoint_id], + )?; + Ok(()) + } + + fn vector(&self) -> Result> { + let conn = lock(&self.conn); + let mut stmt = conn.prepare("SELECT device_id, max_seq FROM sync_vectors")?; + let rows = stmt.query_map([], |row| Ok((row.get(0)?, row.get(1)?)))?; + Ok(rows.collect::>>()?) + } + + fn note_peer_vector(&self, peer_device_id: &str, vector: &BTreeMap) -> Result<()> { + let conn = lock(&self.conn); + for (origin, seq) in vector { + conn.execute( + "INSERT INTO sync_peer_acks + (peer_device_id, origin_device_id, max_seq, updated_at_ms) + VALUES (?1, ?2, ?3, ?4) + ON CONFLICT(peer_device_id, origin_device_id) DO UPDATE SET + max_seq = MAX(max_seq, excluded.max_seq), + updated_at_ms = excluded.updated_at_ms", + params![peer_device_id, origin, seq, now_ms()], + )?; + } + Ok(()) + } + + fn ops_for_peer(&self, peer_device_id: &str) -> Result> { + let conn = lock(&self.conn); + let mut stmt = conn.prepare( + "SELECT o.op_id, o.origin_device_id, o.seq, o.hlc_ms, o.payload_json + FROM sync_ops o + LEFT JOIN sync_peer_acks a + ON a.peer_device_id = ?1 AND a.origin_device_id = o.origin_device_id + WHERE o.seq > COALESCE(a.max_seq, 0) + ORDER BY o.hlc_ms, o.op_id + LIMIT ?2", + )?; + let rows = stmt.query_map(params![peer_device_id, MAX_OPS_PER_BATCH as i64], |row| { + let payload_json: String = row.get(4)?; + Ok(SyncOpWire { + op_id: row.get(0)?, + origin_device_id: row.get(1)?, + seq: row.get(2)?, + hlc_ms: row.get(3)?, + payload: serde_json::from_str(&payload_json).map_err(|err| { + rusqlite::Error::FromSqlConversionFailure( + 4, + rusqlite::types::Type::Text, + Box::new(err), + ) + })?, + }) + })?; + Ok(rows.collect::>>()?) + } + + fn snapshot(&self) -> Result { + let conn = lock(&self.conn); + let mut likes_stmt = conn.prepare( + "SELECT content_id, hlc_ms, op_id + FROM sync_state_likes + WHERE liked = 1 + ORDER BY content_id", + )?; + let likes = likes_stmt + .query_map([], |row| { + Ok(SnapshotLike { + content_id: row.get(0)?, + hlc_ms: row.get(1)?, + op_id: row.get(2)?, + }) + })? + .collect::>>()?; + + let mut playlist_stmt = conn.prepare( + "SELECT playlist_id, title, hlc_ms, op_id + FROM sync_state_playlists + WHERE deleted = 0 + ORDER BY title COLLATE NOCASE", + )?; + let playlist_rows = playlist_stmt + .query_map([], |row| { + Ok(( + row.get::<_, String>(0)?, + row.get::<_, String>(1)?, + row.get::<_, i64>(2)?, + row.get::<_, String>(3)?, + )) + })? + .collect::>>()?; + let mut playlists = Vec::new(); + for (playlist_id, title, hlc_ms, op_id) in playlist_rows { + let mut item_stmt = conn.prepare( + "SELECT content_id, position, hlc_ms, op_id + FROM sync_state_playlist_items + WHERE playlist_id = ?1 AND present = 1 + ORDER BY position, content_id", + )?; + let items = item_stmt + .query_map([&playlist_id], |row| { + Ok(SnapshotPlaylistItem { + content_id: row.get(0)?, + position: row.get(1)?, + hlc_ms: row.get(2)?, + op_id: row.get(3)?, + }) + })? + .collect::>>()?; + playlists.push(SnapshotPlaylist { + playlist_id, + title, + hlc_ms, + op_id, + items, + }); + } + Ok(SyncSnapshot { likes, playlists }) + } + + fn device_profiles(&self) -> Result> { + let conn = lock(&self.conn); + let mut stmt = conn.prepare( + "SELECT device_id, name, client_version, protocol_version, endpoint_id, + endpoint_ticket, revoked_at_ms IS NOT NULL, revoke_cutoff_seq, + COALESCE(last_seen_ms, trusted_at_ms, 0) + FROM sync_devices + WHERE trusted_at_ms IS NOT NULL + ORDER BY device_id", + )?; + let rows = stmt.query_map([], |row| { + Ok(DeviceProfileWire { + device_id: row.get(0)?, + name: row.get(1)?, + client_version: row.get(2)?, + protocol_version: row.get(3)?, + endpoint_id: row.get(4)?, + endpoint_ticket: row.get(5)?, + revoked: row.get::<_, i64>(6)? != 0, + revoke_cutoff_seq: row.get(7)?, + updated_at_ms: row.get(8)?, + }) + })?; + Ok(rows.collect::>>()?) + } + + fn active_remote_devices(&self) -> Result> { + let own = self.ensure_identity()?.device_id; + let conn = lock(&self.conn); + let mut stmt = conn.prepare( + "SELECT device_id, endpoint_ticket + FROM sync_devices + WHERE trusted_at_ms IS NOT NULL + AND revoked_at_ms IS NULL + AND device_id != ?1 + ORDER BY last_seen_ms DESC", + )?; + let rows = stmt.query_map([own], |row| { + Ok(StoredDevice { + device_id: row.get(0)?, + endpoint_ticket: row.get(1)?, + }) + })?; + Ok(rows.collect::>>()?) + } + + fn status_inner(&self) -> Result { + let identity = self.ensure_identity()?; + let conn = lock(&self.conn); + let active_devices: usize = conn.query_row( + "SELECT COUNT(*) FROM sync_devices + WHERE trusted_at_ms IS NOT NULL AND revoked_at_ms IS NULL", + [], + |row| row.get::<_, i64>(0), + )? as usize; + let revoked_devices: usize = conn.query_row( + "SELECT COUNT(*) FROM sync_devices WHERE revoked_at_ms IS NOT NULL", + [], + |row| row.get::<_, i64>(0), + )? as usize; + let pending_requests: usize = conn.query_row( + "SELECT COUNT(*) FROM sync_pending_pairing WHERE status = 'pending'", + [], + |row| row.get::<_, i64>(0), + )? as usize; + let ops_total: usize = conn.query_row("SELECT COUNT(*) FROM sync_ops", [], |row| { + row.get::<_, i64>(0) + })? as usize; + let tombstone_ops: usize = conn.query_row( + "SELECT COUNT(*) FROM sync_ops WHERE tombstone = 1", + [], + |row| row.get::<_, i64>(0), + )? as usize; + let compactable_tombstones = self.compactable_tombstone_count(&conn)?; + let outbox_ops: usize = conn.query_row( + "SELECT COUNT(*) + FROM sync_ops o + WHERE EXISTS ( + SELECT 1 FROM sync_devices d + LEFT JOIN sync_peer_acks a + ON a.peer_device_id = d.device_id + AND a.origin_device_id = o.origin_device_id + WHERE d.trusted_at_ms IS NOT NULL + AND d.revoked_at_ms IS NULL + AND d.device_id != ?1 + AND o.seq > COALESCE(a.max_seq, 0) + )", + [&identity.device_id], + |row| row.get::<_, i64>(0), + )? as usize; + let snapshot_likes: usize = conn.query_row( + "SELECT COUNT(*) FROM sync_state_likes WHERE liked = 1", + [], + |row| row.get::<_, i64>(0), + )? as usize; + let snapshot_playlists: usize = conn.query_row( + "SELECT COUNT(*) FROM sync_state_playlists WHERE deleted = 0", + [], + |row| row.get::<_, i64>(0), + )? as usize; + let snapshot_items: usize = conn.query_row( + "SELECT COUNT(*) FROM sync_state_playlist_items WHERE present = 1", + [], + |row| row.get::<_, i64>(0), + )? as usize; + let unresolved_playlist_items = + self.unresolved_playlist_item_count_with_conn(&conn)? as usize; + let peer_ack_floor = peer_ack_floor_label(&conn)?; + let last_sync = get_meta(&conn, "last_sync")?; + let last_error = get_meta(&conn, "last_error")?; + let mut stmt = conn.prepare( + "SELECT device_id, name, client_version, endpoint_id, last_seen_ms, + revoked_at_ms IS NOT NULL + FROM sync_devices + WHERE trusted_at_ms IS NOT NULL + ORDER BY revoked_at_ms IS NOT NULL, name COLLATE NOCASE, device_id", + )?; + let devices = stmt + .query_map([], |row| { + let device_id: String = row.get(0)?; + Ok(DeviceStatusRow { + is_self: device_id == identity.device_id, + device_id, + name: row.get(1)?, + client_version: row.get(2)?, + endpoint_id: row.get(3)?, + last_seen_ms: row.get(4)?, + revoked: row.get::<_, i64>(5)? != 0, + }) + })? + .collect::>>()?; + Ok(DeviceSyncStatus { + this_device_id: identity.device_id, + this_device_name: identity.name, + group_id: identity.group_id, + active_devices, + revoked_devices, + pending_requests, + ops_total, + tombstone_ops, + compactable_tombstones, + outbox_ops, + snapshot_likes, + snapshot_playlists, + snapshot_items, + unresolved_playlist_items, + peer_ack_floor, + last_sync, + last_error, + devices, + }) + } + + fn compactable_tombstone_count(&self, conn: &Connection) -> Result { + let rows = compactable_tombstone_ids(conn)?; + Ok(rows.len()) + } + + fn unresolved_playlist_item_count_with_conn(&self, conn: &Connection) -> Result { + let mut stmt = conn.prepare( + "SELECT content_id + FROM sync_state_playlist_items + WHERE present = 1", + )?; + let ids = stmt + .query_map([], |row| row.get::<_, String>(0))? + .collect::>>()?; + drop(stmt); + let mut unresolved = 0; + for content_id in ids { + if self.library.track_id_by_content_id(&content_id)?.is_none() { + unresolved += 1; + } + } + Ok(unresolved) + } + + fn set_last_sync(&self, message: Option) -> Result<()> { + let conn = lock(&self.conn); + match message { + Some(message) => set_meta(&conn, "last_sync", &format!("{} · {message}", now_label())), + None => delete_meta(&conn, "last_sync"), + } + } + + fn set_last_error(&self, message: Option) -> Result<()> { + let conn = lock(&self.conn); + match message { + Some(message) => set_meta(&conn, "last_error", &message), + None => delete_meta(&conn, "last_error"), + } + } + + fn notify_library_changed(&self) { + if let Some(tx) = lock(&self.event_tx).as_ref() { + let _ = tx.send(AppEvent::LibraryChanged { message: None }); + let _ = tx.send(AppEvent::DeviceSyncStatus(self.status())); + } + } + + fn gc_tombstones(&self) -> Result<()> { + let conn = lock(&self.conn); + let compactable = compactable_tombstone_ids(&conn)?; + for (op_id, origin, seq) in compactable { + conn.execute( + "INSERT INTO sync_compacted (origin_device_id, max_seq) + VALUES (?1, ?2) + ON CONFLICT(origin_device_id) DO UPDATE SET + max_seq = MAX(max_seq, excluded.max_seq)", + params![origin, seq], + )?; + conn.execute("DELETE FROM sync_ops WHERE op_id = ?1", [op_id])?; + } + conn.execute( + "DELETE FROM sync_state_likes + WHERE liked = 0 + AND EXISTS ( + SELECT 1 FROM sync_compacted c + WHERE c.origin_device_id = substr(sync_state_likes.op_id, 1, instr(sync_state_likes.op_id, ':') - 1) + )", + [], + )?; + Ok(()) + } +} + +#[derive(Debug)] +struct StoredDevice { + device_id: String, + endpoint_ticket: String, +} + +pub async fn serve_peers( + mut acceptor: StreamAcceptor, + sync: Arc, + service: Arc, +) { + while let Some(stream) = acceptor.accept().await { + let sync = Arc::clone(&sync); + let service = Arc::clone(&service); + tokio::spawn(async move { + let peer = stream.peer_id; + if let Err(err) = serve_one(stream, sync, service).await { + tracing::warn!(peer = %peer, "personal sync stream failed: {err:#}"); + } + }); + } +} + +pub async fn sync_loop(sync: Arc, service: Arc) { + let mut interval = tokio::time::interval(SYNC_INTERVAL); + loop { + interval.tick().await; + if let Err(err) = sync.sync_once(Arc::clone(&service)).await { + tracing::debug!("personal sync tick failed: {err:#}"); + } + } +} + +async fn serve_one( + mut stream: ByteStream, + sync: Arc, + service: Arc, +) -> Result<()> { + match read_msg(&mut stream).await? { + WireMessage::PairRequest { + invite_id, + secret, + profile, + vector, + ops, + snapshot, + } => { + handle_pair_request( + stream, sync, service, invite_id, secret, profile, vector, ops, snapshot, + ) + .await + } + WireMessage::Hello { + group_id, + profile, + devices, + vector, + ops, + snapshot, + } => { + handle_hello( + stream, sync, service, group_id, profile, devices, vector, ops, snapshot, + ) + .await + } + _ => anyhow::bail!("unexpected first message"), + } +} + +#[allow(clippy::too_many_arguments)] +async fn handle_pair_request( + mut stream: ByteStream, + sync: Arc, + service: Arc, + invite_id: String, + secret: String, + mut profile: DeviceProfileWire, + vector: BTreeMap, + ops: Vec, + snapshot: SyncSnapshot, +) -> Result<()> { + if !valid_invite(&sync, &invite_id, &secret)? { + tracing::warn!( + peer = %stream.peer_id, + invite_id, + "ignored pairing request with invalid invite secret" + ); + write_msg( + &mut stream, + &WireMessage::PairResponse { + accepted: false, + error: Some("invalid or expired invite".to_string()), + group_id: None, + profile: None, + devices: Vec::new(), + vector: BTreeMap::new(), + ops: Vec::new(), + snapshot: SyncSnapshot::default(), + }, + ) + .await?; + return Ok(()); + } + profile.endpoint_id = stream.peer_id.to_string(); + let request_id = format!("pair_{}", random_hex(8)); + { + let conn = lock(&sync.conn); + conn.execute( + "INSERT INTO sync_pending_pairing + (request_id, device_id, name, client_version, endpoint_id, + endpoint_ticket, invite_id, created_at_ms, status) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, 'pending')", + params![ + request_id, + profile.device_id, + profile.name, + profile.client_version, + profile.endpoint_id, + profile.endpoint_ticket, + invite_id, + now_ms(), + ], + )?; + } + if let Some(tx) = lock(&sync.event_tx).as_ref() { + let _ = tx.send(AppEvent::DevicePairingRequest(PendingPairing { + request_id: request_id.clone(), + device_id: profile.device_id.clone(), + name: profile.name.clone(), + client_version: profile.client_version.clone(), + })); + } + let accepted = wait_pairing_answer(&sync, &request_id).await?; + if !accepted { + write_msg( + &mut stream, + &WireMessage::PairResponse { + accepted: false, + error: Some("pairing denied".to_string()), + group_id: None, + profile: None, + devices: Vec::new(), + vector: BTreeMap::new(), + ops: Vec::new(), + snapshot: SyncSnapshot::default(), + }, + ) + .await?; + return Ok(()); + } + + let own_ticket = service.ticket().await?.to_string(); + let own_profile = sync.own_profile(&own_ticket)?; + let identity = sync.ensure_identity()?; + sync.apply_device_profile(&profile, true)?; + sync.apply_snapshot(snapshot)?; + sync.apply_ops(ops)?; + sync.note_peer_vector(&profile.device_id, &vector)?; + { + let conn = lock(&sync.conn); + conn.execute( + "UPDATE sync_invites SET used_at_ms = ?2 WHERE invite_id = ?1", + params![invite_id, now_ms()], + )?; + } + sync.set_last_sync(Some(format!("paired {}", short_id(&profile.device_id))))?; + let devices = sync.device_profiles()?; + let vector = sync.vector()?; + let ops = sync.ops_for_peer(&profile.device_id)?; + let snapshot = sync.snapshot()?; + write_msg( + &mut stream, + &WireMessage::PairResponse { + accepted: true, + error: None, + group_id: Some(identity.group_id), + profile: Some(own_profile), + devices, + vector, + ops, + snapshot, + }, + ) + .await?; + Ok(()) +} + +#[allow(clippy::too_many_arguments)] +async fn handle_hello( + mut stream: ByteStream, + sync: Arc, + service: Arc, + group_id: String, + mut profile: DeviceProfileWire, + devices: Vec, + vector: BTreeMap, + ops: Vec, + snapshot: SyncSnapshot, +) -> Result<()> { + let identity = sync.ensure_identity()?; + if group_id != identity.group_id { + write_msg( + &mut stream, + &WireMessage::SyncResponse { + accepted: false, + error: Some("sync group mismatch".to_string()), + devices: Vec::new(), + vector: BTreeMap::new(), + ops: Vec::new(), + snapshot: SyncSnapshot::default(), + }, + ) + .await?; + return Ok(()); + } + if !is_active_trusted(&sync, &profile.device_id)? { + write_msg( + &mut stream, + &WireMessage::SyncResponse { + accepted: false, + error: Some("device is not trusted".to_string()), + devices: Vec::new(), + vector: BTreeMap::new(), + ops: Vec::new(), + snapshot: SyncSnapshot::default(), + }, + ) + .await?; + return Ok(()); + } + profile.endpoint_id = stream.peer_id.to_string(); + sync.apply_device_profile(&profile, false)?; + sync.apply_device_profiles(&devices)?; + sync.apply_snapshot(snapshot)?; + sync.apply_ops(ops)?; + sync.note_peer_vector(&profile.device_id, &vector)?; + sync.mark_seen(&profile.device_id, Some(stream.peer_id.to_string()))?; + sync.set_last_sync(Some(format!("synced {}", short_id(&profile.device_id))))?; + + let own_ticket = service.ticket().await?.to_string(); + let own_profile = sync.own_profile(&own_ticket)?; + let devices = sync.device_profiles()?; + let vector = sync.vector()?; + let ops = sync.ops_for_peer(&profile.device_id)?; + let snapshot = sync.snapshot()?; + write_msg( + &mut stream, + &WireMessage::SyncResponse { + accepted: true, + error: None, + devices: { + let mut profiles = devices; + profiles.push(own_profile); + profiles + }, + vector, + ops, + snapshot, + }, + ) + .await?; + sync.gc_tombstones()?; + Ok(()) +} + +async fn wait_pairing_answer(sync: &DeviceSync, request_id: &str) -> Result { + let deadline = now_ms() + PAIRING_WAIT_MS; + loop { + let status: Option = { + let conn = lock(&sync.conn); + conn.query_row( + "SELECT status FROM sync_pending_pairing WHERE request_id = ?1", + [request_id], + |row| row.get(0), + ) + .optional()? + }; + match status.as_deref() { + Some("accepted") => return Ok(true), + Some("denied") => return Ok(false), + _ if now_ms() > deadline => return Ok(false), + _ => tokio::time::sleep(Duration::from_millis(500)).await, + } + } +} + +fn valid_invite(sync: &DeviceSync, invite_id: &str, secret: &str) -> Result { + let conn = lock(&sync.conn); + let expected: Option<(String, i64, Option)> = conn + .query_row( + "SELECT secret_hash, expires_at_ms, used_at_ms + FROM sync_invites + WHERE invite_id = ?1", + [invite_id], + |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)), + ) + .optional()?; + let Some((secret_hash, expires_at_ms, used_at_ms)) = expected else { + return Ok(false); + }; + Ok(used_at_ms.is_none() && expires_at_ms >= now_ms() && secret_hash == hash_secret(secret)) +} + +fn is_active_trusted(sync: &DeviceSync, device_id: &str) -> Result { + let own = sync.ensure_identity()?.device_id; + if device_id == own { + return Ok(true); + } + let conn = lock(&sync.conn); + let found: Option = conn + .query_row( + "SELECT 1 FROM sync_devices + WHERE device_id = ?1 + AND trusted_at_ms IS NOT NULL + AND revoked_at_ms IS NULL", + [device_id], + |row| row.get(0), + ) + .optional()?; + Ok(found.is_some()) +} + +fn init_schema(conn: &Connection) -> Result<()> { + conn.execute_batch( + r#" +CREATE TABLE IF NOT EXISTS sync_meta ( + key TEXT PRIMARY KEY, + value TEXT NOT NULL +); +CREATE TABLE IF NOT EXISTS sync_devices ( + device_id TEXT PRIMARY KEY, + name TEXT NOT NULL DEFAULT '', + client_version TEXT NOT NULL DEFAULT '', + protocol_version INTEGER NOT NULL DEFAULT 1, + endpoint_id TEXT NOT NULL DEFAULT '', + endpoint_ticket TEXT NOT NULL DEFAULT '', + trusted_at_ms INTEGER, + last_seen_ms INTEGER, + revoked_at_ms INTEGER, + revoked_by TEXT, + revoke_cutoff_seq INTEGER +); +CREATE TABLE IF NOT EXISTS sync_invites ( + invite_id TEXT PRIMARY KEY, + secret_hash TEXT NOT NULL, + expires_at_ms INTEGER NOT NULL, + created_at_ms INTEGER NOT NULL, + used_at_ms INTEGER +); +CREATE TABLE IF NOT EXISTS sync_pending_pairing ( + request_id TEXT PRIMARY KEY, + device_id TEXT NOT NULL, + name TEXT NOT NULL, + client_version TEXT NOT NULL, + endpoint_id TEXT NOT NULL, + endpoint_ticket TEXT NOT NULL, + invite_id TEXT NOT NULL, + created_at_ms INTEGER NOT NULL, + answered_at_ms INTEGER, + status TEXT NOT NULL +); +CREATE TABLE IF NOT EXISTS sync_ops ( + op_id TEXT PRIMARY KEY, + origin_device_id TEXT NOT NULL, + seq INTEGER NOT NULL, + kind TEXT NOT NULL, + payload_json TEXT NOT NULL, + hlc_ms INTEGER NOT NULL, + received_at_ms INTEGER NOT NULL, + tombstone INTEGER NOT NULL DEFAULT 0 +); +CREATE INDEX IF NOT EXISTS idx_sync_ops_origin_seq ON sync_ops(origin_device_id, seq); +CREATE INDEX IF NOT EXISTS idx_sync_ops_tombstone ON sync_ops(tombstone); +CREATE TABLE IF NOT EXISTS sync_vectors ( + device_id TEXT PRIMARY KEY, + max_seq INTEGER NOT NULL +); +CREATE TABLE IF NOT EXISTS sync_peer_acks ( + peer_device_id TEXT NOT NULL, + origin_device_id TEXT NOT NULL, + max_seq INTEGER NOT NULL, + updated_at_ms INTEGER NOT NULL, + PRIMARY KEY (peer_device_id, origin_device_id) +); +CREATE TABLE IF NOT EXISTS sync_compacted ( + origin_device_id TEXT PRIMARY KEY, + max_seq INTEGER NOT NULL +); +CREATE TABLE IF NOT EXISTS sync_state_likes ( + content_id TEXT PRIMARY KEY, + liked INTEGER NOT NULL, + hlc_ms INTEGER NOT NULL, + op_id TEXT NOT NULL +); +CREATE TABLE IF NOT EXISTS sync_state_playlists ( + playlist_id TEXT PRIMARY KEY, + title TEXT NOT NULL, + deleted INTEGER NOT NULL DEFAULT 0, + hlc_ms INTEGER NOT NULL, + op_id TEXT NOT NULL +); +CREATE TABLE IF NOT EXISTS sync_state_playlist_items ( + playlist_id TEXT NOT NULL, + content_id TEXT NOT NULL, + present INTEGER NOT NULL DEFAULT 1, + position INTEGER NOT NULL DEFAULT 0, + hlc_ms INTEGER NOT NULL, + op_id TEXT NOT NULL, + PRIMARY KEY (playlist_id, content_id) +); +"#, + )?; + Ok(()) +} + +fn get_meta(conn: &Connection, key: &str) -> Result> { + Ok(conn + .query_row("SELECT value FROM sync_meta WHERE key = ?1", [key], |row| { + row.get(0) + }) + .optional()?) +} + +fn set_meta(conn: &Connection, key: &str, value: &str) -> Result<()> { + conn.execute( + "INSERT INTO sync_meta (key, value) VALUES (?1, ?2) + ON CONFLICT(key) DO UPDATE SET value = excluded.value", + params![key, value], + )?; + Ok(()) +} + +fn delete_meta(conn: &Connection, key: &str) -> Result<()> { + conn.execute("DELETE FROM sync_meta WHERE key = ?1", [key])?; + Ok(()) +} + +fn payload_kind(payload: &SyncOpPayload) -> &'static str { + match payload { + SyncOpPayload::TrackLikeSet { .. } => "track_like_set", + SyncOpPayload::PlaylistCreated { .. } => "playlist_created", + SyncOpPayload::PlaylistRenamed { .. } => "playlist_renamed", + SyncOpPayload::PlaylistDeleted { .. } => "playlist_deleted", + SyncOpPayload::PlaylistTrackAdded { .. } => "playlist_track_added", + SyncOpPayload::PlaylistTrackRemoved { .. } => "playlist_track_removed", + SyncOpPayload::DeviceProfileSet { .. } => "device_profile_set", + SyncOpPayload::DeviceRevoked { .. } => "device_revoked", + } +} + +fn compactable_tombstone_ids(conn: &Connection) -> Result> { + let active_devices = conn + .prepare( + "SELECT device_id FROM sync_devices + WHERE trusted_at_ms IS NOT NULL AND revoked_at_ms IS NULL", + )? + .query_map([], |row| row.get::<_, String>(0))? + .collect::>>()?; + if active_devices.len() <= 1 { + let mut stmt = conn.prepare( + "SELECT op_id, origin_device_id, seq + FROM sync_ops + WHERE tombstone = 1", + )?; + return Ok(stmt + .query_map([], |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)))? + .collect::>>()?); + } + let mut stmt = conn.prepare( + "SELECT op_id, origin_device_id, seq + FROM sync_ops + WHERE tombstone = 1 + ORDER BY received_at_ms", + )?; + let rows = stmt + .query_map([], |row| { + Ok(( + row.get::<_, String>(0)?, + row.get::<_, String>(1)?, + row.get::<_, i64>(2)?, + )) + })? + .collect::>>()?; + let mut out = Vec::new(); + for (op_id, origin, seq) in rows { + let mut all_acked = true; + for device in &active_devices { + let ack: Option = conn + .query_row( + "SELECT max_seq FROM sync_peer_acks + WHERE peer_device_id = ?1 AND origin_device_id = ?2", + params![device, origin], + |row| row.get(0), + ) + .optional()?; + let local_vector: Option = conn + .query_row( + "SELECT max_seq FROM sync_vectors + WHERE device_id = ?1", + [&origin], + |row| row.get(0), + ) + .optional()?; + let seen = ack.or(local_vector).unwrap_or(0); + if seen < seq { + all_acked = false; + break; + } + } + if all_acked { + out.push((op_id, origin, seq)); + } + } + Ok(out) +} + +fn peer_ack_floor_label(conn: &Connection) -> Result { + let count: i64 = conn.query_row("SELECT COUNT(*) FROM sync_peer_acks", [], |row| row.get(0))?; + if count == 0 { + return Ok("none".to_string()); + } + let min: i64 = conn.query_row( + "SELECT COALESCE(MIN(max_seq), 0) FROM sync_peer_acks", + [], + |row| row.get(0), + )?; + let max: i64 = conn.query_row( + "SELECT COALESCE(MAX(max_seq), 0) FROM sync_peer_acks", + [], + |row| row.get(0), + )?; + Ok(format!("{min}..{max} ({count} acks)")) +} + +async fn write_msg(stream: &mut ByteStream, message: &WireMessage) -> Result<()> { + let mut payload = serde_json::to_vec(message)?; + payload.push(b'\n'); + stream.send.write_all(&payload).await?; + Ok(()) +} + +async fn read_msg(stream: &mut ByteStream) -> Result { + let line = read_line(&mut stream.recv).await?; + Ok(serde_json::from_slice(&line)?) +} + +async fn read_line(reader: &mut R) -> Result> { + let mut out = Vec::new(); + loop { + let mut byte = [0u8; 1]; + let read = reader.read(&mut byte).await?; + if read == 0 { + break; + } + if byte[0] == b'\n' { + break; + } + out.push(byte[0]); + if out.len() > MAX_LINE { + anyhow::bail!("protocol line is too large"); + } + } + Ok(out) +} + +fn parse_invite(value: &str) -> Result { + let Some(token) = value.trim().strip_prefix("frid://i/") else { + anyhow::bail!("usage: :connect frid://i/"); + }; + let bytes = base64url_decode(token)?; + let invite: InviteWire = serde_json::from_slice(&bytes)?; + anyhow::ensure!(invite.v == 1, "unsupported invite version"); + Ok(invite) +} + +pub fn invite_network_id(value: &str) -> Result { + let invite = parse_invite(value)?; + let ticket: PeerTicket = invite + .ticket + .parse() + .map_err(|err| anyhow::anyhow!("malformed invite ticket: {err}"))?; + Ok(ticket.network_id) +} + +fn hash_secret(secret: &str) -> String { + blake3::hash(secret.as_bytes()).to_hex().to_string() +} + +fn random_hex(bytes: usize) -> String { + let key = SecretKey::generate(); + let mut seed = key.to_bytes().to_vec(); + while seed.len() < bytes { + seed.extend_from_slice(blake3::hash(&seed).as_bytes()); + } + hex_encode(&seed[..bytes]) +} + +fn hex_encode(bytes: &[u8]) -> String { + bytes.iter().map(|byte| format!("{byte:02x}")).collect() +} + +fn ticket_endpoint_id(ticket: &str) -> Option { + let ticket = PeerTicket::from_str(ticket).ok()?; + Some(ticket.endpoint_id().to_string()) +} + +fn default_device_name(device_id: &str) -> String { + let host = std::env::var("HOSTNAME") + .or_else(|_| std::env::var("COMPUTERNAME")) + .ok() + .filter(|value| !value.trim().is_empty()) + .unwrap_or_else(|| "furumi".to_string()); + format!("{host}-{}", short_id(device_id)) +} + +fn short_id(value: &str) -> String { + value.chars().take(10).collect() +} + +fn base64url_encode(bytes: &[u8]) -> String { + const TABLE: &[u8; 64] = b"ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789-_"; + let mut out = String::new(); + let mut i = 0; + while i < bytes.len() { + let b0 = bytes[i]; + let b1 = bytes.get(i + 1).copied().unwrap_or(0); + let b2 = bytes.get(i + 2).copied().unwrap_or(0); + out.push(TABLE[(b0 >> 2) as usize] as char); + out.push(TABLE[(((b0 & 0b0000_0011) << 4) | (b1 >> 4)) as usize] as char); + if i + 1 < bytes.len() { + out.push(TABLE[(((b1 & 0b0000_1111) << 2) | (b2 >> 6)) as usize] as char); + } + if i + 2 < bytes.len() { + out.push(TABLE[(b2 & 0b0011_1111) as usize] as char); + } + i += 3; + } + out +} + +fn base64url_decode(value: &str) -> Result> { + fn val(byte: u8) -> Option { + match byte { + b'A'..=b'Z' => Some(byte - b'A'), + b'a'..=b'z' => Some(byte - b'a' + 26), + b'0'..=b'9' => Some(byte - b'0' + 52), + b'-' => Some(62), + b'_' => Some(63), + _ => None, + } + } + let bytes = value.as_bytes(); + let mut out = Vec::new(); + let mut i = 0; + while i < bytes.len() { + let a = val(bytes[i]).context("invalid base64url invite")?; + let b = val(*bytes.get(i + 1).context("truncated base64url invite")?) + .context("invalid base64url invite")?; + let c = bytes.get(i + 2).and_then(|byte| val(*byte)); + let d = bytes.get(i + 3).and_then(|byte| val(*byte)); + out.push((a << 2) | (b >> 4)); + if let Some(c) = c { + out.push(((b & 0b0000_1111) << 4) | (c >> 2)); + if let Some(d) = d { + out.push(((c & 0b0000_0011) << 6) | d); + } + } + i += 4; + } + Ok(out) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn base64url_round_trip_without_padding() { + for input in [b"".as_slice(), b"a", b"ab", b"abc", b"abcdef"] { + let encoded = base64url_encode(input); + assert!(!encoded.contains('=')); + assert_eq!(base64url_decode(&encoded).unwrap(), input); + } + } + + #[test] + fn tombstone_detection() { + assert!( + SyncOpPayload::TrackLikeSet { + content_id: "b3:0".into(), + liked: false + } + .is_tombstone() + ); + assert!( + !SyncOpPayload::TrackLikeSet { + content_id: "b3:0".into(), + liked: true + } + .is_tombstone() + ); + } +} diff --git a/src/federation/mod.rs b/src/federation/mod.rs index 6f623fd..5ff01d3 100644 --- a/src/federation/mod.rs +++ b/src/federation/mod.rs @@ -181,11 +181,13 @@ pub struct FedPlayable { struct Running { service: Arc, network_name: String, + network_id: NetworkId, tasks: Vec>, } pub struct Federation { library: Arc, + devices: Arc, data_dir: PathBuf, cache_dir: PathBuf, media_dir: PathBuf, @@ -284,7 +286,7 @@ async fn dht_record_payload_bytes(data_dir: PathBuf, now_ms: u64) -> Result } impl Federation { - pub fn new(library: Arc) -> Arc { + pub fn new(library: Arc, devices: Arc) -> Arc { let dirs = crate::config::project_dirs(); let data_dir = dirs .as_ref() @@ -307,6 +309,7 @@ impl Federation { }); Arc::new(Self { library, + devices, data_dir, cache_dir, media_dir, @@ -363,9 +366,18 @@ impl Federation { /// Starts the DHT node. Idempotent per network name. async fn start(self: &Arc, network_name: String) -> Result<()> { + self.start_with_network_id(NetworkId::from_name(&network_name), network_name) + .await + } + + async fn start_with_network_id( + self: &Arc, + network_id: NetworkId, + network_name: String, + ) -> Result<()> { let mut guard = self.running.lock().await; if let Some(running) = guard.as_ref() { - if running.network_name == network_name { + if running.network_id == network_id { return Ok(()); } stop_running(guard.take()).await; @@ -375,13 +387,15 @@ impl Federation { let config = MusicDhtConfig::builder() .data_dir(&self.data_dir) - .network_id(NetworkId::from_name(&network_name)) + .network_id(network_id) // Peers of the network find each other knowing only its name. .rendezvous(RendezvousConfig::default()) // Peers stream each other's audio over this protocol. .stream_protocol(AUDIO_ALPN) // ...and browse each other's per-artist catalogs over this one. .stream_protocol(CATALOG_ALPN) + // Personal-device sync (likes, playlists, trusted devices). + .stream_protocol(crate::devices::SYNC_ALPN) .build() .map_err(|err| anyhow::anyhow!("invalid federation config: {err}"))?; let (service, mut events) = MusicDhtService::start(config) @@ -428,11 +442,32 @@ impl Federation { Arc::clone(&self.library), service.endpoint_id(), )); + let sync_acceptor = service + .stream_acceptor(crate::devices::SYNC_ALPN) + .map_err(|err| anyhow::anyhow!("failed to take the device-sync acceptor: {err}"))?; + let device_sync_task = tokio::spawn(crate::devices::serve_peers( + sync_acceptor, + Arc::clone(&self.devices), + Arc::clone(&service), + )); + let device_sync = Arc::clone(&self.devices); + let device_service = Arc::clone(&service); + let device_tick_task = tokio::spawn(async move { + crate::devices::sync_loop(device_sync, device_service).await; + }); *guard = Some(Running { service, network_name, - tasks: vec![event_task, sync_task, audio_task, catalog_task], + network_id, + tasks: vec![ + event_task, + sync_task, + audio_task, + catalog_task, + device_sync_task, + device_tick_task, + ], }); self.set_error(None); Ok(()) @@ -870,6 +905,38 @@ impl Federation { Ok(ticket.to_string()) } + pub async fn device_invite(self: &Arc) -> Result { + if self.running.lock().await.is_none() { + let status = self.devices.status(); + let network_name = format!("furumi-device-sync:{}", status.group_id); + self.start_with_network_id(NetworkId::from_name(&network_name), "device-sync".into()) + .await?; + } + let service = self.service().await?; + self.devices.create_invite(service).await + } + + pub async fn device_connect(self: &Arc, invite: &str) -> Result { + let network_id = crate::devices::invite_network_id(invite)?; + let needs_start = self + .running + .lock() + .await + .as_ref() + .is_none_or(|running| running.network_id != network_id); + if needs_start { + self.start_with_network_id(network_id, "device-invite".to_string()) + .await?; + } + let service = self.service().await?; + self.devices.connect_invite(service, invite).await + } + + pub async fn device_sync_now(self: &Arc) -> Result<()> { + let service = self.service().await?; + self.devices.sync_once(service).await + } + pub async fn connect(&self, ticket: &str) -> Result { let service = self.service().await?; let ticket: PeerTicket = ticket diff --git a/src/library/mod.rs b/src/library/mod.rs index a95e19c..30972a8 100644 --- a/src/library/mod.rs +++ b/src/library/mod.rs @@ -68,6 +68,7 @@ CREATE TABLE IF NOT EXISTS track_artists ( ); CREATE TABLE IF NOT EXISTS playlists ( id INTEGER PRIMARY KEY, + sync_id TEXT UNIQUE, title TEXT NOT NULL, description TEXT, created_at TEXT NOT NULL DEFAULT (datetime('now')) @@ -669,8 +670,14 @@ impl Library { pub fn create_playlist(&self, title: &str) -> Result { let conn = self.lock(); conn.execute("INSERT INTO playlists (title) VALUES (?1)", [title])?; + let id = conn.last_insert_rowid(); + let sync_id = make_playlist_sync_id(id, title); + conn.execute( + "UPDATE playlists SET sync_id = ?2 WHERE id = ?1", + params![id, sync_id], + )?; Ok(PlaylistCard { - id: conn.last_insert_rowid(), + id, title: title.to_string(), track_count: 0, kind: "normal".to_string(), @@ -692,6 +699,67 @@ impl Library { Ok(()) } + pub fn delete_playlist_by_sync_id(&self, sync_id: &str) -> Result<()> { + let conn = self.lock(); + conn.execute("DELETE FROM playlists WHERE sync_id = ?1", [sync_id])?; + Ok(()) + } + + pub fn playlist_sync_id(&self, id: i64) -> Result> { + let conn = self.lock(); + Ok(conn + .query_row("SELECT sync_id FROM playlists WHERE id = ?1", [id], |row| { + row.get(0) + }) + .optional()?) + } + + pub fn ensure_playlist_sync_id(&self, id: i64) -> Result { + let conn = self.lock(); + let existing: Option = conn + .query_row("SELECT sync_id FROM playlists WHERE id = ?1", [id], |row| { + row.get(0) + }) + .optional()? + .flatten(); + if let Some(sync_id) = existing { + return Ok(sync_id); + } + let title: String = + conn.query_row("SELECT title FROM playlists WHERE id = ?1", [id], |row| { + row.get(0) + })?; + let sync_id = make_playlist_sync_id(id, &title); + conn.execute( + "UPDATE playlists SET sync_id = ?2 WHERE id = ?1", + params![id, sync_id], + )?; + Ok(sync_id) + } + + pub fn upsert_synced_playlist(&self, sync_id: &str, title: &str) -> Result { + let conn = self.lock(); + if let Some(id) = conn + .query_row( + "SELECT id FROM playlists WHERE sync_id = ?1", + [sync_id], + |row| row.get(0), + ) + .optional()? + { + conn.execute( + "UPDATE playlists SET title = ?2 WHERE id = ?1", + params![id, title], + )?; + return Ok(id); + } + conn.execute( + "INSERT INTO playlists (sync_id, title) VALUES (?1, ?2)", + params![sync_id, title], + )?; + Ok(conn.last_insert_rowid()) + } + pub fn add_tracks_to_playlist(&self, playlist_id: i64, track_ids: &[i64]) -> Result<()> { let mut conn = self.lock(); let tx = conn.transaction()?; @@ -725,6 +793,125 @@ impl Library { Ok(()) } + pub fn track_content_id_by_id(&self, track_id: i64) -> Result> { + let conn = self.lock(); + let content_id: Option = conn + .query_row( + "SELECT content_id FROM tracks WHERE id = ?1", + [track_id], + |row| row.get::<_, Option>(0), + ) + .optional()? + .flatten() + .and_then(|value| music_dht::normalize_content_id(&value)); + Ok(content_id) + } + + pub fn track_content_ids(&self, track_ids: &[i64]) -> Result> { + let conn = self.lock(); + let mut out = Vec::new(); + for &track_id in track_ids { + let content_id: Option = conn + .query_row( + "SELECT content_id FROM tracks WHERE id = ?1", + [track_id], + |row| row.get::<_, Option>(0), + ) + .optional()? + .flatten() + .and_then(|value| music_dht::normalize_content_id(&value)); + if let Some(content_id) = content_id { + out.push(content_id); + } + } + Ok(out) + } + + pub fn track_id_by_content_id(&self, content_id: &str) -> Result> { + let Some(content_id) = music_dht::normalize_content_id(content_id) else { + return Ok(None); + }; + let conn = self.lock(); + Ok(conn + .query_row( + "SELECT id FROM tracks WHERE content_id = ?1 LIMIT 1", + [content_id], + |row| row.get(0), + ) + .optional()?) + } + + pub fn set_like(&self, track_id: i64, liked: bool) -> Result<()> { + let conn = self.lock(); + if liked { + conn.execute( + "INSERT OR IGNORE INTO likes (track_id) VALUES (?1)", + [track_id], + )?; + } else { + conn.execute("DELETE FROM likes WHERE track_id = ?1", [track_id])?; + } + Ok(()) + } + + pub fn add_content_id_to_synced_playlist( + &self, + playlist_sync_id: &str, + content_id: &str, + ) -> Result { + let Some(track_id) = self.track_id_by_content_id(content_id)? else { + return Ok(false); + }; + let conn = self.lock(); + let playlist_id: Option = conn + .query_row( + "SELECT id FROM playlists WHERE sync_id = ?1", + [playlist_sync_id], + |row| row.get::<_, i64>(0), + ) + .optional()?; + let Some(playlist_id) = playlist_id else { + return Ok(false); + }; + let next: i64 = conn.query_row( + "SELECT COALESCE(MAX(position), -1) + 1 FROM playlist_tracks WHERE playlist_id = ?1", + [playlist_id], + |row| row.get(0), + )?; + conn.execute( + "INSERT OR IGNORE INTO playlist_tracks (playlist_id, track_id, position) + VALUES (?1, ?2, ?3)", + params![playlist_id, track_id, next], + )?; + Ok(true) + } + + pub fn remove_content_id_from_synced_playlist( + &self, + playlist_sync_id: &str, + content_id: &str, + ) -> Result<()> { + let Some(track_id) = self.track_id_by_content_id(content_id)? else { + return Ok(()); + }; + let conn = self.lock(); + let playlist_id: Option = conn + .query_row( + "SELECT id FROM playlists WHERE sync_id = ?1", + [playlist_sync_id], + |row| row.get::<_, i64>(0), + ) + .optional()?; + let Some(playlist_id) = playlist_id else { + return Ok(()); + }; + conn.execute( + "DELETE FROM playlist_tracks WHERE playlist_id = ?1 AND track_id = ?2", + params![playlist_id, track_id], + )?; + Ok(()) + } + /// Liked federated tracks, newest first, as playable references. pub fn fed_likes(&self) -> Result> { let conn = self.lock(); @@ -1222,6 +1409,28 @@ fn ensure_schema_migrations(conn: &Connection) -> Result<()> { if !track_columns.iter().any(|column| column == "content_id") { conn.execute("ALTER TABLE tracks ADD COLUMN content_id TEXT", [])?; } + let playlist_columns = table_columns(conn, "playlists")?; + if !playlist_columns.iter().any(|column| column == "sync_id") { + conn.execute("ALTER TABLE playlists ADD COLUMN sync_id TEXT", [])?; + conn.execute( + "CREATE UNIQUE INDEX IF NOT EXISTS idx_playlists_sync_id + ON playlists(sync_id)", + [], + )?; + } + let mut rows = conn.prepare("SELECT id, title FROM playlists WHERE sync_id IS NULL")?; + let missing = rows + .query_map([], |row| { + Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?)) + })? + .collect::>>()?; + drop(rows); + for (id, title) in missing { + conn.execute( + "UPDATE playlists SET sync_id = ?2 WHERE id = ?1", + params![id, make_playlist_sync_id(id, &title)], + )?; + } Ok(()) } @@ -1232,6 +1441,15 @@ fn table_columns(conn: &Connection, table: &str) -> Result> { .collect::>>()?) } +fn make_playlist_sync_id(id: i64, title: &str) -> String { + let now = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|d| d.as_nanos()) + .unwrap_or(0); + let seed = format!("playlist:{id}:{title}:{now}:{}", std::process::id()); + format!("pl_{}", &blake3::hash(seed.as_bytes()).to_hex()[..24]) +} + #[cfg(test)] mod tests { use super::*; diff --git a/src/main.rs b/src/main.rs index ca4f793..8313963 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,6 +1,7 @@ mod app; mod art; mod config; +mod devices; mod federation; mod library; mod media; diff --git a/src/ui/federation.rs b/src/ui/federation.rs index 10836bc..ef0e12a 100644 --- a/src/ui/federation.rs +++ b/src/ui/federation.rs @@ -6,7 +6,7 @@ use ratatui::text::{Line, Span}; use ratatui::widgets::{Block, Paragraph}; use super::theme; -use crate::app::state::{AppState, FedRow}; +use crate::app::state::{AppState, FedRow, settings_rows}; pub fn draw(frame: &mut Frame, area: Rect, state: &AppState) { let block = Block::bordered() @@ -16,7 +16,7 @@ pub fn draw(frame: &mut Frame, area: Rect, state: &AppState) { let inner = block.inner(area); frame.render_widget(block, area); - let rows_height = (FedRow::ALL.len() + state.visualizer.scripts.len() + 6) as u16; + let rows_height = (settings_rows(state).len() + 5) as u16; let [rows_area, _, status_area] = Layout::vertical([ Constraint::Length(rows_height.min(inner.height)), Constraint::Length(1), @@ -66,6 +66,83 @@ fn draw_settings_rows(frame: &mut Frame, area: Rect, state: &AppState) { cursor += 1; } + y = y.saturating_add(1); + draw_section(frame, area, &mut y, "Connected Devices"); + let devices = state.federation.devices.as_ref(); + draw_row( + frame, + area, + &mut y, + cursor, + state.settings_cursor, + "This device name", + devices + .map(|status| status.this_device_name.clone()) + .unwrap_or_else(|| "loading…".to_string()), + ); + cursor += 1; + draw_row( + frame, + area, + &mut y, + cursor, + state.settings_cursor, + "Generate device invite", + "↵".to_string(), + ); + cursor += 1; + draw_row( + frame, + area, + &mut y, + cursor, + state.settings_cursor, + "Connect device by invite…", + "↵".to_string(), + ); + cursor += 1; + draw_row( + frame, + area, + &mut y, + cursor, + state.settings_cursor, + "Sync devices now", + "↵".to_string(), + ); + cursor += 1; + if let Some(status) = devices { + for device in &status.devices { + let label = if device.is_self { + format!("* {}", device.name) + } else if device.revoked { + format!(" {} (revoked)", device.name) + } else { + format!(" {}", device.name) + }; + let version = if device.client_version.is_empty() { + "unknown".to_string() + } else { + format!("v{}", device.client_version) + }; + let value = if device.is_self || device.revoked { + version + } else { + format!("{version} · revoke ↵") + }; + draw_row( + frame, + area, + &mut y, + cursor, + state.settings_cursor, + &label, + value, + ); + cursor += 1; + } + } + y = y.saturating_add(1); draw_section(frame, area, &mut y, "Visualizations"); draw_row( @@ -273,5 +350,53 @@ fn draw_status(frame: &mut Frame, area: Rect, state: &AppState) { } } } + lines.push(Line::default()); + lines.push(Line::styled("Connected Devices", theme::header())); + match &state.federation.devices { + None => lines.push(Line::styled("loading…", theme::dim())), + Some(status) => { + lines.push(status_line("This device", status.this_device_id.clone())); + lines.push(status_line("Sync group", status.group_id.clone())); + lines.push(status_line( + "Active devices", + status.active_devices.to_string(), + )); + lines.push(status_line( + "Revoked devices", + status.revoked_devices.to_string(), + )); + lines.push(status_line( + "Pending requests", + status.pending_requests.to_string(), + )); + lines.push(status_line("Ops in log", status.ops_total.to_string())); + lines.push(status_line( + "Tombstones", + format!( + "{} ({} compactable)", + status.tombstone_ops, status.compactable_tombstones + ), + )); + lines.push(status_line("Outbox ops", status.outbox_ops.to_string())); + lines.push(status_line( + "Snapshot", + format!( + "{} likes, {} playlists, {} items", + status.snapshot_likes, status.snapshot_playlists, status.snapshot_items + ), + )); + lines.push(status_line( + "Unresolved items", + status.unresolved_playlist_items.to_string(), + )); + lines.push(status_line("Peer ack floor", status.peer_ack_floor.clone())); + if let Some(last_sync) = &status.last_sync { + lines.push(status_line("Last device sync", last_sync.clone())); + } + if let Some(last_error) = &status.last_error { + lines.push(status_line("Device error", last_error.clone())); + } + } + } frame.render_widget(Paragraph::new(lines), area); } diff --git a/src/ui/popup.rs b/src/ui/popup.rs index 8831180..2cd1836 100644 --- a/src/ui/popup.rs +++ b/src/ui/popup.rs @@ -36,6 +36,15 @@ pub fn draw(frame: &mut Frame, state: &AppState) { Some(Popup::LogDetail(entry)) => draw_log_detail(frame, entry), Some(Popup::FedInput { field, input }) => draw_fed_input(frame, field.title(), input), Some(Popup::FedText { title, text }) => draw_fed_text(frame, title, text), + Some(Popup::DevicePairing { + device_id, + name, + client_version, + .. + }) => draw_device_pairing(frame, device_id, name, client_version), + Some(Popup::ConfirmDeviceRevoke { device_id, name }) => { + draw_device_revoke(frame, device_id, name) + } None => {} } } @@ -127,6 +136,58 @@ fn draw_fed_text(frame: &mut Frame, title: &str, text: &str) { ); } +fn draw_device_pairing(frame: &mut Frame, device_id: &str, name: &str, client_version: &str) { + let area = centered(frame.area(), 64, 8); + let block = Block::bordered() + .title(" Pair device ") + .title_style(theme::header()) + .border_style(theme::accent()); + let inner = block.inner(area); + frame.render_widget(Clear, area); + frame.render_widget(block, area); + let lines = vec![ + Line::from(vec![ + Span::styled("Name ", theme::dim()), + Span::raw(name.to_string()), + ]), + Line::from(vec![ + Span::styled("Version ", theme::dim()), + Span::raw(client_version.to_string()), + ]), + Line::from(vec![ + Span::styled("Device ", theme::dim()), + Span::raw(device_id.chars().take(24).collect::()), + ]), + Line::default(), + Line::styled("enter/y accept · n/esc deny", theme::dim()), + ]; + frame.render_widget(Paragraph::new(lines), inner); +} + +fn draw_device_revoke(frame: &mut Frame, device_id: &str, name: &str) { + let area = centered(frame.area(), 64, 7); + let block = Block::bordered() + .title(" Revoke device ") + .title_style(theme::header()) + .border_style(theme::accent()); + let inner = block.inner(area); + frame.render_widget(Clear, area); + frame.render_widget(block, area); + let lines = vec![ + Line::from(vec![ + Span::styled("Device ", theme::dim()), + Span::raw(name.to_string()), + ]), + Line::from(vec![ + Span::styled("ID ", theme::dim()), + Span::raw(device_id.chars().take(24).collect::()), + ]), + Line::default(), + Line::styled("enter/y revoke · n/esc cancel", theme::dim()), + ]; + frame.render_widget(Paragraph::new(lines), inner); +} + /// Metadata edit form: one bordered input per field, the focused field gets /// the accent border and a cursor block. fn draw_edit(