diff --git a/Cargo.lock b/Cargo.lock index fe7ab71..a2e9cd5 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1845,12 +1845,14 @@ checksum = "e6d5a32815ae3f33302d95fdcb2ce17862f8c65363dcfd29360480ba1001fc9c" [[package]] name = "furumusic" -version = "0.8.3" +version = "0.8.4" dependencies = [ "anyhow", + "async-stream", "async-trait", "base64 0.22.1", "blake3", + "bytes", "chrono", "cot", "croner", diff --git a/Cargo.toml b/Cargo.toml index d102879..64e7cc4 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "furumusic" -version = "0.8.4" +version = "0.9.0" edition = "2024" description = "Reusable web-app boilerplate: auth, OIDC/SSO, admin panel, user management, i18n, PostgreSQL" @@ -14,6 +14,8 @@ serde = { version = "1", features = ["derive"] } openidconnect = "4.0" reqwest = { version = "0.12", default-features = false, features = ["rustls-tls", "json"] } tokio = { version = "1", features = ["sync", "fs", "io-util"] } +async-stream = "0.3" +bytes = "1" tower = "0.5" base64 = "0.22" blake3 = "1" diff --git a/src/admin/v2.rs b/src/admin/v2.rs index 003ad60..d4526c7 100644 --- a/src/admin/v2.rs +++ b/src/admin/v2.rs @@ -452,6 +452,8 @@ struct AdminSettingsValues { federation_enabled: bool, #[serde(default)] federation_network_id: String, + #[serde(default)] + federation_save_on_listen: bool, } #[derive(Debug, Clone, Serialize, JsonSchema)] @@ -478,6 +480,7 @@ struct AdminSettingsSources { agent_concurrency: &'static str, federation_enabled: &'static str, federation_network_id: &'static str, + federation_save_on_listen: &'static str, } #[derive(Debug, Deserialize)] @@ -506,6 +509,8 @@ pub(super) struct UpdateSettingsRequest { federation_enabled: bool, #[serde(default)] federation_network_id: String, + #[serde(default)] + federation_save_on_listen: bool, } #[derive(Debug, Serialize, JsonSchema)] @@ -993,6 +998,10 @@ pub async fn update_settings( "federation_network_id", body.federation_network_id.trim().to_string(), ), + ( + "federation_save_on_listen", + body.federation_save_on_listen.to_string(), + ), ]; for (key, value) in fields { let mut entry = ConfigEntry::new(key.to_string(), value); @@ -1143,6 +1152,7 @@ fn settings_dto(config: AppConfig, sources: ConfigSources) -> AdminSettingsDto { agent_concurrency: config.agent_concurrency.to_string(), federation_enabled: config.federation_enabled, federation_network_id: config.federation_network_id, + federation_save_on_listen: config.federation_save_on_listen, }, sources: AdminSettingsSources { auth_password_enabled: sources.auth_password_enabled.code(), @@ -1167,6 +1177,7 @@ fn settings_dto(config: AppConfig, sources: ConfigSources) -> AdminSettingsDto { agent_concurrency: sources.agent_concurrency.code(), federation_enabled: sources.federation_enabled.code(), federation_network_id: sources.federation_network_id.code(), + federation_save_on_listen: sources.federation_save_on_listen.code(), }, } } diff --git a/src/config.rs b/src/config.rs index 61724d9..3d000e7 100644 --- a/src/config.rs +++ b/src/config.rs @@ -137,6 +137,7 @@ pub struct ConfigSources { pub lastfm_shared_secret: ConfigSource, pub federation_enabled: ConfigSource, pub federation_network_id: ConfigSource, + pub federation_save_on_listen: ConfigSource, } impl Default for ConfigSources { @@ -166,6 +167,7 @@ impl Default for ConfigSources { lastfm_shared_secret: ConfigSource::Default, federation_enabled: ConfigSource::Default, federation_network_id: ConfigSource::Default, + federation_save_on_listen: ConfigSource::Default, } } } @@ -280,6 +282,9 @@ pub struct AppConfig { /// Federation network id — the shared secret every peer of the network /// uses to find the others. pub federation_network_id: String, + /// Whether a federated track requested for playback is imported into the + /// shared local library. This is a server-wide administrator policy. + pub federation_save_on_listen: bool, } impl Default for AppConfig { @@ -309,6 +314,7 @@ impl Default for AppConfig { lastfm_shared_secret: String::new(), federation_enabled: false, federation_network_id: String::new(), + federation_save_on_listen: false, } } } @@ -339,6 +345,7 @@ impl_env_overrides!( lastfm_shared_secret, federation_enabled, federation_network_id, + federation_save_on_listen, ); impl AppConfig { @@ -468,6 +475,7 @@ impl AppConfig { apply_db_field!(lastfm_shared_secret); apply_db_field!(federation_enabled); apply_db_field!(federation_network_id); + apply_db_field!(federation_save_on_listen); } } diff --git a/src/federation/client.rs b/src/federation/client.rs new file mode 100644 index 0000000..041925d --- /dev/null +++ b/src/federation/client.rs @@ -0,0 +1,536 @@ +//! Receiving side of music federation. +//! +//! User-facing identity is content-addressed. An `(owner, item_id)` pair is +//! only a source locator and several locators may resolve the same track. + +use std::collections::HashMap; + +use anyhow::{Context, Result}; +use music_dht::{ItemKind, LibraryItem, normalize_content_id}; +use serde::Serialize; +use serde_json::{Value, json}; +use sqlx::Row as _; +use tokio::io::AsyncReadExt; + +use super::{Federation, now_iso}; + +const MAX_CATALOG_BYTES: u64 = 4 * 1024 * 1024; + +#[derive(Debug, Clone, Serialize)] +pub struct TrackKeyDto { + pub content_id: String, +} + +#[derive(Debug, Clone, Serialize)] +pub struct ArtistKeyDto { + pub normalized_name: String, +} + +#[derive(Debug, Clone, Serialize)] +pub struct ArtistRefDto { + pub key: ArtistKeyDto, + pub name: String, + pub local_id: Option, +} + +#[derive(Debug, Clone, Serialize)] +pub struct ReleaseKeyDto { + pub normalized_title: String, + pub primary_artists: Vec, + pub release_type: Option, + pub year: Option, +} + +#[derive(Debug, Clone, Serialize)] +pub struct ReleaseRefDto { + pub key: ReleaseKeyDto, + pub local_id: Option, + pub title: String, +} + +#[derive(Debug, Clone, Serialize)] +pub struct FederationSourceDto { + pub owner: String, + pub item_id: String, +} + +#[derive(Debug, Clone, Serialize)] +pub struct LocalAvailabilityDto { + pub track_id: i64, + pub stream_url: String, +} + +#[derive(Debug, Clone, Serialize)] +pub struct TrackMetadataDto { + pub title: String, + pub artists: Vec, + pub featured_artists: Vec, + pub release: Option, + pub year: Option, + pub duration_seconds: Option, + pub track_number: Option, + pub disc_number: Option, + pub cover_url: Option, +} + +#[derive(Debug, Clone, Serialize)] +pub struct TrackAvailabilityDto { + pub state: &'static str, + pub local: Option, + pub federation: Vec, +} + +#[derive(Debug, Clone, Serialize)] +pub struct TrackDto { + pub key: TrackKeyDto, + pub metadata: TrackMetadataDto, + pub availability: TrackAvailabilityDto, +} + +#[derive(Debug, Clone, Serialize)] +pub struct SearchEvent { + pub search_id: String, + pub sequence: u64, + pub kind: &'static str, + pub peer: Option, + pub entity_key: Value, + pub entity: Value, +} + +impl Federation { + pub fn stream_artist_catalogs( + self: &std::sync::Arc, + name: String, + ) -> tokio::sync::mpsc::UnboundedReceiver> + { + let (sender, receiver) = tokio::sync::mpsc::unbounded_channel(); + let federation = std::sync::Arc::clone(self); + tokio::spawn(async move { + let result = async { + let service = federation.service().await?; + let normalized = music_dht::normalize_name(&name); + let outcome = service + .search_network(&name) + .await + .map_err(|err| anyhow::anyhow!("federated artist search failed: {err}"))?; + let owners: std::collections::HashSet<_> = outcome + .network_results + .iter() + .filter(|item| { + (item.kind == ItemKind::Artist && item.normalized_name == normalized) + || item + .artist_names + .iter() + .chain(&item.featured_artist_names) + .any(|artist| music_dht::normalize_name(artist) == normalized) + }) + .map(|item| item.owner) + .collect(); + for owner in owners { + let service = std::sync::Arc::clone(&service); + let sender = sender.clone(); + let name = name.clone(); + tokio::spawn(async move { + let result = tokio::time::timeout( + std::time::Duration::from_secs(8), + fetch_artist_catalog(&service, owner, &name), + ) + .await + .map_err(|_| anyhow::anyhow!("catalog request timed out")) + .and_then(|result| result) + .map(|artist| (owner.to_string(), artist)); + let _ = sender.send(result); + }); + } + Ok::<(), anyhow::Error>(()) + } + .await; + if let Err(err) = result { + let _ = sender.send(Err(err)); + } + }); + receiver + } + + /// Performs one bounded DHT search and returns entity upserts. The HTTP + /// layer streams each upsert independently; catalog fan-out can append + /// events to the same contract without changing the browser model. + pub async fn search_events(&self, search_id: &str, query: &str) -> Result> { + let query = query.trim(); + anyhow::ensure!(!query.is_empty(), "search query is empty"); + anyhow::ensure!(query.chars().count() <= 200, "search query is too long"); + + let started = std::time::Instant::now(); + let service = self.service().await?; + tracing::info!( + search_id, + query, + connected_peers = service.connected_peers().len(), + known_contacts = service.known_peers().len(), + "federated search started" + ); + let own = service.endpoint_id(); + let result = tokio::time::timeout( + std::time::Duration::from_secs(20), + service.search_network(query), + ) + .await + .map_err(|_| anyhow::anyhow!("federated search timed out after 20 seconds"))? + .map_err(|err| anyhow::anyhow!("federated search failed: {err}"))?; + tracing::info!( + search_id, + query, + local_results = result.local_results.len(), + network_results = result.network_results.len(), + queried_nodes = result.queried_nodes, + elapsed_ms = started.elapsed().as_millis() as u64, + "federated DHT search finished" + ); + let all_items: Vec = result + .local_results + .into_iter() + .chain(result.network_results) + .collect(); + + let pool = self.pool().await?; + let mut by_content: HashMap = HashMap::new(); + for item in all_items.iter().filter(|item| item.kind == ItemKind::Track) { + let Some(content_id) = item.content_id.as_deref().and_then(normalize_content_id) else { + // A globally usable track reference must be verifiable. + continue; + }; + let local = local_availability(&pool, &content_id).await?; + let source = FederationSourceDto { + owner: item.owner.to_string(), + item_id: hex(item.id.as_bytes()), + }; + let entry = by_content.entry(content_id.clone()).or_insert_with(|| { + track_from_item(content_id.clone(), item, local, item.owner == own) + }); + if !entry.availability.federation.iter().any(|candidate| { + candidate.owner == source.owner && candidate.item_id == source.item_id + }) { + entry.availability.federation.push(source); + } + if entry.availability.local.is_some() { + entry.availability.state = "local"; + } + } + + let mut tracks: Vec<_> = by_content.into_values().collect(); + tracks.sort_by(|left, right| { + left.metadata + .title + .to_lowercase() + .cmp(&right.metadata.title.to_lowercase()) + }); + + let mut events = Vec::with_capacity(all_items.len()); + for (index, track) in tracks.into_iter().enumerate() { + persist_track_ref(&pool, &track).await?; + let peer = track + .availability + .federation + .first() + .map(|source| source.owner.clone()); + events.push(SearchEvent { + search_id: search_id.to_owned(), + sequence: index as u64 + 1, + kind: "federation.track", + peer, + entity_key: serde_json::to_value(&track.key)?, + entity: serde_json::to_value(track)?, + }); + } + let mut artist_peers: HashMap)> = HashMap::new(); + let mut releases: HashMap = HashMap::new(); + for item in &all_items { + match item.kind { + ItemKind::Artist => { + let key = music_dht::normalize_name(&item.name); + let entry = artist_peers + .entry(key) + .or_insert_with(|| (item.name.clone(), Vec::new())); + let owner = item.owner.to_string(); + if !entry.1.contains(&owner) { + entry.1.push(owner); + } + } + ItemKind::Release => { + let artist_keys: Vec = item + .artist_names + .iter() + .map(|name| music_dht::normalize_name(name)) + .collect(); + let normalized_title = music_dht::normalize_name(&item.name); + let cover_url = all_items + .iter() + .find(|track| { + track.kind == ItemKind::Track + && track.release_title.as_deref().is_some_and(|title| { + music_dht::normalize_name(title) == normalized_title + }) + && track.year == item.year + }) + .map(|track| { + format!( + "/api/player/federation/tracks/artwork?owner={}&item_id={}", + track.owner, + hex(track.id.as_bytes()) + ) + }); + let key = format!( + "{}|{}|{}", + normalized_title, + artist_keys.join(","), + item.year.map_or_else(String::new, |year| year.to_string()) + ); + releases.entry(key.clone()).or_insert_with(|| { + json!({ + "key": { + "normalized_title": music_dht::normalize_name(&item.name), + "primary_artists": artist_keys, + "release_type": null, + "year": item.year, + }, + "title": item.name, + "artists": item.artist_names, + "year": item.year, + "cover_url": cover_url, + "sources": [{ + "owner": item.owner.to_string(), + "item_id": hex(item.id.as_bytes()), + }], + }) + }); + } + ItemKind::Track => {} + } + } + for (key, (name, peers)) in artist_peers { + let sequence = events.len() as u64 + 1; + events.push(SearchEvent { + search_id: search_id.to_owned(), + sequence, + kind: "federation.artist", + peer: peers.first().cloned(), + entity_key: json!({ "normalized_name": key }), + entity: json!({ + "key": { "normalized_name": key }, + "name": name, + "image_url": null, + "peers": peers, + }), + }); + } + for (key, release) in releases { + let sequence = events.len() as u64 + 1; + events.push(SearchEvent { + search_id: search_id.to_owned(), + sequence, + kind: "federation.release", + peer: None, + entity_key: json!({ "composite": key }), + entity: release, + }); + } + tracing::info!( + search_id, + query, + events = events.len(), + elapsed_ms = started.elapsed().as_millis() as u64, + "federated search response ready" + ); + Ok(events) + } +} + +async fn fetch_artist_catalog( + service: &music_dht::MusicDhtService, + owner: music_dht::EndpointId, + artist: &str, +) -> Result { + let mut stream = service + .open_stream(owner, super::CATALOG_ALPN) + .await + .map_err(|err| anyhow::anyhow!("cannot reach catalog peer: {err}"))?; + let mut request = serde_json::to_vec(&music_dht::catalog::CatalogRequest { + artist: artist.to_owned(), + want: Some("catalog".to_owned()), + ..Default::default() + })?; + request.push(b'\n'); + stream.send.write_all(&request).await?; + stream.send.finish()?; + let mut payload = Vec::new(); + stream + .recv + .take(MAX_CATALOG_BYTES + 1) + .read_to_end(&mut payload) + .await?; + anyhow::ensure!( + payload.len() as u64 <= MAX_CATALOG_BYTES, + "catalog response is too large" + ); + let response: music_dht::catalog::CatalogResponse = + serde_json::from_slice(&payload).context("invalid catalog response")?; + anyhow::ensure!( + response.ok, + "peer refused catalog: {}", + response.error.unwrap_or_else(|| "unknown error".to_owned()) + ); + response.artist.context("peer returned no artist catalog") +} + +fn track_from_item( + content_id: String, + item: &LibraryItem, + local: Option, + own: bool, +) -> TrackDto { + let owner = item.owner.to_string(); + let item_id = hex(item.id.as_bytes()); + let artists = artist_refs(&item.artist_names); + let featured_artists = artist_refs(&item.featured_artist_names); + let release = item.release_title.as_ref().map(|title| ReleaseRefDto { + key: ReleaseKeyDto { + normalized_title: music_dht::normalize_name(title), + primary_artists: item + .artist_names + .iter() + .map(|artist| music_dht::normalize_name(artist)) + .collect(), + release_type: None, + year: item.year, + }, + local_id: None, + title: title.clone(), + }); + let state = if local.is_some() || own { + "local" + } else { + "federated" + }; + TrackDto { + key: TrackKeyDto { content_id }, + metadata: TrackMetadataDto { + title: item.name.clone(), + artists, + featured_artists, + release, + year: item.year, + duration_seconds: item.duration_seconds, + track_number: item.track_number, + disc_number: item.disc_number, + cover_url: Some(format!( + "/api/player/federation/tracks/artwork?owner={owner}&item_id={item_id}" + )), + }, + availability: TrackAvailabilityDto { + state, + local, + federation: vec![FederationSourceDto { owner, item_id }], + }, + } +} + +fn artist_refs(names: &[String]) -> Vec { + names + .iter() + .map(|name| ArtistRefDto { + key: ArtistKeyDto { + normalized_name: music_dht::normalize_name(name), + }, + name: name.clone(), + local_id: None, + }) + .collect() +} + +async fn local_availability( + pool: &sqlx::PgPool, + content_id: &str, +) -> Result> { + let row = sqlx::query( + "SELECT t.id + FROM furumusic__federation_content_id_cache c + JOIN furumusic__track t ON t.audio_file_id = c.media_file_id + WHERE c.content_id = $1 AND t.is_hidden = false + LIMIT 1", + ) + .bind(content_id) + .fetch_optional(pool) + .await?; + Ok(row.map(|row| { + let track_id: i64 = row.get(0); + LocalAvailabilityDto { + track_id, + stream_url: format!("/api/player/stream/{track_id}"), + } + })) +} + +async fn persist_track_ref(pool: &sqlx::PgPool, track: &TrackDto) -> Result<()> { + let metadata = serde_json::to_value(&track.metadata)?; + let local_id = track + .availability + .local + .as_ref() + .map(|local| local.track_id); + let row = sqlx::query( + "INSERT INTO furumusic__track_ref + (content_id, local_track_id, title, release_title, year, + duration_seconds, metadata_json, metadata_authority, created_at, updated_at) + VALUES ($1, $2, $3, $4, $5, $6, $7, 'federation', $8, $8) + ON CONFLICT (content_id) DO UPDATE SET + local_track_id = COALESCE(furumusic__track_ref.local_track_id, EXCLUDED.local_track_id), + title = EXCLUDED.title, + release_title = EXCLUDED.release_title, + year = EXCLUDED.year, + duration_seconds = EXCLUDED.duration_seconds, + metadata_json = EXCLUDED.metadata_json, + updated_at = EXCLUDED.updated_at + RETURNING id", + ) + .bind(&track.key.content_id) + .bind(local_id) + .bind(&track.metadata.title) + .bind( + track + .metadata + .release + .as_ref() + .map(|release| &release.title), + ) + .bind(track.metadata.year) + .bind(track.metadata.duration_seconds) + .bind(metadata) + .bind(now_iso()) + .fetch_one(pool) + .await + .context("persisting content-addressed track reference failed")?; + let track_ref_id: i64 = row.get(0); + for source in &track.availability.federation { + sqlx::query( + "INSERT INTO furumusic__federation_track_source + (track_ref_id, owner_peer_id, item_id, last_seen_ms, metadata_json) + VALUES ($1, $2, $3, $4, $5) + ON CONFLICT (owner_peer_id, item_id) DO UPDATE SET + track_ref_id = EXCLUDED.track_ref_id, + last_seen_ms = EXCLUDED.last_seen_ms, + metadata_json = EXCLUDED.metadata_json", + ) + .bind(track_ref_id) + .bind(&source.owner) + .bind(&source.item_id) + .bind(chrono::Utc::now().timestamp_millis()) + .bind(json!({ "track": track.metadata })) + .execute(pool) + .await?; + } + Ok(()) +} + +fn hex(bytes: &[u8]) -> String { + bytes.iter().map(|byte| format!("{byte:02x}")).collect() +} diff --git a/src/federation/devices.rs b/src/federation/devices.rs index 2d43b68..5186832 100644 --- a/src/federation/devices.rs +++ b/src/federation/devices.rs @@ -771,6 +771,149 @@ pub async fn record_track_like( Ok(()) } +/// Records a like for a content-addressed track that need not be local yet. +/// The ordinary local likes table remains a projection for materialized +/// tracks; the durable user intent lives under `content_id`. +pub async fn record_content_like( + pool: &sqlx::PgPool, + user_id: i64, + content_id: &str, + liked: bool, + fed: Option, +) -> Result<()> { + let content_id = + music_dht::normalize_content_id(content_id).context("invalid track content id")?; + let fed = fed + .map(serde_json::from_value::) + .transpose() + .context("invalid federated track metadata")?; + record_local_op( + pool, + user_id, + SyncOpPayload::TrackLikeSet { + content_id, + liked, + fed: liked.then_some(fed).flatten(), + }, + ) + .await +} + +/// Adds a content-addressed track to a user's playlist without requiring a +/// local numeric track row. Materialization later fills `local_track_id` +/// without replacing this playlist item. +pub async fn record_content_playlist_add( + pool: &sqlx::PgPool, + user_id: i64, + playlist_id: i64, + content_id: &str, + position: i64, + fed: Option, +) -> Result<()> { + let content_id = + music_dht::normalize_content_id(content_id).context("invalid track content id")?; + let fed = fed + .map(serde_json::from_value::) + .transpose() + .context("invalid federated track metadata")?; + let title = playlist_title(pool, playlist_id) + .await? + .context("playlist not found")?; + let sync_id = ensure_local_playlist_sync_id(pool, user_id, playlist_id, &title).await?; + record_local_op( + pool, + user_id, + SyncOpPayload::PlaylistTrackAdded { + playlist_id: sync_id, + content_id, + position, + fed, + }, + ) + .await +} + +/// Reconciles content-addressed user state after a track becomes local. +/// History is intentionally not touched: only actual web playback writes it. +pub async fn materialize_content_state( + pool: &sqlx::PgPool, + content_id: &str, + track_id: i64, +) -> Result<()> { + let content_id = + music_dht::normalize_content_id(content_id).context("invalid track content id")?; + let likes = sqlx::query( + "SELECT user_id, liked, hlc_ms + FROM furumusic__fed_state_like WHERE content_id = $1", + ) + .bind(&content_id) + .fetch_all(pool) + .await?; + for row in likes { + let user_id: i64 = row.get("user_id"); + let liked: bool = row.get("liked"); + if liked { + sqlx::query( + "INSERT INTO furumusic__user_liked_track + (user_id, track_id, created_at) + VALUES ($1, $2, $3) + ON CONFLICT (user_id, track_id) DO NOTHING", + ) + .bind(user_id) + .bind(track_id) + .bind(iso_from_ms(row.get("hlc_ms"))) + .execute(pool) + .await?; + } + } + let playlist_items = sqlx::query( + "SELECT i.user_id, i.playlist_id, i.position + FROM furumusic__fed_state_playlist_item i + WHERE i.content_id = $1 AND i.present = true", + ) + .bind(&content_id) + .fetch_all(pool) + .await?; + for row in playlist_items { + let user_id: i64 = row.get("user_id"); + let sync_id: String = row.get("playlist_id"); + let playlist_id = ensure_playlist_for_item(pool, user_id, &sync_id).await?; + sqlx::query( + "INSERT INTO furumusic__playlist_track + (playlist_id, track_id, position, added_at, added_by_user_id) + SELECT $1, $2, $3, $4, $5 + WHERE NOT EXISTS ( + SELECT 1 FROM furumusic__playlist_track + WHERE playlist_id = $1 AND track_id = $2 + )", + ) + .bind(playlist_id) + .bind(track_id) + .bind(row.get::("position") as i32) + .bind(now_iso()) + .bind(user_id) + .execute(pool) + .await?; + } + sqlx::query( + "UPDATE furumusic__fed_state_like + SET local_track_id = $2 WHERE content_id = $1", + ) + .bind(&content_id) + .bind(track_id) + .execute(pool) + .await?; + sqlx::query( + "UPDATE furumusic__fed_state_playlist_item + SET local_track_id = $2 WHERE content_id = $1", + ) + .bind(&content_id) + .bind(track_id) + .execute(pool) + .await?; + Ok(()) +} + pub async fn record_playlist_created( pool: &sqlx::PgPool, user_id: i64, @@ -2664,6 +2807,31 @@ pub async fn record_web_active_transfer( Ok(()) } +pub async fn record_web_active_takeover( + pool: &sqlx::PgPool, + user_id: i64, + previous_device_id: &str, + state: serde_json::Value, +) -> Result<()> { + ensure_web_playback_target(pool, user_id, previous_device_id).await?; + let identity = ensure_identity(pool, user_id, "").await?; + let wire = playback_state_from_browser_json(pool, state).await?; + record_local_op( + pool, + user_id, + SyncOpPayload::PlaybackCommand { + target_device_id: previous_device_id.to_string(), + command: PlaybackCommand::ActiveChanged { + active_device_id: identity.device_id, + active_device_name: identity.name, + state: wire, + }, + }, + ) + .await?; + Ok(()) +} + async fn ensure_web_playback_target( pool: &sqlx::PgPool, user_id: i64, diff --git a/src/federation/mod.rs b/src/federation/mod.rs index ead27de..85c1b61 100644 --- a/src/federation/mod.rs +++ b/src/federation/mod.rs @@ -9,10 +9,12 @@ //! or download from other peers. //! //! Settings are the regular admin config entries (`federation_enabled`, -//! `federation_network_id`) and apply on the fly — saving the settings +//! `federation_network_id`, `federation_save_on_listen`) and apply on the fly — saving the settings //! starts, stops or re-joins the node without a server restart. +pub mod client; pub mod devices; +mod receive; mod serve; mod storage; @@ -51,6 +53,13 @@ struct ContentHashJob { file_path: String, } +#[derive(Clone)] +struct CachedArtwork { + bytes: Vec, + mime: String, + fetched_at: std::time::Instant, +} + #[derive(Debug, Clone)] struct TransportSample { at: String, @@ -207,8 +216,12 @@ pub struct Federation { data_dir: PathBuf, database_url: std::sync::Mutex, storage_dir: std::sync::Mutex, + save_on_listen: std::sync::atomic::AtomicBool, content_cache: std::sync::Mutex>, content_pending: std::sync::Mutex>, + prepared_cache: std::sync::Mutex>, + artwork_cache: std::sync::Mutex>, + download_locks: std::sync::Mutex>>>, pool: tokio::sync::OnceCell, running: tokio::sync::Mutex>, last_sync: std::sync::Mutex>, @@ -234,8 +247,12 @@ pub fn handle() -> Arc { data_dir: PathBuf::from(crate::media_paths::resolve_config_path("federation")), database_url: std::sync::Mutex::new(String::new()), storage_dir: std::sync::Mutex::new(String::new()), + save_on_listen: std::sync::atomic::AtomicBool::new(false), content_cache: std::sync::Mutex::new(Default::default()), content_pending: std::sync::Mutex::new(Default::default()), + prepared_cache: std::sync::Mutex::new(Default::default()), + artwork_cache: std::sync::Mutex::new(Default::default()), + download_locks: std::sync::Mutex::new(Default::default()), pool: tokio::sync::OnceCell::new(), running: tokio::sync::Mutex::new(None), last_sync: std::sync::Mutex::new(None), @@ -285,7 +302,8 @@ impl Federation { let mut effective = config.clone(); let rows = sqlx::query( "SELECT key, value FROM furumusic__config_entry - WHERE key IN ('federation_enabled', 'federation_network_id', 'agent_storage_dir')", + WHERE key IN ('federation_enabled', 'federation_network_id', + 'federation_save_on_listen', 'agent_storage_dir')", ) .fetch_all(&pool) .await @@ -304,6 +322,11 @@ impl Federation { } } "federation_network_id" => effective.federation_network_id = value, + "federation_save_on_listen" => { + if let Ok(parsed) = value.parse() { + effective.federation_save_on_listen = parsed; + } + } "agent_storage_dir" => { effective.agent_storage_dir = crate::media_paths::resolve_config_path(&value); } @@ -318,6 +341,10 @@ impl Federation { pub async fn apply(self: &Arc, config: &AppConfig) { *lock(&self.database_url) = config.database_url.clone(); *lock(&self.storage_dir) = config.agent_storage_dir.clone(); + self.save_on_listen.store( + config.federation_save_on_listen, + std::sync::atomic::Ordering::Relaxed, + ); let network = config.federation_network_id.trim().to_string(); if config.federation_enabled && !network.is_empty() { if let Err(err) = self.start(network, config.agent_storage_dir.clone()).await { @@ -905,6 +932,16 @@ impl Federation { ) .await } + + pub async fn fed_device_web_active_takeover( + &self, + user_id: i64, + previous_device_id: &str, + state: serde_json::Value, + ) -> Result<()> { + let pool = self.pool().await?; + devices::record_web_active_takeover(&pool, user_id, previous_device_id, state).await + } } async fn persist_content_id( diff --git a/src/federation/receive.rs b/src/federation/receive.rs new file mode 100644 index 0000000..abcfd58 --- /dev/null +++ b/src/federation/receive.rs @@ -0,0 +1,914 @@ +//! Verified federated audio download and trusted materialization. +//! +//! This module never writes inbox, processing-task, or review tables. Peer +//! metadata is the authority for this import path. + +use std::path::{Path, PathBuf}; +use std::str::FromStr; + +use anyhow::{Context, Result}; +use music_dht::EndpointId; +use serde::{Deserialize, Serialize}; +use sha2::{Digest as _, Sha256}; +use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt}; + +use super::{Federation, now_iso}; + +const MAX_LINE: usize = 4096; +const MAX_AUDIO_BYTES: u64 = 4 * 1024 * 1024 * 1024; +const MAX_IMAGE_BYTES: u64 = 16 * 1024 * 1024; +const ARTWORK_CACHE_TTL: std::time::Duration = std::time::Duration::from_secs(60 * 60); + +#[derive(Debug, Clone, Deserialize, Serialize)] +struct TrackMetadata { + title: String, + #[serde(default)] + artists: Vec, + #[serde(default)] + featured_artists: Vec, + #[serde(default)] + album_artists: Vec, + #[serde(default)] + release_title: String, + release_type: Option, + year: Option, + track_number: Option, + disc_number: Option, +} + +#[derive(Serialize)] +struct AudioRequest<'a> { + item_id: &'a str, + offset: u64, + want_cover: bool, +} + +#[derive(Deserialize)] +struct AudioHeader { + ok: bool, + #[serde(default)] + error: Option, + #[serde(default)] + mime_type: String, + #[serde(default)] + total_size: u64, + #[serde(default)] + metadata: Option, + #[serde(default)] + cover_size: u64, + #[serde(default)] + cover_mime: String, + #[serde(default)] + artist_image_size: u64, + #[serde(default)] + artist_image_mime: String, +} + +pub struct PreparedTrack { + pub local_track_id: Option, + pub stream_url: String, +} + +#[derive(Debug, Clone, Serialize)] +pub struct DownloadProgress { + pub phase: &'static str, + pub received: u64, + pub total: u64, +} + +struct Downloaded { + path: PathBuf, + mime: String, + metadata: TrackMetadata, + cover: Option<(Vec, String)>, + artist_image: Option<(Vec, String)>, +} + +impl Federation { + pub async fn discover_catalog_artwork( + &self, + artist: &str, + release: Option<&str>, + ) -> Result, String)>> { + let service = self.service().await?; + let normalized = music_dht::normalize_name(artist); + let outcome = service + .search_network(artist) + .await + .map_err(|err| anyhow::anyhow!("artwork peer discovery failed: {err}"))?; + let mut owners = Vec::new(); + for item in outcome.local_results.iter().chain(&outcome.network_results) { + let matches = (item.kind == music_dht::ItemKind::Artist + && item.normalized_name == normalized) + || item + .artist_names + .iter() + .chain(&item.featured_artist_names) + .any(|name| music_dht::normalize_name(name) == normalized); + let owner = item.owner.to_string(); + if matches && !owners.contains(&owner) { + owners.push(owner); + } + } + for owner in owners { + match tokio::time::timeout( + std::time::Duration::from_secs(5), + self.catalog_artwork(&owner, artist, release), + ) + .await + { + Ok(Ok(Some(artwork))) => return Ok(Some(artwork)), + Ok(Ok(None)) => {} + Ok(Err(err)) => { + tracing::debug!(%owner, %artist, ?release, "catalog artwork peer failed: {err:#}"); + } + Err(_) => { + tracing::debug!(%owner, %artist, ?release, "catalog artwork peer timed out"); + } + } + } + Ok(None) + } + + pub async fn catalog_artwork( + &self, + owner: &str, + artist: &str, + release: Option<&str>, + ) -> Result, String)>> { + let cache_key = format!("catalog:{owner}:{artist}:{}", release.unwrap_or_default()); + if let Some(cached) = cached_artwork(self, &cache_key) { + return Ok(Some(cached)); + } + let owner = EndpointId::from_str(owner).context("invalid federation owner")?; + let service = self.service().await?; + let mut stream = service + .open_stream(owner, super::CATALOG_ALPN) + .await + .map_err(|err| anyhow::anyhow!("cannot reach catalog peer: {err}"))?; + let mut request = serde_json::to_vec(&music_dht::catalog::CatalogRequest { + artist: artist.to_owned(), + want: Some(if release.is_some() { + "release_cover".to_owned() + } else { + "artist_image".to_owned() + }), + release: release.map(str::to_owned), + ..Default::default() + })?; + request.push(b'\n'); + stream.send.write_all(&request).await?; + stream.send.finish()?; + let header: music_dht::catalog::CatalogImageHeader = + serde_json::from_slice(&read_line(&mut stream.recv).await?) + .context("invalid catalog artwork response")?; + if !header.ok || header.size == 0 { + return Ok(None); + } + anyhow::ensure!( + header.size <= MAX_IMAGE_BYTES, + "catalog artwork is too large" + ); + anyhow::ensure!( + header.mime_type.starts_with("image/"), + "invalid catalog artwork mime type" + ); + let mut bytes = vec![0; header.size as usize]; + stream.recv.read_exact(&mut bytes).await?; + let artwork = (bytes, header.mime_type); + cache_artwork(self, cache_key, &artwork); + Ok(Some(artwork)) + } + + pub async fn track_artwork( + &self, + owner: &str, + item_id: &str, + ) -> Result, String)>> { + let cache_key = format!("{owner}:{item_id}"); + if let Some(cached) = cached_artwork(self, &cache_key) { + return Ok(Some(cached)); + } + let owner = EndpointId::from_str(owner).context("invalid federation owner")?; + anyhow::ensure!(item_id.len() == 64, "invalid federation item id"); + let service = self.service().await?; + let mut stream = service + .open_stream(owner, super::AUDIO_ALPN) + .await + .map_err(|err| anyhow::anyhow!("cannot reach track owner: {err}"))?; + write_line( + &mut stream.send, + &AudioRequest { + item_id, + offset: 0, + want_cover: true, + }, + ) + .await?; + stream.send.finish()?; + let header: AudioHeader = serde_json::from_slice(&read_line(&mut stream.recv).await?) + .context("invalid artwork response")?; + anyhow::ensure!( + header.ok, + "peer refused artwork: {}", + header.error.unwrap_or_else(|| "unknown error".to_string()) + ); + let artwork = read_segment( + &mut stream.recv, + header.cover_size, + &header.cover_mime, + "cover", + ) + .await?; + if let Some((bytes, mime)) = artwork { + let artwork = (bytes, mime); + cache_artwork(self, cache_key, &artwork); + Ok(Some(artwork)) + } else { + Ok(None) + } + } + + pub async fn prepare_content_with_progress( + self: &std::sync::Arc, + content_id: &str, + owner: &str, + item_id: &str, + mut progress: F, + ) -> Result + where + F: FnMut(DownloadProgress) + Send, + { + progress(DownloadProgress { + phase: "checking", + received: 0, + total: 0, + }); + let content_id = + music_dht::normalize_content_id(content_id).context("invalid content id")?; + let pool = self.pool().await?; + let token = content_id.trim_start_matches("b3:").to_owned(); + let download_lock = { + let mut locks = super::lock(&self.download_locks); + std::sync::Arc::clone( + locks + .entry(content_id.clone()) + .or_insert_with(|| std::sync::Arc::new(tokio::sync::Mutex::new(()))), + ) + }; + let _download_guard = download_lock.lock().await; + if let Some(track_id) = local_track_id(&pool, &content_id).await? { + progress(DownloadProgress { + phase: "ready", + received: 1, + total: 1, + }); + return Ok(PreparedTrack { + local_track_id: Some(track_id), + stream_url: format!("/api/player/stream/{track_id}"), + }); + } + + let save = self + .save_on_listen + .load(std::sync::atomic::Ordering::Relaxed); + if !save + && let Some((_path, _mime)) = super::lock(&self.prepared_cache).get(&token).cloned() + { + return Ok(PreparedTrack { + local_track_id: None, + stream_url: format!("/api/player/federation/cache/{token}"), + }); + } + let owner = EndpointId::from_str(owner).context("invalid federation owner")?; + anyhow::ensure!(item_id.len() == 64, "invalid federation item id"); + let service = self.service().await?; + let storage_root = if save { + PathBuf::from(super::lock(&self.storage_dir).clone()) + } else { + PathBuf::from(crate::media_paths::resolve_config_path("federation-cache")) + }; + anyhow::ensure!( + !storage_root.as_os_str().is_empty(), + "media storage directory is not configured" + ); + let dir = storage_root.join("federation"); + tokio::fs::create_dir_all(&dir).await?; + let downloaded = + download(&service, owner, item_id, &content_id, &dir, &mut progress).await?; + if !save { + super::lock(&self.prepared_cache) + .insert(token.clone(), (downloaded.path, downloaded.mime)); + return Ok(PreparedTrack { + local_track_id: None, + stream_url: format!("/api/player/federation/cache/{token}"), + }); + } + progress(DownloadProgress { + phase: "saving", + received: 1, + total: 1, + }); + let track_id = materialize(&pool, &storage_root, &content_id, downloaded).await?; + // The normal periodic sync will publish it; this immediate sync keeps + // save-on-listen useful to the federation without waiting a minute. + if let Err(err) = self.sync_now().await { + tracing::warn!(track_id, "post-import federation publish failed: {err:#}"); + } + Ok(PreparedTrack { + local_track_id: Some(track_id), + stream_url: format!("/api/player/stream/{track_id}"), + }) + } + + pub fn prepared_cache_file(&self, token: &str) -> Option<(PathBuf, String)> { + if token.len() != 64 || !token.bytes().all(|byte| byte.is_ascii_hexdigit()) { + return None; + } + super::lock(&self.prepared_cache).get(token).cloned() + } +} + +fn cached_artwork(federation: &Federation, key: &str) -> Option<(Vec, String)> { + let mut cache = super::lock(&federation.artwork_cache); + let cached = cache.get(key).cloned()?; + if cached.fetched_at.elapsed() > ARTWORK_CACHE_TTL { + cache.remove(key); + return None; + } + Some((cached.bytes, cached.mime)) +} + +fn cache_artwork(federation: &Federation, key: String, artwork: &(Vec, String)) { + let mut cache = super::lock(&federation.artwork_cache); + if cache.len() >= 512 { + cache.retain(|_, value| value.fetched_at.elapsed() <= ARTWORK_CACHE_TTL); + if cache.len() >= 512 + && let Some(oldest) = cache + .iter() + .min_by_key(|(_, value)| value.fetched_at) + .map(|(key, _)| key.clone()) + { + cache.remove(&oldest); + } + } + cache.insert( + key, + super::CachedArtwork { + bytes: artwork.0.clone(), + mime: artwork.1.clone(), + fetched_at: std::time::Instant::now(), + }, + ); +} + +async fn download( + service: &music_dht::MusicDhtService, + owner: EndpointId, + item_id: &str, + content_id: &str, + dir: &Path, + progress: &mut (impl FnMut(DownloadProgress) + Send), +) -> Result { + progress(DownloadProgress { + phase: "connecting", + received: 0, + total: 0, + }); + let mut stream = service + .open_stream(owner, super::AUDIO_ALPN) + .await + .map_err(|err| anyhow::anyhow!("cannot reach track owner: {err}"))?; + write_line( + &mut stream.send, + &AudioRequest { + item_id, + offset: 0, + want_cover: true, + }, + ) + .await?; + stream.send.finish()?; + let header: AudioHeader = serde_json::from_slice(&read_line(&mut stream.recv).await?) + .context("invalid audio response")?; + anyhow::ensure!( + header.ok, + "peer refused audio: {}", + header.error.unwrap_or_else(|| "unknown error".to_string()) + ); + anyhow::ensure!( + header.total_size > 0 && header.total_size <= MAX_AUDIO_BYTES, + "invalid federated audio size" + ); + progress(DownloadProgress { + phase: "downloading", + received: 0, + total: header.total_size, + }); + let metadata = header.metadata.context("peer returned no track metadata")?; + let cover = read_segment( + &mut stream.recv, + header.cover_size, + &header.cover_mime, + "cover", + ) + .await?; + let artist_image = read_segment( + &mut stream.recv, + header.artist_image_size, + &header.artist_image_mime, + "artist image", + ) + .await?; + + let stem = content_id.trim_start_matches("b3:"); + let extension = audio_extension(&header.mime_type); + let final_path = dir.join(format!("{stem}.{extension}")); + let part_path = dir.join(format!(".{stem}.{extension}.part")); + let mut file = tokio::fs::File::create(&part_path).await?; + let mut hasher = blake3::Hasher::new(); + let mut received = 0u64; + let mut buf = vec![0u8; 64 * 1024]; + while received < header.total_size { + let remaining = (header.total_size - received).min(buf.len() as u64) as usize; + let count = stream.recv.read(&mut buf[..remaining]).await?.unwrap_or(0); + anyhow::ensure!(count > 0, "audio stream ended early"); + file.write_all(&buf[..count]).await?; + hasher.update(&buf[..count]); + received += count as u64; + progress(DownloadProgress { + phase: "downloading", + received, + total: header.total_size, + }); + } + file.flush().await?; + drop(file); + let actual = format!("b3:{}", hasher.finalize().to_hex()); + progress(DownloadProgress { + phase: "verifying", + received, + total: header.total_size, + }); + if actual != content_id { + let _ = tokio::fs::remove_file(&part_path).await; + anyhow::bail!("downloaded audio content id mismatch"); + } + tokio::fs::rename(&part_path, &final_path).await?; + Ok(Downloaded { + path: final_path, + mime: header.mime_type, + metadata, + cover, + artist_image, + }) +} + +async fn materialize( + pool: &sqlx::PgPool, + storage_root: &Path, + content_id: &str, + downloaded: Downloaded, +) -> Result { + if let Some(track_id) = local_track_id(pool, content_id).await? { + return Ok(track_id); + } + let bytes = tokio::fs::read(&downloaded.path).await?; + let sha256 = format!("{:x}", Sha256::digest(&bytes)); + let relative = downloaded + .path + .strip_prefix(storage_root) + .unwrap_or(&downloaded.path) + .to_string_lossy() + .into_owned(); + let mut tx = pool.begin().await?; + let existing_media: Option = sqlx::query_scalar( + "SELECT id FROM furumusic__media_file + WHERE file_type = 'audio' AND sha256_hash = $1 LIMIT 1", + ) + .bind(&sha256) + .fetch_optional(&mut *tx) + .await?; + let media_id = match existing_media { + Some(id) => id, + None => { + sqlx::query_scalar( + "INSERT INTO furumusic__media_file + (file_type, file_path, original_filename, mime_type, + file_size_bytes, sha256_hash, audio_format, + uploaded_by_user_id, uploader_name, created_at) + VALUES ('audio', $1, $2, $3, $4, $5, $6, NULL, 'Federation', $7) + RETURNING id", + ) + .bind(&relative) + .bind( + downloaded + .path + .file_name() + .and_then(|name| name.to_str()) + .unwrap_or("federated-audio"), + ) + .bind(&downloaded.mime) + .bind(bytes.len() as i64) + .bind(&sha256) + .bind(audio_extension(&downloaded.mime)) + .bind(now_iso()) + .fetch_one(&mut *tx) + .await? + } + }; + sqlx::query( + "INSERT INTO furumusic__federation_content_id_cache + (media_file_id, sha256_hash, content_id, updated_at) + VALUES ($1, $2, $3, $4) + ON CONFLICT (media_file_id) DO UPDATE SET + sha256_hash = EXCLUDED.sha256_hash, + content_id = EXCLUDED.content_id, + updated_at = EXCLUDED.updated_at", + ) + .bind(media_id) + .bind(&sha256) + .bind(content_id) + .bind(now_iso()) + .execute(&mut *tx) + .await?; + if let Some(track_id) = sqlx::query_scalar::<_, i64>( + "SELECT id FROM furumusic__track WHERE audio_file_id = $1 LIMIT 1", + ) + .bind(media_id) + .fetch_optional(&mut *tx) + .await? + { + tx.commit().await?; + return Ok(track_id); + } + + let release_title = nonempty(&downloaded.metadata.release_title).unwrap_or("Unknown release"); + let release_sort = music_dht::normalize_name(release_title); + let release_id: i64 = if let Some(id) = sqlx::query_scalar( + "SELECT id FROM furumusic__release + WHERE title_sort = $1 AND year IS NOT DISTINCT FROM $2 + ORDER BY id LIMIT 1", + ) + .bind(&release_sort) + .bind(downloaded.metadata.year) + .fetch_optional(&mut *tx) + .await? + { + id + } else { + sqlx::query_scalar( + "INSERT INTO furumusic__release + (title, title_sort, release_type, year, is_hidden, model_name, + created_at, updated_at) + VALUES ($1, $2, $3, $4, false, NULL, $5, $5) + RETURNING id", + ) + .bind(release_title) + .bind(&release_sort) + .bind( + downloaded + .metadata + .release_type + .as_deref() + .unwrap_or("album"), + ) + .bind(downloaded.metadata.year) + .bind(now_iso()) + .fetch_one(&mut *tx) + .await? + }; + let duration: f64 = sqlx::query_scalar( + "SELECT COALESCE(duration_seconds, 0) + FROM furumusic__track_ref WHERE content_id = $1", + ) + .bind(content_id) + .fetch_optional(&mut *tx) + .await? + .unwrap_or(0.0); + let track_id: i64 = sqlx::query_scalar( + "INSERT INTO furumusic__track + (title, title_sort, release_id, track_number, disc_number, + duration_seconds, audio_file_id, year, is_hidden, model_name, + created_at, updated_at) + VALUES ($1, $2, $3, $4, $5, $6, $7, $8, false, NULL, $9, $9) + RETURNING id", + ) + .bind(&downloaded.metadata.title) + .bind(music_dht::normalize_name(&downloaded.metadata.title)) + .bind(release_id) + .bind(downloaded.metadata.track_number) + .bind(downloaded.metadata.disc_number) + .bind(duration) + .bind(media_id) + .bind(downloaded.metadata.year) + .bind(now_iso()) + .fetch_one(&mut *tx) + .await?; + + let main_artists = if downloaded.metadata.artists.is_empty() { + &downloaded.metadata.album_artists + } else { + &downloaded.metadata.artists + }; + let mut main_artist_ids = Vec::new(); + for (position, name) in main_artists.iter().enumerate() { + let artist_id = ensure_artist(&mut tx, name).await?; + main_artist_ids.push(artist_id); + link_track_artist(&mut tx, track_id, artist_id, "main", position as i32).await?; + link_release_artist(&mut tx, release_id, artist_id, position as i32).await?; + } + for (position, name) in downloaded.metadata.featured_artists.iter().enumerate() { + let artist_id = ensure_artist(&mut tx, name).await?; + link_track_artist(&mut tx, track_id, artist_id, "featuring", position as i32).await?; + } + sqlx::query( + "UPDATE furumusic__track_ref + SET local_track_id = $2, metadata_authority = 'federation', updated_at = $3 + WHERE content_id = $1", + ) + .bind(content_id) + .bind(track_id) + .bind(now_iso()) + .execute(&mut *tx) + .await?; + sqlx::query( + "UPDATE furumusic__fed_state_like + SET local_track_id = $2 WHERE content_id = $1", + ) + .bind(content_id) + .bind(track_id) + .execute(&mut *tx) + .await?; + sqlx::query( + "UPDATE furumusic__fed_state_playlist_item + SET local_track_id = $2 WHERE content_id = $1", + ) + .bind(content_id) + .bind(track_id) + .execute(&mut *tx) + .await?; + tx.commit().await?; + super::devices::materialize_content_state(pool, content_id, track_id).await?; + + // Artwork is non-authoritative for identity and may be installed after + // the audio transaction. Failure does not invalidate a verified track. + if let Some((bytes, mime)) = downloaded.cover + && let Err(err) = + install_release_artwork(pool, storage_root, release_id, &bytes, &mime).await + { + tracing::warn!( + release_id, + "failed to install federated release artwork: {err:#}" + ); + } + if let Some((bytes, mime)) = downloaded.artist_image + && let Err(err) = + install_artist_artwork(pool, storage_root, &main_artist_ids, &bytes, &mime).await + { + tracing::warn!("failed to install federated artist artwork: {err:#}"); + } + Ok(track_id) +} + +async fn install_release_artwork( + pool: &sqlx::PgPool, + storage_root: &Path, + release_id: i64, + bytes: &[u8], + mime: &str, +) -> Result<()> { + let media_id = persist_artwork(pool, storage_root, bytes, mime).await?; + sqlx::query( + "UPDATE furumusic__release + SET cover_file_id = $1, updated_at = $3 + WHERE id = $2 AND cover_file_id IS NULL", + ) + .bind(media_id) + .bind(release_id) + .bind(now_iso()) + .execute(pool) + .await?; + Ok(()) +} + +async fn install_artist_artwork( + pool: &sqlx::PgPool, + storage_root: &Path, + artist_ids: &[i64], + bytes: &[u8], + mime: &str, +) -> Result<()> { + if artist_ids.is_empty() { + return Ok(()); + } + let media_id = persist_artwork(pool, storage_root, bytes, mime).await?; + sqlx::query( + "UPDATE furumusic__artist + SET image_file_id = $1, updated_at = $3 + WHERE id = ANY($2) AND image_file_id IS NULL", + ) + .bind(media_id) + .bind(artist_ids) + .bind(now_iso()) + .execute(pool) + .await?; + Ok(()) +} + +async fn persist_artwork( + pool: &sqlx::PgPool, + storage_root: &Path, + bytes: &[u8], + mime: &str, +) -> Result { + anyhow::ensure!(!bytes.is_empty() && bytes.len() as u64 <= MAX_IMAGE_BYTES); + anyhow::ensure!(mime.starts_with("image/"), "invalid artwork mime type"); + let hash = format!("{:x}", Sha256::digest(bytes)); + if let Some(id) = sqlx::query_scalar( + "SELECT id FROM furumusic__media_file + WHERE file_type = 'cover_art' AND sha256_hash = $1 LIMIT 1", + ) + .bind(&hash) + .fetch_optional(pool) + .await? + { + return Ok(id); + } + let extension = image_extension(mime); + let filename = format!("federation-artwork-{}.{}", &hash[..16], extension); + let dir = storage_root.join("federation").join("artwork"); + tokio::fs::create_dir_all(&dir).await?; + let path = dir.join(&filename); + tokio::fs::write(&path, bytes).await?; + let relative = path + .strip_prefix(storage_root) + .unwrap_or(&path) + .to_string_lossy() + .into_owned(); + let id = sqlx::query_scalar( + "INSERT INTO furumusic__media_file + (file_type, file_path, original_filename, mime_type, + file_size_bytes, sha256_hash, uploaded_by_user_id, + uploader_name, created_at) + VALUES ('cover_art', $1, $2, $3, $4, $5, NULL, 'Federation', $6) + RETURNING id", + ) + .bind(&relative) + .bind(&filename) + .bind(mime) + .bind(bytes.len() as i64) + .bind(&hash) + .bind(now_iso()) + .fetch_one(pool) + .await?; + if let Err(err) = crate::agent::cover_variants::ensure_cover_variants(&path).await { + tracing::warn!( + media_id = id, + "failed to generate federated artwork variants: {err}" + ); + } + Ok(id) +} + +fn image_extension(mime: &str) -> &'static str { + match mime { + "image/png" => "png", + "image/webp" => "webp", + "image/gif" => "gif", + "image/avif" => "avif", + _ => "jpg", + } +} + +async fn ensure_artist(tx: &mut sqlx::Transaction<'_, sqlx::Postgres>, name: &str) -> Result { + let normalized = music_dht::normalize_name(name); + if let Some(id) = sqlx::query_scalar( + "SELECT id FROM furumusic__artist WHERE name_sort = $1 ORDER BY id LIMIT 1", + ) + .bind(&normalized) + .fetch_optional(&mut **tx) + .await? + { + return Ok(id); + } + Ok(sqlx::query_scalar( + "INSERT INTO furumusic__artist + (name, name_sort, is_hidden, model_name, created_at, updated_at) + VALUES ($1, $2, false, NULL, $3, $3) RETURNING id", + ) + .bind(name) + .bind(normalized) + .bind(now_iso()) + .fetch_one(&mut **tx) + .await?) +} + +async fn link_track_artist( + tx: &mut sqlx::Transaction<'_, sqlx::Postgres>, + track_id: i64, + artist_id: i64, + role: &str, + position: i32, +) -> Result<()> { + sqlx::query( + "INSERT INTO furumusic__track_artist + (track_id, artist_id, role, position) VALUES ($1, $2, $3, $4)", + ) + .bind(track_id) + .bind(artist_id) + .bind(role) + .bind(position) + .execute(&mut **tx) + .await?; + Ok(()) +} + +async fn link_release_artist( + tx: &mut sqlx::Transaction<'_, sqlx::Postgres>, + release_id: i64, + artist_id: i64, + position: i32, +) -> Result<()> { + sqlx::query( + "INSERT INTO furumusic__release_artist + (release_id, artist_id, position) + SELECT $1, $2, $3 WHERE NOT EXISTS ( + SELECT 1 FROM furumusic__release_artist + WHERE release_id = $1 AND artist_id = $2 + )", + ) + .bind(release_id) + .bind(artist_id) + .bind(position) + .execute(&mut **tx) + .await?; + Ok(()) +} + +async fn local_track_id(pool: &sqlx::PgPool, content_id: &str) -> Result> { + Ok(sqlx::query_scalar( + "SELECT t.id + FROM furumusic__federation_content_id_cache c + JOIN furumusic__track t ON t.audio_file_id = c.media_file_id + WHERE c.content_id = $1 AND t.is_hidden = false LIMIT 1", + ) + .bind(content_id) + .fetch_optional(pool) + .await?) +} + +async fn read_line(reader: &mut R) -> Result> { + let mut line = Vec::new(); + let mut byte = [0u8; 1]; + loop { + anyhow::ensure!( + reader.read_exact(&mut byte).await.is_ok(), + "stream ended early" + ); + if byte[0] == b'\n' { + return Ok(line); + } + line.push(byte[0]); + anyhow::ensure!(line.len() <= MAX_LINE, "protocol line too large"); + } +} + +async fn write_line(writer: &mut W, value: &impl Serialize) -> Result<()> { + let mut line = serde_json::to_vec(value)?; + line.push(b'\n'); + writer.write_all(&line).await?; + Ok(()) +} + +async fn read_segment( + reader: &mut R, + size: u64, + mime: &str, + label: &str, +) -> Result, String)>> { + if size == 0 { + return Ok(None); + } + anyhow::ensure!(size <= MAX_IMAGE_BYTES, "{label} is too large"); + let mut bytes = vec![0u8; size as usize]; + reader.read_exact(&mut bytes).await?; + Ok(Some((bytes, mime.to_owned()))) +} + +fn audio_extension(mime: &str) -> &'static str { + match mime { + "audio/mpeg" => "mp3", + "audio/flac" => "flac", + "audio/ogg" => "ogg", + "audio/opus" => "opus", + "audio/wav" => "wav", + "audio/mp4" => "m4a", + "audio/aac" => "aac", + _ => "bin", + } +} + +fn nonempty(value: &str) -> Option<&str> { + (!value.trim().is_empty()).then_some(value.trim()) +} diff --git a/src/music/mod.rs b/src/music/mod.rs index 86d5dfd..9b44fe4 100644 --- a/src/music/mod.rs +++ b/src/music/mod.rs @@ -2223,6 +2223,105 @@ pub mod db_migrations { &[Operation::custom(ensure_federation_content_id_cache).build()]; } + #[cot::db::migrations::migration_op] + async fn create_content_addressed_music_refs( + ctx: migrations::MigrationContext<'_>, + ) -> cot::db::Result<()> { + // A track reference is durable user-facing identity. `local_track_id` + // is availability, not identity: it may become non-NULL after a + // federated track is materialized without changing likes/playlists. + ctx.db + .raw( + "CREATE TABLE IF NOT EXISTS furumusic__track_ref ( + id BIGSERIAL PRIMARY KEY, + content_id TEXT NOT NULL UNIQUE, + local_track_id BIGINT UNIQUE, + title TEXT NOT NULL, + release_title TEXT, + year INTEGER, + duration_seconds DOUBLE PRECISION, + metadata_json JSONB NOT NULL DEFAULT '{}'::jsonb, + metadata_authority TEXT NOT NULL DEFAULT 'local', + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL + )", + ) + .await?; + ctx.db + .raw( + "CREATE INDEX IF NOT EXISTS idx_track_ref_local_track + ON furumusic__track_ref (local_track_id) + WHERE local_track_id IS NOT NULL", + ) + .await?; + ctx.db + .raw( + "CREATE TABLE IF NOT EXISTS furumusic__federation_track_source ( + track_ref_id BIGINT NOT NULL REFERENCES furumusic__track_ref(id) + ON DELETE CASCADE, + owner_peer_id TEXT NOT NULL, + item_id TEXT NOT NULL, + last_seen_ms BIGINT NOT NULL, + metadata_json JSONB NOT NULL DEFAULT '{}'::jsonb, + PRIMARY KEY (owner_peer_id, item_id) + )", + ) + .await?; + ctx.db + .raw( + "CREATE INDEX IF NOT EXISTS idx_federation_track_source_ref + ON furumusic__federation_track_source (track_ref_id, last_seen_ms DESC)", + ) + .await?; + ctx.db + .raw( + "ALTER TABLE furumusic__user_liked_track + ADD COLUMN IF NOT EXISTS track_ref_id BIGINT + REFERENCES furumusic__track_ref(id)", + ) + .await?; + ctx.db + .raw( + "CREATE UNIQUE INDEX IF NOT EXISTS idx_user_liked_track_ref_uniq + ON furumusic__user_liked_track (user_id, track_ref_id) + WHERE track_ref_id IS NOT NULL", + ) + .await?; + ctx.db + .raw( + "ALTER TABLE furumusic__playlist_track + ADD COLUMN IF NOT EXISTS track_ref_id BIGINT + REFERENCES furumusic__track_ref(id)", + ) + .await?; + ctx.db + .raw( + "CREATE INDEX IF NOT EXISTS idx_playlist_track_ref + ON furumusic__playlist_track (track_ref_id) + WHERE track_ref_id IS NOT NULL", + ) + .await?; + // History deliberately remains local-track based. Only the web + // player's existing playback report records history and triggers + // Last.fm scrobbling. + Ok(()) + } + + #[derive(Debug, Copy, Clone)] + pub struct M0040CreateContentAddressedMusicRefs; + + impl migrations::Migration for M0040CreateContentAddressedMusicRefs { + const APP_NAME: &'static str = "furumusic"; + const MIGRATION_NAME: &'static str = "m_0040_create_content_addressed_music_refs"; + const DEPENDENCIES: &'static [migrations::MigrationDependency] = + &[migrations::MigrationDependency::migration( + "furumusic", + "m_0039_ensure_federation_content_id_cache", + )]; + const OPERATIONS: &'static [Operation] = + &[Operation::custom(create_content_addressed_music_refs).build()]; + } + pub const MIGRATIONS: &[&SyncDynMigration] = &[ &M0006CreateMediaFile, &M0007CreateArtist, @@ -2253,5 +2352,6 @@ pub mod db_migrations { &M0037CreatePlaylistShareLinks, &M0038CreateFedDeviceSync, &M0039EnsureFederationContentIdCache, + &M0040CreateContentAddressedMusicRefs, ]; } diff --git a/src/player/dto.rs b/src/player/dto.rs index 0899b83..b7457b7 100644 --- a/src/player/dto.rs +++ b/src/player/dto.rs @@ -51,6 +51,7 @@ pub(super) struct ArtistRef { #[derive(Debug, Clone, Serialize, JsonSchema)] pub(super) struct TrackItem { pub(super) id: i64, + pub(super) content_id: Option, pub(super) title: String, pub(super) track_number: Option, pub(super) disc_number: Option, @@ -561,6 +562,46 @@ pub(super) struct LikeStatus { pub(super) liked: bool, } +#[derive(Debug, Deserialize, JsonSchema)] +pub(super) struct ContentTrackMutation { + pub(super) content_id: String, + pub(super) liked: Option, + pub(super) playlist_id: Option, + pub(super) position: Option, + pub(super) federation: Option, +} + +#[derive(Debug, Deserialize, JsonSchema)] +pub(super) struct PrepareFederatedTrackRequest { + pub(super) content_id: String, + pub(super) owner: String, + pub(super) item_id: String, +} + +#[derive(Debug, Deserialize, JsonSchema)] +pub(super) struct FederationArtworkQuery { + pub(super) owner: String, + pub(super) item_id: String, +} + +#[derive(Debug, Deserialize, JsonSchema)] +pub(super) struct FederationArtistQuery { + pub(super) name: String, +} + +#[derive(Debug, Deserialize, JsonSchema)] +pub(super) struct FederationCatalogArtworkQuery { + pub(super) owner: String, + pub(super) artist: String, + pub(super) release: Option, +} + +#[derive(Debug, Deserialize, JsonSchema)] +pub(super) struct FederationArtworkDiscoveryQuery { + pub(super) artist: String, + pub(super) release: Option, +} + #[derive(Debug, Serialize, JsonSchema)] pub(super) struct LikedIds { pub(super) track_ids: Vec, diff --git a/src/player/mod.rs b/src/player/mod.rs index 5d00b6f..546c5d6 100644 --- a/src/player/mod.rs +++ b/src/player/mod.rs @@ -1,6 +1,7 @@ use std::collections::{HashMap, HashSet, VecDeque}; use std::sync::{Arc, Mutex, OnceLock}; +use bytes::Bytes; use cot::db::Database; use cot::http::StatusCode; use cot::http::header::{ @@ -13,6 +14,7 @@ use cot::router::method::{delete, get, post}; use cot::router::{Route, Router}; use cot::session::Session; use cot::{App, Body, Template}; +use sqlx::Row as _; use crate::auth; use crate::config::AppConfig; @@ -223,6 +225,22 @@ impl PlayerDeviceHub { last_seen_ms: now, }, ); + // Match the trusted-device playback contract used by the TUI: a + // background/stale active snapshot must not steal playback from a + // browser that is actively playing. An explicit web handoff changes + // `active_device_by_user` to the federated virtual device before the + // snapshot arrives, so it still passes through here. + let local_playback_is_protected = state + .active_device_by_user + .get(&user_id) + .is_some_and(|active_id| !is_fed_virtual_device_id(active_id)) + && state + .playback_state_by_user + .get(&user_id) + .is_some_and(|playback| playback.track.is_some() && !playback.paused); + if active && local_playback_is_protected { + return Ok(()); + } let should_update_playback = active || state .active_device_by_user @@ -1097,6 +1115,76 @@ mod device_tests { Some("fed:dev_e5ffc3b65642770c26c53ecf".to_string()) ); } + + #[test] + fn federated_snapshot_does_not_steal_active_browser_playback() { + let hub = PlayerDeviceHub::default(); + let user_id = 7; + { + let mut state = hub.state.lock().expect("device hub"); + state.devices_by_user.entry(user_id).or_default().insert( + "browser".to_string(), + PlayerDevice { + id: "browser".to_string(), + name: "Browser".to_string(), + kind: "computer".to_string(), + last_seen_ms: current_millis(), + }, + ); + state + .active_device_by_user + .insert(user_id, "browser".to_string()); + state.playback_state_by_user.insert( + user_id, + PlayerDevicePlaybackStateDto { + track: Some(serde_json::json!({"id": 1})), + tracks: vec![], + index: 0, + position_seconds: 10.0, + duration_seconds: 100.0, + paused: false, + shuffle: false, + repeat_mode: "off".to_string(), + volume: 0.7, + updated_at_ms: current_millis(), + }, + ); + } + + hub.apply_fed_playback_state_json( + user_id, + "remote", + "Remote", + true, + serde_json::json!({ + "track": {"id": 2}, + "tracks": [], + "index": 0, + "position_seconds": 0.0, + "duration_seconds": 100.0, + "paused": false, + "shuffle": false, + "repeat_mode": "off", + "volume": 0.7 + }), + ) + .expect("valid snapshot"); + + let state = hub.state.lock().expect("device hub"); + assert_eq!( + state + .active_device_by_user + .get(&user_id) + .map(String::as_str), + Some("browser") + ); + assert!( + state + .devices_by_user + .get(&user_id) + .is_some_and(|devices| devices.contains_key("fed:remote")) + ); + } } #[derive(Debug, sqlx::FromRow)] @@ -1153,12 +1241,16 @@ async fn me_handler( return Ok(json_error(StatusCode::UNAUTHORIZED, "not authenticated")); }; - let liked_tracks: (i64,) = - sqlx::query_as("SELECT COUNT(*) FROM furumusic__user_liked_track WHERE user_id = $1") - .bind(user.id) - .fetch_one(pool) - .await - .map_err(|e| cot::Error::internal(e.to_string()))?; + let liked_tracks: (i64,) = sqlx::query_as( + "SELECT + (SELECT COUNT(*) FROM furumusic__user_liked_track WHERE user_id = $1) + + (SELECT COUNT(*) FROM furumusic__fed_state_like + WHERE user_id = $1 AND liked = true AND local_track_id IS NULL)", + ) + .bind(user.id) + .fetch_one(pool) + .await + .map_err(|e| cot::Error::internal(e.to_string()))?; let playlists: (i64,) = sqlx::query_as("SELECT COUNT(*) FROM furumusic__playlist WHERE owner_id = $1") @@ -2683,6 +2775,7 @@ async fn load_user_upload_tracks( UserUploadTrack { track: TrackItem { id: row.id, + content_id: None, title: row.title, track_number: row.track_number, disc_number: row.disc_number, @@ -3733,12 +3826,16 @@ async fn playlists_handler( }; // Count liked tracks for the virtual Likes playlist - let likes_count: (i64,) = - sqlx::query_as("SELECT COUNT(*) FROM furumusic__user_liked_track WHERE user_id = $1") - .bind(user.id) - .fetch_one(pool) - .await - .map_err(|e| cot::Error::internal(e.to_string()))?; + let likes_count: (i64,) = sqlx::query_as( + "SELECT + (SELECT COUNT(*) FROM furumusic__user_liked_track WHERE user_id = $1) + + (SELECT COUNT(*) FROM furumusic__fed_state_like + WHERE user_id = $1 AND liked = true AND local_track_id IS NULL)", + ) + .bind(user.id) + .fetch_one(pool) + .await + .map_err(|e| cot::Error::internal(e.to_string()))?; let mut cards = vec![PlaylistCard { id: -1, @@ -3778,16 +3875,31 @@ async fn playlists_handler( .await .map_err(|e| cot::Error::internal(e.to_string()))?; - cards.extend(rows.into_iter().map(|r| PlaylistCard { - id: r.id, - title: r.title, - track_count: r.track_count, - is_own: r.is_own, - owner_name: Some(r.owner_name), - is_public: r.is_public, - is_saved: r.is_saved, - kind: "user".to_string(), - })); + for r in rows { + let federated_count: i64 = sqlx::query_scalar( + "SELECT COUNT(*) + FROM furumusic__fed_state_playlist_item i + JOIN furumusic__fed_state_playlist p + ON p.user_id = i.user_id AND p.playlist_id = i.playlist_id + WHERE i.user_id = $1 AND p.local_playlist_id = $2 + AND i.present = true AND i.local_track_id IS NULL", + ) + .bind(user.id) + .bind(r.id) + .fetch_one(pool) + .await + .unwrap_or(0); + cards.push(PlaylistCard { + id: r.id, + title: r.title, + track_count: r.track_count + federated_count, + is_own: r.is_own, + owner_name: Some(r.owner_name), + is_public: r.is_public, + is_saved: r.is_saved, + kind: "user".to_string(), + }); + } Json(cards).into_response() } @@ -3888,6 +4000,25 @@ async fn build_track_items( pool: &sqlx::PgPool, ) -> cot::Result> { let track_ids: Vec = tracks.iter().map(|t| t.id).collect(); + let mut content_ids: HashMap = if track_ids.is_empty() { + HashMap::new() + } else { + sqlx::query( + "SELECT t.id AS track_id, c.content_id + FROM furumusic__track t + JOIN furumusic__media_file m ON m.id = t.audio_file_id + JOIN furumusic__federation_content_id_cache c + ON c.media_file_id = m.id AND c.sha256_hash = m.sha256_hash + WHERE t.id = ANY($1)", + ) + .bind(&track_ids) + .fetch_all(pool) + .await + .map_err(|e| cot::Error::internal(e.to_string()))? + .into_iter() + .map(|row| (row.get("track_id"), row.get("content_id"))) + .collect() + }; let track_artists = if track_ids.is_empty() { Vec::new() @@ -3934,6 +4065,7 @@ async fn build_track_items( let tid = t.id; TrackItem { id: t.id, + content_id: content_ids.remove(&tid), title: t.title, track_number: t.track_number, disc_number: t.disc_number, @@ -4300,6 +4432,56 @@ async fn stream_handler( Ok(response) } +async fn federation_cache_stream_handler( + auth_ctx: auth::AuthContext, + session: Session, + db: Database, + request: cot::request::Request, + Path(path): Path, +) -> cot::Result> { + let Some(_user) = auth::get_request_user(&auth_ctx, &session, &db).await else { + return Ok(json_error(StatusCode::UNAUTHORIZED, "not authenticated")); + }; + let Some((full_path, mime)) = crate::federation::handle().prepared_cache_file(&path.id) else { + return Ok(json_error( + StatusCode::NOT_FOUND, + "federated cache entry not found", + )); + }; + let file_size = tokio::fs::metadata(&full_path) + .await + .map_err(|err| cot::Error::internal(err.to_string()))? + .len(); + if let Some(range) = request + .headers() + .get(RANGE) + .and_then(|value| value.to_str().ok()) + .and_then(|value| parse_range(value, file_size)) + { + let (start, end) = range; + let chunk_size = end - start + 1; + let data = read_file_range(&full_path, start, chunk_size).await?; + return Ok(cot::http::Response::builder() + .status(StatusCode::PARTIAL_CONTENT) + .header(CONTENT_TYPE, mime) + .header(ACCEPT_RANGES, "bytes") + .header(CONTENT_RANGE, format!("bytes {start}-{end}/{file_size}")) + .header(CONTENT_LENGTH, chunk_size.to_string()) + .body(Body::fixed(data)) + .expect("valid response")); + } + let data = tokio::fs::read(full_path) + .await + .map_err(|err| cot::Error::internal(err.to_string()))?; + Ok(cot::http::Response::builder() + .status(StatusCode::OK) + .header(CONTENT_TYPE, mime) + .header(ACCEPT_RANGES, "bytes") + .header(CONTENT_LENGTH, file_size.to_string()) + .body(Body::fixed(data)) + .expect("valid response")) +} + async fn local_upload_handler( auth_ctx: auth::AuthContext, session: Session, @@ -4710,6 +4892,21 @@ async fn devices_select_handler( { return Ok(json_error(StatusCode::BAD_REQUEST, &format!("{err}"))); } + } else if let Some(previous_fed_device_id) = previous_active_id + .as_deref() + .and_then(fed_device_id_from_virtual) + && let Some(playback_state) = response.playback_state.clone() + { + let state = match serde_json::to_value(playback_state) { + Ok(state) => state, + Err(err) => return Ok(json_error(StatusCode::BAD_REQUEST, &format!("{err}"))), + }; + if let Err(err) = crate::federation::handle() + .fed_device_web_active_takeover(user.id, previous_fed_device_id, state) + .await + { + return Ok(json_error(StatusCode::BAD_REQUEST, &format!("{err}"))); + } } Json(response).into_response() } @@ -5429,6 +5626,7 @@ async fn history_list_handler( let tid = row.id; TrackItem { id: row.id, + content_id: None, title: row.title, track_number: row.track_number, disc_number: row.disc_number, @@ -5798,6 +5996,7 @@ async fn search_handler( let tid = t.id; TrackItem { id: t.id, + content_id: None, title: t.title, track_number: t.track_number, disc_number: t.disc_number, @@ -5835,6 +6034,421 @@ async fn search_handler( .into_response() } +// --------------------------------------------------------------------------- +// GET /api/player/federation/search/events?q=... +// --------------------------------------------------------------------------- + +async fn federation_search_events_handler( + auth_ctx: auth::AuthContext, + session: Session, + db: Database, + query: cot::request::extractors::UrlQuery, +) -> cot::Result { + let Some(_user) = auth::get_request_user(&auth_ctx, &session, &db).await else { + return Ok(json_error(StatusCode::UNAUTHORIZED, "not authenticated")); + }; + let q = query.0.q.trim().to_owned(); + if q.is_empty() || q.chars().count() > 200 { + return Ok(json_error(StatusCode::BAD_REQUEST, "invalid search query")); + } + + let search_id = uuid::Uuid::new_v4().to_string(); + let body_search_id = search_id.clone(); + let body = Body::streaming(async_stream::try_stream! { + let started = serde_json::json!({ + "search_id": body_search_id, + "sequence": 0, + "kind": "search.started", + "entity_key": null, + "entity": { "query": q }, + }); + yield Bytes::from(format!("event: search.started\ndata: {started}\n\n")); + + match crate::federation::handle().search_events(&body_search_id, &q).await { + Ok(events) => { + let mut last_sequence = 0; + for event in events { + last_sequence = event.sequence; + let data = serde_json::to_string(&event) + .map_err(|err| cot::Error::internal(err.to_string()))?; + yield Bytes::from(format!("event: {}\ndata: {data}\n\n", event.kind)); + } + let completed = serde_json::json!({ + "search_id": body_search_id, + "sequence": last_sequence + 1, + "kind": "search.completed", + "entity_key": null, + "entity": { "partial": false }, + }); + yield Bytes::from(format!("event: search.completed\ndata: {completed}\n\n")); + } + Err(err) => { + let failed = serde_json::json!({ + "search_id": body_search_id, + "sequence": 1, + "kind": "search.failed", + "entity_key": null, + "entity": { "message": err.to_string() }, + }); + yield Bytes::from(format!("event: search.failed\ndata: {failed}\n\n")); + } + } + }); + cot::http::Response::builder() + .status(StatusCode::OK) + .header(CONTENT_TYPE, "text/event-stream") + .header("cache-control", "no-cache, no-transform") + .header("x-accel-buffering", "no") + .body(body) + .map_err(|err| cot::Error::internal(err.to_string())) +} + +async fn federation_artist_events_handler( + auth_ctx: auth::AuthContext, + session: Session, + db: Database, + UrlQuery(query): UrlQuery, +) -> cot::Result { + let Some(_user) = auth::get_request_user(&auth_ctx, &session, &db).await else { + return Ok(json_error(StatusCode::UNAUTHORIZED, "not authenticated")); + }; + let name = query.name.trim().to_owned(); + if name.is_empty() || name.chars().count() > 200 { + return Ok(json_error(StatusCode::BAD_REQUEST, "invalid artist name")); + } + let mut catalogs = crate::federation::handle().stream_artist_catalogs(name.clone()); + let body = Body::streaming(async_stream::try_stream! { + let started = serde_json::json!({ "kind": "artist.started", "entity": { "name": name } }); + yield Bytes::from(format!("event: artist.started\ndata: {started}\n\n")); + let mut sequence = 0u64; + while let Some(result) = catalogs.recv().await { + sequence += 1; + match result { + Ok((owner, artist)) => { + let event = serde_json::json!({ + "kind": "federation.artist.catalog", + "sequence": sequence, + "peer": owner, + "entity": artist, + }); + yield Bytes::from(format!("event: federation.artist.catalog\ndata: {event}\n\n")); + } + Err(err) => { + tracing::debug!(artist = %name, "federated artist catalog skipped: {err:#}"); + } + } + } + let completed = serde_json::json!({ + "kind": "artist.completed", + "sequence": sequence + 1, + "entity": { "partial": false }, + }); + yield Bytes::from(format!("event: artist.completed\ndata: {completed}\n\n")); + }); + cot::http::Response::builder() + .status(StatusCode::OK) + .header(CONTENT_TYPE, "text/event-stream") + .header("cache-control", "no-cache, no-transform") + .header("x-accel-buffering", "no") + .body(body) + .map_err(|err| cot::Error::internal(err.to_string())) +} + +async fn content_like_handler( + auth_ctx: auth::AuthContext, + session: Session, + db: Database, + pool: &sqlx::PgPool, + Json(body): Json, +) -> cot::Result { + let Some(user) = auth::get_request_user(&auth_ctx, &session, &db).await else { + return Ok(json_error(StatusCode::UNAUTHORIZED, "not authenticated")); + }; + let liked = body.liked.unwrap_or(true); + crate::federation::devices::record_content_like( + pool, + user.id, + &body.content_id, + liked, + body.federation, + ) + .await + .map_err(|err| cot::Error::internal(err.to_string()))?; + Json(serde_json::json!({ "ok": true, "liked": liked })).into_response() +} + +async fn content_likes_handler( + auth_ctx: auth::AuthContext, + session: Session, + db: Database, + pool: &sqlx::PgPool, +) -> cot::Result { + let Some(user) = auth::get_request_user(&auth_ctx, &session, &db).await else { + return Ok(json_error(StatusCode::UNAUTHORIZED, "not authenticated")); + }; + let ids = sqlx::query_scalar::<_, String>( + "SELECT content_id FROM furumusic__fed_state_like + WHERE user_id = $1 AND liked = true ORDER BY hlc_ms DESC", + ) + .bind(user.id) + .fetch_all(pool) + .await + .map_err(|err| cot::Error::internal(err.to_string()))?; + Json(ids).into_response() +} + +async fn content_playlist_add_handler( + auth_ctx: auth::AuthContext, + session: Session, + db: Database, + pool: &sqlx::PgPool, + Json(body): Json, +) -> cot::Result { + let Some(user) = auth::get_request_user(&auth_ctx, &session, &db).await else { + return Ok(json_error(StatusCode::UNAUTHORIZED, "not authenticated")); + }; + let Some(playlist_id) = body.playlist_id else { + return Ok(json_error(StatusCode::BAD_REQUEST, "missing playlist id")); + }; + let owns: bool = sqlx::query_scalar( + "SELECT EXISTS ( + SELECT 1 FROM furumusic__playlist WHERE id = $1 AND owner_id = $2 + )", + ) + .bind(playlist_id) + .bind(user.id) + .fetch_one(pool) + .await + .map_err(|err| cot::Error::internal(err.to_string()))?; + if !owns { + return Ok(json_error(StatusCode::FORBIDDEN, "not your playlist")); + } + let position = match body.position { + Some(position) => position, + None => sqlx::query_scalar::<_, i64>( + "SELECT COALESCE(MAX(position), -1) + 1 + FROM furumusic__fed_state_playlist_item i + JOIN furumusic__fed_state_playlist p + ON p.user_id = i.user_id AND p.playlist_id = i.playlist_id + WHERE p.user_id = $1 AND p.local_playlist_id = $2 AND i.present = true", + ) + .bind(user.id) + .bind(playlist_id) + .fetch_one(pool) + .await + .unwrap_or(0), + }; + crate::federation::devices::record_content_playlist_add( + pool, + user.id, + playlist_id, + &body.content_id, + position, + body.federation, + ) + .await + .map_err(|err| cot::Error::internal(err.to_string()))?; + Json(serde_json::json!({ "ok": true })).into_response() +} + +async fn federation_playlist_tracks_handler( + auth_ctx: auth::AuthContext, + session: Session, + db: Database, + pool: &sqlx::PgPool, + Path(path): Path, +) -> cot::Result { + let Some(user) = auth::get_request_user(&auth_ctx, &session, &db).await else { + return Ok(json_error(StatusCode::UNAUTHORIZED, "not authenticated")); + }; + let rows = if path.id == -1 { + sqlx::query( + "SELECT content_id, 0::bigint AS position, fed_json + FROM furumusic__fed_state_like + WHERE user_id = $1 AND liked = true + AND local_track_id IS NULL AND fed_json IS NOT NULL + ORDER BY hlc_ms DESC", + ) + .bind(user.id) + .fetch_all(pool) + .await + } else { + let accessible: bool = sqlx::query_scalar( + "SELECT EXISTS ( + SELECT 1 FROM furumusic__playlist p + WHERE p.id = $1 + AND (p.owner_id = $2 OR p.is_public = true OR EXISTS ( + SELECT 1 FROM furumusic__saved_playlist sp + WHERE sp.user_id = $2 AND sp.playlist_id = p.id + )) + )", + ) + .bind(path.id) + .bind(user.id) + .fetch_one(pool) + .await + .map_err(|err| cot::Error::internal(err.to_string()))?; + if !accessible { + return Ok(json_error(StatusCode::NOT_FOUND, "playlist not found")); + } + sqlx::query( + "SELECT i.content_id, i.position, i.fed_json + FROM furumusic__fed_state_playlist_item i + JOIN furumusic__fed_state_playlist p + ON p.user_id = i.user_id AND p.playlist_id = i.playlist_id + WHERE i.user_id = $1 AND p.local_playlist_id = $2 + AND i.present = true AND i.local_track_id IS NULL + AND i.fed_json IS NOT NULL + ORDER BY i.position", + ) + .bind(user.id) + .bind(path.id) + .fetch_all(pool) + .await + } + .map_err(|err| cot::Error::internal(err.to_string()))?; + let tracks = rows + .into_iter() + .map(|row| { + serde_json::json!({ + "content_id": row.get::("content_id"), + "position": row.get::("position"), + "federation": row.get::("fed_json"), + }) + }) + .collect::>(); + Json(tracks).into_response() +} + +async fn prepare_federated_track_handler( + auth_ctx: auth::AuthContext, + session: Session, + db: Database, + Json(body): Json, +) -> cot::Result { + let Some(_user) = auth::get_request_user(&auth_ctx, &session, &db).await else { + return Ok(json_error(StatusCode::UNAUTHORIZED, "not authenticated")); + }; + let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel::(); + let progress_sender = sender.clone(); + tokio::spawn(async move { + let result = crate::federation::handle() + .prepare_content_with_progress( + &body.content_id, + &body.owner, + &body.item_id, + move |progress| { + let _ = progress_sender.send(serde_json::json!({ + "kind": "progress", + "phase": progress.phase, + "received": progress.received, + "total": progress.total, + })); + }, + ) + .await; + let event = match result { + Ok(prepared) => serde_json::json!({ + "kind": "completed", + "ok": true, + "local_track_id": prepared.local_track_id, + "stream_url": prepared.stream_url, + }), + Err(err) => serde_json::json!({ + "kind": "failed", + "error": err.to_string(), + }), + }; + let _ = sender.send(event); + }); + let response_body = Body::streaming(async_stream::try_stream! { + while let Some(event) = receiver.recv().await { + let line = serde_json::to_string(&event) + .map_err(|err| cot::Error::internal(err.to_string()))?; + yield Bytes::from(format!("{line}\n")); + } + }); + cot::http::Response::builder() + .status(StatusCode::OK) + .header(CONTENT_TYPE, "application/x-ndjson") + .header("cache-control", "no-cache, no-transform") + .header("x-accel-buffering", "no") + .body(response_body) + .map_err(|err| cot::Error::internal(err.to_string())) +} + +async fn federation_track_artwork_handler( + auth_ctx: auth::AuthContext, + session: Session, + db: Database, + UrlQuery(query): UrlQuery, +) -> cot::Result { + let Some(_user) = auth::get_request_user(&auth_ctx, &session, &db).await else { + return Ok(json_error(StatusCode::UNAUTHORIZED, "not authenticated")); + }; + match crate::federation::handle() + .track_artwork(&query.owner, &query.item_id) + .await + { + Ok(Some((bytes, mime))) => Ok(cot::http::Response::builder() + .status(StatusCode::OK) + .header(CONTENT_TYPE, mime) + .header("cache-control", "private, max-age=86400") + .body(Body::fixed(bytes)) + .expect("valid response")), + Ok(None) => Ok(json_error(StatusCode::NOT_FOUND, "artwork not available")), + Err(err) => Ok(json_error(StatusCode::BAD_GATEWAY, &err.to_string())), + } +} + +async fn federation_catalog_artwork_handler( + auth_ctx: auth::AuthContext, + session: Session, + db: Database, + UrlQuery(query): UrlQuery, +) -> cot::Result { + let Some(_user) = auth::get_request_user(&auth_ctx, &session, &db).await else { + return Ok(json_error(StatusCode::UNAUTHORIZED, "not authenticated")); + }; + match crate::federation::handle() + .catalog_artwork(&query.owner, &query.artist, query.release.as_deref()) + .await + { + Ok(Some((bytes, mime))) => Ok(cot::http::Response::builder() + .status(StatusCode::OK) + .header(CONTENT_TYPE, mime) + .header("cache-control", "private, max-age=86400") + .body(Body::fixed(bytes)) + .expect("valid response")), + Ok(None) => Ok(json_error(StatusCode::NOT_FOUND, "artwork not available")), + Err(err) => Ok(json_error(StatusCode::BAD_GATEWAY, &err.to_string())), + } +} + +async fn federation_artwork_discovery_handler( + auth_ctx: auth::AuthContext, + session: Session, + db: Database, + UrlQuery(query): UrlQuery, +) -> cot::Result { + let Some(_user) = auth::get_request_user(&auth_ctx, &session, &db).await else { + return Ok(json_error(StatusCode::UNAUTHORIZED, "not authenticated")); + }; + match crate::federation::handle() + .discover_catalog_artwork(&query.artist, query.release.as_deref()) + .await + { + Ok(Some((bytes, mime))) => Ok(cot::http::Response::builder() + .status(StatusCode::OK) + .header(CONTENT_TYPE, mime) + .header("cache-control", "private, max-age=300") + .body(Body::fixed(bytes)) + .expect("valid response")), + Ok(None) => Ok(json_error(StatusCode::NOT_FOUND, "artwork not available")), + Err(err) => Ok(json_error(StatusCode::BAD_GATEWAY, &err.to_string())), + } +} + // --------------------------------------------------------------------------- // POST /api/player/playlists — create playlist // --------------------------------------------------------------------------- @@ -8927,6 +9541,195 @@ impl App for PlayerApp { }), "player_search", ), + Route::with_handler_and_name( + "/federation/search/events", + get( + move |auth_ctx: auth::AuthContext, + session: Session, + db: Database, + query: cot::request::extractors::UrlQuery| async move { + federation_search_events_handler(auth_ctx, session, db, query).await + }, + ), + "player_federation_search_events", + ), + Route::with_handler_and_name( + "/tracks/content/like", + get({ + let pool = Arc::clone(&pool); + let pool_config = Arc::clone(&pool_config); + move |auth_ctx: auth::AuthContext, session: Session, db: Database| { + let pool = Arc::clone(&pool); + let pool_config = Arc::clone(&pool_config); + async move { + let pg_pool = pool + .get_or_init(|| async { + sqlx::postgres::PgPoolOptions::new() + .max_connections(5) + .connect(&pool_config.database_url) + .await + .expect("player pool") + }) + .await; + content_likes_handler(auth_ctx, session, db, pg_pool).await + } + } + }) + .post({ + let pool = Arc::clone(&pool); + let pool_config = Arc::clone(&pool_config); + move |auth_ctx: auth::AuthContext, + session: Session, + db: Database, + json: Json| { + let pool = Arc::clone(&pool); + let pool_config = Arc::clone(&pool_config); + async move { + let pg_pool = pool + .get_or_init(|| async { + sqlx::postgres::PgPoolOptions::new() + .max_connections(5) + .connect(&pool_config.database_url) + .await + .expect("player pool") + }) + .await; + content_like_handler(auth_ctx, session, db, pg_pool, json).await + } + } + }), + "player_content_like", + ), + Route::with_handler_and_name( + "/tracks/content/playlist", + post({ + let pool = Arc::clone(&pool); + let pool_config = Arc::clone(&pool_config); + move |auth_ctx: auth::AuthContext, + session: Session, + db: Database, + json: Json| { + let pool = Arc::clone(&pool); + let pool_config = Arc::clone(&pool_config); + async move { + let pg_pool = pool + .get_or_init(|| async { + sqlx::postgres::PgPoolOptions::new() + .max_connections(5) + .connect(&pool_config.database_url) + .await + .expect("player pool") + }) + .await; + content_playlist_add_handler(auth_ctx, session, db, pg_pool, json).await + } + } + }), + "player_content_playlist", + ), + Route::with_handler_and_name( + "/federation/tracks/prepare", + post( + move |auth_ctx: auth::AuthContext, + session: Session, + db: Database, + json: Json| async move { + prepare_federated_track_handler(auth_ctx, session, db, json).await + }, + ), + "player_federation_track_prepare", + ), + Route::with_handler_and_name( + "/playlists/{id}/federation-tracks", + get({ + let pool = Arc::clone(&pool); + let pool_config = Arc::clone(&pool_config); + move |auth_ctx: auth::AuthContext, + session: Session, + db: Database, + path: Path| { + let pool = Arc::clone(&pool); + let pool_config = Arc::clone(&pool_config); + async move { + let pg_pool = pool + .get_or_init(|| async { + sqlx::postgres::PgPoolOptions::new() + .max_connections(5) + .connect(&pool_config.database_url) + .await + .expect("player pool") + }) + .await; + federation_playlist_tracks_handler( + auth_ctx, session, db, pg_pool, path, + ) + .await + } + } + }), + "player_federation_playlist_tracks", + ), + Route::with_handler_and_name( + "/federation/artists/events", + get( + move |auth_ctx: auth::AuthContext, + session: Session, + db: Database, + query: UrlQuery| async move { + federation_artist_events_handler(auth_ctx, session, db, query).await + }, + ), + "player_federation_artist_events", + ), + Route::with_handler_and_name( + "/federation/tracks/artwork", + get( + move |auth_ctx: auth::AuthContext, + session: Session, + db: Database, + query: UrlQuery| async move { + federation_track_artwork_handler(auth_ctx, session, db, query).await + }, + ), + "player_federation_track_artwork", + ), + Route::with_handler_and_name( + "/federation/catalog/artwork", + get( + move |auth_ctx: auth::AuthContext, + session: Session, + db: Database, + query: UrlQuery| async move { + federation_catalog_artwork_handler(auth_ctx, session, db, query).await + }, + ), + "player_federation_catalog_artwork", + ), + Route::with_handler_and_name( + "/federation/catalog/artwork/discover", + get( + move |auth_ctx: auth::AuthContext, + session: Session, + db: Database, + query: UrlQuery| async move { + federation_artwork_discovery_handler(auth_ctx, session, db, query).await + }, + ), + "player_federation_artwork_discovery", + ), + Route::with_handler_and_name( + "/federation/cache/{id}", + get( + move |auth_ctx: auth::AuthContext, + session: Session, + db: Database, + request: cot::request::Request, + path: Path| async move { + federation_cache_stream_handler(auth_ctx, session, db, request, path).await + }, + ), + "player_federation_cache_stream", + ), // -- Tracks by IDs -- Route::with_handler_and_name( "/tracks-by-ids", diff --git a/templates/admin/v2.html b/templates/admin/v2.html index c08f53a..9aecbd9 100644 --- a/templates/admin/v2.html +++ b/templates/admin/v2.html @@ -2293,6 +2293,17 @@ tbody tr:hover {
Every peer using the same id finds the others automatically.
+
+ +
+ + +
+
Server-wide policy. Imported tracks become available to every user and are published by this peer. Federation metadata is trusted and bypasses the AI agent.
+
@@ -2924,7 +2935,8 @@ function adminV2() { agent_context_limit: '', agent_concurrency: '', federation_enabled: false, - federation_network_id: '' + federation_network_id: '', + federation_save_on_listen: false }, settingsProbe: { status: 'idle', ok: false }, settingsProbeLoading: false, diff --git a/templates/player/modals.html b/templates/player/modals.html index db8f73b..a4efdab 100644 --- a/templates/player/modals.html +++ b/templates/player/modals.html @@ -636,6 +636,144 @@
+ + +
+
+

+ Federation + LIVE +

+ + +
+ +
+
+ +
+ +
@@ -580,7 +706,9 @@
@@ -1041,14 +1365,19 @@ - -
@@ -1118,6 +1449,31 @@ + + + +
@@ -1125,7 +1481,9 @@
@@ -1134,7 +1492,7 @@
-