diff --git a/Cargo.lock b/Cargo.lock index 0ebd346..6284232 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2159,6 +2159,7 @@ dependencies = [ "rand 0.9.5", "rusqlite", "serde", + "serde_json", "tempfile", "thiserror 2.0.19", "tokio", @@ -3299,6 +3300,19 @@ dependencies = [ "syn 3.0.2", ] +[[package]] +name = "serde_json" +version = "1.0.151" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c841b55ecdae098c80dcae9cf767f6f8a0c2cdb3416bbef72181df4d0fe73f14" +dependencies = [ + "itoa", + "memchr", + "serde", + "serde_core", + "zmij", +] + [[package]] name = "serdect" version = "0.4.3" @@ -4681,3 +4695,9 @@ dependencies = [ "quote", "syn 2.0.119", ] + +[[package]] +name = "zmij" +version = "1.0.23" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "29666d0abbfad1e3dc4dcf6144730dd3a3ab225bbbdac83319345b1b44ccfc1b" diff --git a/Cargo.toml b/Cargo.toml index 4b27dca..f249e9e 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -21,6 +21,7 @@ iroh-tickets = "1" mainline = "7" tokio = { version = "1", features = ["rt-multi-thread", "macros", "sync", "time", "io-util", "io-std", "signal"] } serde = { version = "1", features = ["derive"] } +serde_json = "1" postcard = { version = "1", features = ["alloc"] } blake3 = "1" futures = "0.3" diff --git a/crates/music-dht/Cargo.toml b/crates/music-dht/Cargo.toml index c9d8111..166a4d3 100644 --- a/crates/music-dht/Cargo.toml +++ b/crates/music-dht/Cargo.toml @@ -10,6 +10,7 @@ rust-version.workspace = true federation-net = { workspace = true } tokio = { workspace = true } serde = { workspace = true } +serde_json = { workspace = true } postcard = { workspace = true } blake3 = { workspace = true } futures = { workspace = true } diff --git a/crates/music-dht/src/device_sync.rs b/crates/music-dht/src/device_sync.rs new file mode 100644 index 0000000..d60beeb --- /dev/null +++ b/crates/music-dht/src/device_sync.rs @@ -0,0 +1,733 @@ +//! Wire protocol shared by furumi personal-device synchronization clients. +//! +//! The storage and UI policy live in applications. This module only defines +//! the stable JSON-lines messages spoken over the auxiliary iroh stream ALPN, +//! plus small helpers for opaque `frid://i/...` invites. + +use std::collections::BTreeMap; +use std::str::FromStr; + +use federation_net::{ByteStream, NetworkId, PeerTicket, SecretKey}; +use serde::{Deserialize, Serialize}; +use tokio::io::{AsyncRead, AsyncReadExt}; + +use crate::error::{MusicDhtError, Result}; + +/// Auxiliary ALPN used for personal-device sync streams. +pub const SYNC_ALPN: &[u8] = b"furumi/sync/1"; + +/// Version of the personal-device sync wire protocol. +pub const DEVICE_SYNC_PROTOCOL_VERSION: u16 = 1; + +/// Default invite lifetime used by applications unless they need a custom TTL. +pub const DEFAULT_INVITE_TTL_MS: i64 = 10 * 60 * 1000; + +/// Maximum JSON line accepted by the protocol. +pub const MAX_SYNC_LINE: usize = 8 * 1024 * 1024; + +/// Opaque invite payload encoded into `frid://i/`. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct InviteWire { + /// Invite schema version. + pub v: u16, + /// Transport ticket for the inviting endpoint. + #[serde(rename = "t")] + pub ticket: String, + /// Device id of the inviter. + #[serde(rename = "d")] + pub device_id: String, + /// Short invite id stored by the inviter. + #[serde(rename = "i")] + pub invite_id: String, + /// Runtime pairing secret. + #[serde(rename = "s")] + pub secret: String, + /// Expiration timestamp in unix milliseconds. + #[serde(rename = "e")] + pub expires_at_ms: i64, +} + +/// Public device profile exchanged inside a trusted sync group. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct DeviceProfileWire { + /// Stable device id. + pub device_id: String, + /// User-visible device name. + pub name: String, + /// Application version. + pub client_version: String, + /// Wire protocol version. + pub protocol_version: u16, + /// Transport endpoint id, as string. + pub endpoint_id: String, + /// Transport ticket used for direct reconnects. + pub endpoint_ticket: String, + /// Whether the device has been revoked. + #[serde(default)] + pub revoked: bool, + /// Last accepted seq from this device after revocation. + #[serde(default)] + pub revoke_cutoff_seq: Option, + /// Profile update timestamp in unix milliseconds. + #[serde(default)] + pub updated_at_ms: i64, +} + +/// Portable track reference embedded into sync payloads for unresolved tracks. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct SyncedFedTrack { + /// DHT item id, when known. + pub item_id: String, + /// DHT owner endpoint id, when known. + pub owner: String, + /// Track title. + pub title: String, + /// Main artist names. + #[serde(default)] + pub artist_names: Vec, + /// Featured artist names. + #[serde(default)] + pub featured_artist_names: Vec, + /// Release or track year. + pub year: Option, + /// Duration in seconds. + pub duration_seconds: Option, + /// Stable content id (`b3:`). + pub content_id: String, + /// Release title, when known. + pub release_title: Option, + /// Track number. + pub track_number: Option, + /// Disc number. + pub disc_number: Option, +} + +/// Portable playback queue track. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct PlaybackTrack { + /// Local or placeholder id on the sending device. + pub id: i64, + /// Track title. + pub title: String, + /// Track number. + pub track_number: Option, + /// Disc number. + pub disc_number: Option, + /// Duration in seconds. + pub duration_seconds: f64, + /// Main artist names. + #[serde(default)] + pub artist_names: Vec, + /// Featured artist names. + #[serde(default)] + pub featured_artist_names: Vec, + /// Local or placeholder release id. + pub release_id: i64, + /// Release title. + pub release_title: String, + /// Release year. + pub release_year: Option, + /// Legacy compatibility only. Paths are device-local and should be empty. + #[serde(default)] + pub file_path: String, + /// Stable content id. + pub content_id: Option, + /// Audio format label. + pub audio_format: Option, + /// Audio bitrate. + pub audio_bitrate: Option, + /// Audio sample rate. + pub audio_sample_rate: Option, + /// Audio bit depth. + pub audio_bit_depth: Option, + /// File size in bytes. + pub file_size_bytes: Option, + /// Sender-side play count. + #[serde(default)] + pub play_count: i64, + /// Federation metadata for unresolved tracks. + #[serde(default)] + pub fed: Option, +} + +/// Playback repeat mode. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)] +#[serde(rename_all = "snake_case")] +pub enum PlaybackRepeat { + /// Repeat disabled. + #[default] + Off, + /// Repeat current track. + One, + /// Repeat the whole queue. + All, +} + +/// Portable playback state. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct PlaybackStateWire { + /// Queue tracks. + #[serde(default)] + pub queue: Vec, + /// Current queue index. + #[serde(default)] + pub queue_pos: usize, + /// Whether playback is active. + pub playing: bool, + /// Whether playback is paused. + pub paused: bool, + /// Active device idle timestamp, when paused/stopped. + #[serde(default)] + pub idle_since_ms: Option, + /// Current position in seconds. + pub position_secs: f64, + /// Volume 0..100. + #[serde(default)] + pub volume: u8, + /// Shuffle mode. + pub shuffle: bool, + /// Repeat mode. + pub repeat: PlaybackRepeat, +} + +/// Playback status broadcast by a device during sync. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct PlaybackSnapshot { + /// Device id. + pub device_id: String, + /// Device name. + pub device_name: String, + /// Whether this device considers itself active. + pub active: bool, + /// Update timestamp in unix milliseconds. + pub updated_at_ms: i64, + /// Playback state. + pub state: PlaybackStateWire, +} + +/// Playback command delivered through the sync op log. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(tag = "kind", rename_all = "snake_case")] +pub enum PlaybackCommand { + /// Replace target state. `seek` means the target should also seek audio. + SetState { + /// New state. + state: PlaybackStateWire, + /// Whether to seek the audio sink to `state.position_secs`. + #[serde(default)] + seek: bool, + }, + /// Another device became active. + ActiveChanged { + /// New active device id. + active_device_id: String, + /// New active device name. + active_device_name: String, + /// Transferred playback state. + state: PlaybackStateWire, + }, +} + +/// One operation in the device-sync log. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct SyncOpWire { + /// Unique op id, usually `:`. + pub op_id: String, + /// Origin device id. + pub origin_device_id: String, + /// Monotonic sequence inside the origin device log. + pub seq: i64, + /// Hybrid/logical timestamp in unix milliseconds. + pub hlc_ms: i64, + /// Operation payload. + pub payload: SyncOpPayload, +} + +/// Personal-device sync operation payload. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(tag = "kind", rename_all = "snake_case")] +pub enum SyncOpPayload { + /// Set like state for a content id. + TrackLikeSet { + /// Content id. + content_id: String, + /// Desired like state. + liked: bool, + /// Optional unresolved federated metadata. + #[serde(default)] + fed: Option, + }, + /// Create playlist. + PlaylistCreated { + /// Stable playlist sync id. + playlist_id: String, + /// Playlist title. + title: String, + }, + /// Rename playlist. + PlaylistRenamed { + /// Stable playlist sync id. + playlist_id: String, + /// Playlist title. + title: String, + }, + /// Delete playlist. + PlaylistDeleted { + /// Stable playlist sync id. + playlist_id: String, + }, + /// Add track to playlist. + PlaylistTrackAdded { + /// Stable playlist sync id. + playlist_id: String, + /// Content id. + content_id: String, + /// Playlist position. + position: i64, + /// Optional unresolved federated metadata. + #[serde(default)] + fed: Option, + }, + /// Remove track from playlist. + PlaylistTrackRemoved { + /// Stable playlist sync id. + playlist_id: String, + /// Content id. + content_id: String, + }, + /// Update origin device profile. + DeviceProfileSet { + /// Display name. + name: String, + /// Client version. + client_version: String, + /// Endpoint ticket. + endpoint_ticket: String, + /// Endpoint id. + endpoint_id: String, + }, + /// Trust a device. + DeviceTrusted { + /// Trusted device id. + target_device_id: String, + }, + /// Revoke a device. + DeviceRevoked { + /// Revoked device id. + target_device_id: String, + /// Max seq seen from the target at revoke time. + target_max_seq_seen: i64, + }, + /// Playback command for a target device. + PlaybackCommand { + /// Target device id. + target_device_id: String, + /// Command. + command: PlaybackCommand, + }, +} + +impl SyncOpPayload { + /// Returns true for operations that remove state and may eventually compact. + pub fn is_tombstone(&self) -> bool { + matches!( + self, + SyncOpPayload::TrackLikeSet { liked: false, .. } + | SyncOpPayload::PlaylistDeleted { .. } + | SyncOpPayload::PlaylistTrackRemoved { .. } + | SyncOpPayload::DeviceRevoked { .. } + ) + } +} + +/// Materialized sync snapshot sent with every handshake. +#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)] +pub struct SyncSnapshot { + /// Current likes. + #[serde(default)] + pub likes: Vec, + /// Like tombstones. + #[serde(default)] + pub unlikes: Vec, + /// Current playlists. + #[serde(default)] + pub playlists: Vec, + /// Deleted playlists. + #[serde(default)] + pub deleted_playlists: Vec, + /// Removed playlist items. + #[serde(default)] + pub removed_playlist_items: Vec, +} + +/// Snapshot like row. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct SnapshotLike { + /// Content id. + pub content_id: String, + /// Last write timestamp. + pub hlc_ms: i64, + /// Last write op id. + pub op_id: String, + /// Optional unresolved federated metadata. + #[serde(default)] + pub fed: Option, +} + +/// Snapshot unlike tombstone. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct SnapshotLikeTombstone { + /// Content id. + pub content_id: String, + /// Last write timestamp. + pub hlc_ms: i64, + /// Last write op id. + pub op_id: String, +} + +/// Snapshot playlist. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct SnapshotPlaylist { + /// Stable playlist sync id. + pub playlist_id: String, + /// Title. + pub title: String, + /// Last write timestamp. + pub hlc_ms: i64, + /// Last write op id. + pub op_id: String, + /// Current items. + #[serde(default)] + pub items: Vec, +} + +/// Snapshot deleted playlist tombstone. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct SnapshotPlaylistTombstone { + /// Stable playlist sync id. + pub playlist_id: String, + /// Last write timestamp. + pub hlc_ms: i64, + /// Last write op id. + pub op_id: String, +} + +/// Snapshot playlist item. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct SnapshotPlaylistItem { + /// Content id. + pub content_id: String, + /// Playlist position. + pub position: i64, + /// Last write timestamp. + pub hlc_ms: i64, + /// Last write op id. + pub op_id: String, + /// Optional unresolved federated metadata. + #[serde(default)] + pub fed: Option, +} + +/// Snapshot removed playlist item tombstone. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct SnapshotPlaylistItemTombstone { + /// Stable playlist sync id. + pub playlist_id: String, + /// Content id. + pub content_id: String, + /// Last write timestamp. + pub hlc_ms: i64, + /// Last write op id. + pub op_id: String, +} + +/// Top-level message spoken on [`SYNC_ALPN`]. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(tag = "type", rename_all = "snake_case")] +pub enum WireMessage { + /// Pairing request sent by the invite consumer. + PairRequest { + /// Invite id. + invite_id: String, + /// Invite secret. + secret: String, + /// Requester device profile. + profile: DeviceProfileWire, + /// Requester sync group id. + #[serde(default)] + group_id: Option, + /// Number of active devices in requester group. + #[serde(default)] + group_active_devices: usize, + /// Requester trusted devices. + #[serde(default)] + devices: Vec, + /// Requester vector clock. + vector: BTreeMap, + /// Requester pending operations. + ops: Vec, + /// Requester materialized snapshot. + snapshot: SyncSnapshot, + /// Requester playback snapshot. + #[serde(default)] + playback: Option, + }, + /// Pairing response sent by the inviter. + PairResponse { + /// Whether pairing was accepted. + accepted: bool, + /// Whether the request is waiting for user confirmation. + #[serde(default)] + pending: bool, + /// Optional rejection error. + #[serde(default)] + error: Option, + /// Accepted group id. + #[serde(default)] + group_id: Option, + /// Inviter profile. + #[serde(default)] + profile: Option, + /// Inviter trusted devices. + #[serde(default)] + devices: Vec, + /// Inviter vector clock. + #[serde(default)] + vector: BTreeMap, + /// Inviter pending operations. + #[serde(default)] + ops: Vec, + /// Inviter materialized snapshot. + #[serde(default)] + snapshot: SyncSnapshot, + /// Inviter playback snapshot. + #[serde(default)] + playback: Option, + }, + /// Regular sync hello. + Hello { + /// Sync group id. + group_id: String, + /// Sender profile. + profile: DeviceProfileWire, + /// Sender trusted devices. + devices: Vec, + /// Sender vector clock. + vector: BTreeMap, + /// Sender pending operations. + ops: Vec, + /// Sender materialized snapshot. + snapshot: SyncSnapshot, + /// Sender playback snapshot. + #[serde(default)] + playback: Option, + }, + /// Regular sync response. + SyncResponse { + /// Whether sync was accepted. + accepted: bool, + /// Optional rejection error. + #[serde(default)] + error: Option, + /// Receiver trusted devices. + #[serde(default)] + devices: Vec, + /// Receiver vector clock. + #[serde(default)] + vector: BTreeMap, + /// Receiver pending operations. + #[serde(default)] + ops: Vec, + /// Receiver materialized snapshot. + #[serde(default)] + snapshot: SyncSnapshot, + /// Receiver playback snapshot. + #[serde(default)] + playback: Option, + }, +} + +/// Encodes an invite as a compact `frid://i/...` link. +pub fn encode_invite(invite: &InviteWire) -> Result { + let bytes = serde_json::to_vec(invite).map_err(protocol_err)?; + Ok(format!("frid://i/{}", base64url_encode(&bytes))) +} + +/// Parses a `frid://i/...` invite. +pub fn parse_invite(value: &str) -> Result { + let Some(token) = value.trim().strip_prefix("frid://i/") else { + return Err(MusicDhtError::Protocol( + "usage: frid://i/".to_string(), + )); + }; + let bytes = base64url_decode(token)?; + let invite: InviteWire = serde_json::from_slice(&bytes).map_err(protocol_err)?; + if invite.v != 1 { + return Err(MusicDhtError::Protocol( + "unsupported invite version".to_string(), + )); + } + Ok(invite) +} + +/// Extracts the network id carried by an invite's endpoint ticket. +pub fn invite_network_id(value: &str) -> Result { + let invite = parse_invite(value)?; + let ticket = PeerTicket::from_str(&invite.ticket).map_err(protocol_err)?; + Ok(ticket.network_id) +} + +/// Hashes an invite secret for durable storage. +pub fn hash_secret(secret: &str) -> String { + blake3::hash(secret.as_bytes()).to_hex().to_string() +} + +/// Generates a random lowercase hex token. +pub 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]) +} + +/// Extracts endpoint id from a ticket string. +pub fn ticket_endpoint_id(ticket: &str) -> Option { + let ticket = PeerTicket::from_str(ticket).ok()?; + Some(ticket.endpoint_id().to_string()) +} + +/// Writes a JSON-lines sync message. +pub async fn write_msg(stream: &mut ByteStream, message: &WireMessage) -> Result<()> { + let mut payload = serde_json::to_vec(message).map_err(protocol_err)?; + payload.push(b'\n'); + stream.send.write_all(&payload).await.map_err(network_err)?; + Ok(()) +} + +/// Finishes the sending side. +pub async fn finish_send(stream: &mut ByteStream) -> Result<()> { + stream.send.finish().map_err(network_err)?; + Ok(()) +} + +/// Finishes the response side and waits briefly for peer acknowledgement. +pub async fn finish_response( + stream: &mut ByteStream, + drain_timeout: std::time::Duration, +) -> Result<()> { + stream.send.finish().map_err(network_err)?; + let _ = tokio::time::timeout(drain_timeout, stream.send.stopped()).await; + Ok(()) +} + +/// Reads a JSON-lines sync message. +pub async fn read_msg(stream: &mut ByteStream) -> Result { + let line = read_line(&mut stream.recv).await?; + serde_json::from_slice(&line).map_err(protocol_err) +} + +/// Reads one newline-terminated protocol line. +pub async fn read_line(reader: &mut R) -> Result> { + read_line_limited(reader, MAX_SYNC_LINE).await +} + +/// Reads one newline-terminated protocol line with a custom size limit. +pub async fn read_line_limited( + reader: &mut R, + max_line: usize, +) -> Result> { + let mut out = Vec::new(); + loop { + let mut byte = [0u8; 1]; + let read = reader.read(&mut byte).await.map_err(network_err)?; + if read == 0 { + break; + } + if byte[0] == b'\n' { + break; + } + out.push(byte[0]); + if out.len() > max_line { + return Err(MusicDhtError::Protocol( + "protocol line is too large".to_string(), + )); + } + } + Ok(out) +} + +/// Reads one sync message from any async reader. +pub async fn read_msg_from(reader: &mut R) -> Result { + let line = read_line(reader).await?; + serde_json::from_slice(&line).map_err(protocol_err) +} + +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::with_capacity(bytes.len() * 3 / 4); + let mut i = 0; + while i < bytes.len() { + let a = val(bytes[i]) + .ok_or_else(|| MusicDhtError::Protocol("invalid base64url invite".to_string()))?; + let b = val(*bytes + .get(i + 1) + .ok_or_else(|| MusicDhtError::Protocol("truncated base64url invite".to_string()))?) + .ok_or_else(|| MusicDhtError::Protocol("invalid base64url invite".to_string()))?; + 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 << 4) | (c >> 2)); + if let Some(d) = d { + out.push((c << 6) | d); + } + } + i += 4; + } + Ok(out) +} + +fn hex_encode(bytes: &[u8]) -> String { + bytes.iter().map(|byte| format!("{byte:02x}")).collect() +} + +fn protocol_err(err: impl std::fmt::Display) -> MusicDhtError { + MusicDhtError::Protocol(err.to_string()) +} + +fn network_err(err: impl std::fmt::Display) -> MusicDhtError { + MusicDhtError::Network(err.to_string()) +} diff --git a/crates/music-dht/src/lib.rs b/crates/music-dht/src/lib.rs index c72beee..bcc61ac 100644 --- a/crates/music-dht/src/lib.rs +++ b/crates/music-dht/src/lib.rs @@ -79,6 +79,7 @@ mod config; mod database; +pub mod device_sync; mod dht; mod error; mod message;