//! 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?; finish_send(&mut stream).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?; finish_send(&mut stream).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?; finish_send(&mut stream).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?; finish_send(&mut stream).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?; finish_send(&mut stream).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?; finish_send(&mut stream).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?; finish_send(&mut stream).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?; finish_send(&mut stream).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 finish_send(stream: &mut ByteStream) -> Result<()> { stream.send.finish()?; 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() ); } }