diff --git a/Cargo.lock b/Cargo.lock index 6414ad0..d0a8b5f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1793,9 +1793,9 @@ dependencies = [ [[package]] name = "federation-net" -version = "0.2.0" +version = "0.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "15a8707baeccb46b5935138f9cb3df3c988c0730b807a2634d901f26b39250d6" +checksum = "c3e690b370c505d153bef214b21a8f2aa55d667367ac1e16bde8bc0de88963c2" dependencies = [ "blake3", "data-encoding", @@ -1938,7 +1938,7 @@ checksum = "e6d5a32815ae3f33302d95fdcb2ce17862f8c65363dcfd29360480ba1001fc9c" [[package]] name = "furumusic" -version = "0.9.8" +version = "0.10.1" dependencies = [ "anyhow", "async-stream", @@ -3775,9 +3775,9 @@ dependencies = [ [[package]] name = "music-dht" -version = "0.3.1" +version = "0.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "91592b40de9f3c2158a39da105c821e9bf17f461fe142a56a8607fb0faf56a9c" +checksum = "0c5b429b90a8f1b0980b3a35a6fa5445d7a275c737eb04db18db4d7f14c81478" dependencies = [ "async-trait", "blake3", diff --git a/Cargo.toml b/Cargo.toml index a83da2d..6685dbb 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "furumusic" -version = "0.10.0" +version = "0.10.1" edition = "2024" description = "Reusable web-app boilerplate: auth, OIDC/SSO, admin panel, user management, i18n, PostgreSQL" @@ -43,4 +43,4 @@ uuid = "1" librqbit = { version = "8.1.1", features = ["disable-upload"] } # P2P federation: publishes the library into a shared DHT and serves audio / # catalogs to furumi peers (TUI clients) over the frid stack. -music-dht = "0.3.1" +music-dht = "0.4.0" diff --git a/README.md b/README.md index 6ddfd4a..2e5695c 100644 --- a/README.md +++ b/README.md @@ -1,5 +1,12 @@ # furumusic +Furumusic can join the decentralized Furumi federation while remaining a +complete local web player. Optional similarity search stores versioned audio +embeddings in PostgreSQL and uses signed two-level LSH summaries in a separate +DHT to discover compatible peers without a central recommendation index. The +shared `music-dht` layer owns routing and wire compatibility; model inference +and exact cosine ranking stay local to each instance. + Reusable web-app boilerplate: auth, OIDC/SSO, admin panel, user management, i18n, PostgreSQL. Built with Rust ([cot](https://cot.rs) framework). diff --git a/src/federation/capabilities.rs b/src/federation/capabilities.rs index a2fcb96..26cf5c5 100644 --- a/src/federation/capabilities.rs +++ b/src/federation/capabilities.rs @@ -69,6 +69,12 @@ mod tests { manifest.protocols.get(SIMILARITY_ID), Some(&music_dht::similarity::SIMILARITY_PROTOCOL_VERSION) ); + assert_eq!( + manifest + .protocols + .get(music_dht::capabilities::SIMILARITY_DHT_ID), + Some(&music_dht::similarity_lsh::SIMILARITY_DHT_PROTOCOL_VERSION) + ); manifest.validate().unwrap(); } } diff --git a/src/federation/mod.rs b/src/federation/mod.rs index a1f1a15..6ce5917 100644 --- a/src/federation/mod.rs +++ b/src/federation/mod.rs @@ -29,6 +29,8 @@ use std::time::Duration; use anyhow::{Context, Result}; use music_dht::capabilities::CAPABILITIES_ALPN; +use music_dht::similarity_dht::SimilarityDht; +use music_dht::similarity_lsh::SIMILARITY_DHT_ALPN; use music_dht::{ ByteStream, ByteStreamConnectionStats, ItemKind, ItemSpec, MusicDhtConfig, MusicDhtService, NetworkId, PeerTicket, PublishStats, RendezvousConfig, SyncStats, @@ -49,6 +51,7 @@ const TRANSPORT_SAMPLE_LIMIT: usize = 16; struct Running { service: Arc, + similarity_dht: Arc, network_name: String, tasks: Vec>, } @@ -221,7 +224,8 @@ pub fn record_stream_transport( } pub struct Federation { - /// Transport data directory; server-side DHT state and identity live in PostgreSQL. + /// Transport files and replaceable similarity-routing cache. Durable + /// catalog DHT state and identity live in PostgreSQL. data_dir: PathBuf, database_url: std::sync::Mutex, storage_dir: std::sync::Mutex, @@ -380,6 +384,9 @@ impl Federation { let dht_storage = Arc::new(PostgresFederationStorage::new(pool.clone()).await?); let secret_key = dht_storage.load_or_create_secret_key().await?; self.transport_stats.reset(); + tokio::fs::create_dir_all(&self.data_dir) + .await + .with_context(|| format!("creating {}", self.data_dir.display()))?; let config = MusicDhtConfig::builder() .data_dir(&self.data_dir) @@ -389,7 +396,8 @@ impl Federation { .stream_protocol(AUDIO_ALPN) .stream_protocol(CATALOG_ALPN) .stream_protocol(devices::SYNC_ALPN) - .stream_protocol(SIMILARITY_ALPN) + .schema_independent_stream_protocol(SIMILARITY_ALPN) + .schema_independent_stream_protocol(SIMILARITY_DHT_ALPN) .schema_independent_stream_protocol(CAPABILITIES_ALPN) .build() .map_err(|err| anyhow::anyhow!("invalid federation config: {err}"))?; @@ -404,6 +412,25 @@ impl Federation { "federation started" ); + let similarity_dht = SimilarityDht::open( + Arc::clone(&service), + self.data_dir.join("similarity-routing.sqlite3"), + ) + .await + .map_err(|error| anyhow::anyhow!("failed to start the similarity DHT: {error}"))?; + let similarity_dht_acceptor = service + .stream_acceptor(SIMILARITY_DHT_ALPN) + .map_err(|error| anyhow::anyhow!("failed to take similarity DHT acceptor: {error}"))?; + let similarity_dht_serve_task = + tokio::spawn(Arc::clone(&similarity_dht).serve(similarity_dht_acceptor)); + let similarity_dht_maintenance_task = + tokio::spawn(Arc::clone(&similarity_dht).maintenance()); + let similarity_manager = crate::similarity::handle(); + let similarity_dht_sync_task = tokio::spawn(similarity_route_sync_loop( + Arc::clone(&similarity_dht), + Arc::clone(&similarity_manager), + )); + // Drain DHT events into the log; the channel is bounded. let event_task = tokio::spawn(async move { while let Some(event) = events.recv().await { @@ -467,13 +494,14 @@ impl Federation { .map_err(|err| anyhow::anyhow!("failed to take the similarity acceptor: {err}"))?; let similarity_task = tokio::spawn(similarity::serve_peers( similarity_acceptor, - crate::similarity::handle(), + similarity_manager, service.endpoint_id(), Arc::clone(&self.transport_stats), )); *guard = Some(Running { service, + similarity_dht, network_name, tasks: vec![ event_task, @@ -484,6 +512,9 @@ impl Federation { device_sync_task, capabilities_task, similarity_task, + similarity_dht_serve_task, + similarity_dht_maintenance_task, + similarity_dht_sync_task, ], }); self.set_error(None); @@ -507,6 +538,20 @@ impl Federation { .context("federation is not running") } + async fn similarity_services(&self) -> Result<(Arc, Arc)> { + self.running + .lock() + .await + .as_ref() + .map(|running| { + ( + Arc::clone(&running.service), + Arc::clone(&running.similarity_dht), + ) + }) + .context("federation is not running") + } + async fn spawn_sync_soon(self: &Arc) { if let Ok(service) = self.service().await { let fed = Arc::clone(self); @@ -810,6 +855,7 @@ impl Federation { "endpoint_id": service.endpoint_id().to_string(), "connected_peers": peers, "known_contacts": service.known_peers().len(), + "similarity_routing_peers": running.similarity_dht.known_peers(), "published_items": published, "transport": self.transport_stats.snapshot(), }) @@ -854,8 +900,15 @@ impl Federation { crate::similarity::handle().enabled(), "similarity search is disabled" ); - let service = self.service().await?; - similarity::search(service, query, limit, Arc::clone(&self.transport_stats)).await + let (service, similarity_dht) = self.similarity_services().await?; + similarity::search( + service, + similarity_dht, + query, + limit, + Arc::clone(&self.transport_stats), + ) + .await } pub async fn fed_device_status( @@ -983,6 +1036,65 @@ impl Federation { } } +async fn similarity_route_sync_loop( + routing: Arc, + manager: Arc, +) { + let mut interval = tokio::time::interval(SYNC_INTERVAL); + interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + let mut published_marker: Option<(String, blake3::Hash)> = None; + loop { + interval.tick().await; + if !manager.enabled() { + if published_marker.take().is_some() { + routing.clear_local_signatures(); + tracing::info!("local similarity DHT publication disabled"); + } + continue; + } + let status = manager.status(); + let Some(profile_id) = status.active_profile else { + continue; + }; + if status.phase != crate::similarity::Phase::Ready { + continue; + } + let signatures = match manager.routing_signatures(&profile_id).await { + Ok(signatures) => signatures, + Err(error) => { + tracing::warn!(%error, %profile_id, "similarity routing signatures unavailable"); + continue; + } + }; + let mut hasher = blake3::Hasher::new(); + for signature in &signatures { + hasher.update(signature); + } + let marker = (profile_id.clone(), hasher.finalize()); + if published_marker.as_ref() == Some(&marker) { + continue; + } + match routing + .sync_local_signatures(profile_id.clone(), signatures) + .await + { + Ok(stats) => { + tracing::info!( + profile = %profile_id, + records = stats.records, + keys = stats.keys, + remote_nodes = stats.remote_nodes, + "local similarity DHT index synchronized" + ); + published_marker = Some(marker); + } + Err(error) => { + tracing::warn!(%error, %profile_id, "similarity DHT synchronization failed"); + } + } + } +} + async fn persist_content_id( pool: &PgPool, media_file_id: i64, diff --git a/src/federation/similarity.rs b/src/federation/similarity.rs index 512c0ec..7a68aa1 100644 --- a/src/federation/similarity.rs +++ b/src/federation/similarity.rs @@ -7,7 +7,10 @@ use std::time::Duration; use anyhow::{Context as _, Result}; use futures_util::stream::{self, StreamExt as _}; use music_dht::similarity::{self as wire, SimilarityHit, SimilarityRequest, SimilarityResponse}; -use music_dht::{ByteStream, EndpointId, ItemId, ItemKind, MusicDhtService, StreamAcceptor}; +use music_dht::similarity_dht::SimilarityDht; +use music_dht::{ + ByteStream, EndpointId, ItemId, ItemKind, MusicDhtService, PeerTicket, StreamAcceptor, +}; use crate::similarity::{Manager, QueryVector}; @@ -15,9 +18,11 @@ use super::TransportStats; pub use music_dht::similarity::SIMILARITY_ALPN; -const MAX_QUERY_PEERS: usize = 16; -const QUERY_CONCURRENCY: usize = 6; +const INITIAL_QUERY_PEERS: usize = 16; +const MAX_QUERY_PEERS: usize = 48; +const QUERY_CONCURRENCY: usize = 8; const QUERY_TIMEOUT: Duration = Duration::from_secs(5); +const ROUTING_TIMEOUT: Duration = Duration::from_secs(5); const MAX_PER_ARTIST: usize = 3; const MAX_NEAR_DUPLICATE_SIGNATURE_DISTANCE: u32 = 8; @@ -142,20 +147,49 @@ async fn serve_one( pub async fn search( service: Arc, + routing: Arc, query: QueryVector, limit: usize, transport: Arc, ) -> Result> { let own = service.endpoint_id(); - let mut peers = Vec::new(); + let routed = match tokio::time::timeout( + ROUTING_TIMEOUT, + routing.find_peers(&query.profile_id, &query.vector, MAX_QUERY_PEERS), + ) + .await + { + Ok(Ok(peers)) => peers, + Err(_) => { + tracing::debug!("similarity DHT lookup timed out; using known peers"); + Vec::new() + } + Ok(Err(error)) => { + tracing::debug!(%error, "similarity DHT lookup unavailable; using known peers"); + Vec::new() + } + }; let mut seen = HashSet::new(); + let mut peers: Vec = routed + .into_iter() + .filter_map(|ticket| { + let owner = ticket.endpoint_id(); + (owner != own && seen.insert(owner)).then_some(QueryPeer { + owner, + ticket: Some(ticket), + }) + }) + .collect(); for peer in service .connected_peers() .into_iter() .chain(service.known_peers().into_iter().map(|peer| peer.peer_id)) { if peer != own && seen.insert(peer) { - peers.push(peer); + peers.push(QueryPeer { + owner: peer, + ticket: None, + }); } if peers.len() >= MAX_QUERY_PEERS { break; @@ -168,30 +202,40 @@ pub async fn search( limit.clamp(1, wire::MAX_SIMILARITY_RESULTS), )?); - let responses = stream::iter(peers.into_iter().map(|peer| { - let service = Arc::clone(&service); - let request = Arc::clone(&request); - let transport = Arc::clone(&transport); - async move { - tokio::time::timeout( - QUERY_TIMEOUT, - query_peer(service, peer, &request, transport), - ) - .await - .map_err(|_| anyhow::anyhow!("similarity peer timed out"))? - } - })) - .buffer_unordered(QUERY_CONCURRENCY) - .collect::>() - .await; - let mut hits = Vec::new(); + let initial = peers.len().min(INITIAL_QUERY_PEERS); + let responses = query_peers( + Arc::clone(&service), + &peers[..initial], + Arc::clone(&request), + Arc::clone(&transport), + ) + .await; + let mut successful = 0usize; for response in responses { match response { - Ok(peer_hits) => hits.extend(peer_hits), + Ok(peer_hits) => { + successful += 1; + hits.extend(peer_hits); + } Err(error) => tracing::debug!(%error, "similarity peer query skipped"), } } + if initial < peers.len() && (hits.len() < limit || successful < initial.min(4)) { + for response in query_peers( + Arc::clone(&service), + &peers[initial..], + Arc::clone(&request), + Arc::clone(&transport), + ) + .await + { + match response { + Ok(peer_hits) => hits.extend(peer_hits), + Err(error) => tracing::debug!(%error, "fallback similarity peer query skipped"), + } + } + } hits.sort_by(|left, right| right.1.total_cmp(&left.1)); let mut dedup = HashSet::new(); let mut signatures = vec![query_signature]; @@ -241,22 +285,54 @@ pub async fn search( Ok(tracks) } +type PeerHits = Vec<( + RemoteSimilarityTrack, + f32, + Option<[u8; wire::SIMILARITY_SIGNATURE_BYTES]>, +)>; + +#[derive(Clone)] +struct QueryPeer { + owner: EndpointId, + ticket: Option, +} + +async fn query_peers( + service: Arc, + peers: &[QueryPeer], + request: Arc, + transport: Arc, +) -> Vec> { + stream::iter(peers.iter().cloned().map(|peer| { + let service = Arc::clone(&service); + let request = Arc::clone(&request); + let transport = Arc::clone(&transport); + async move { + tokio::time::timeout( + QUERY_TIMEOUT, + query_peer(service, peer, &request, transport), + ) + .await + .map_err(|_| anyhow::anyhow!("similarity peer timed out"))? + } + })) + .buffer_unordered(QUERY_CONCURRENCY) + .collect() + .await +} + async fn query_peer( service: Arc, - owner: EndpointId, + peer: QueryPeer, request: &SimilarityRequest, transport: Arc, -) -> Result< - Vec<( - RemoteSimilarityTrack, - f32, - Option<[u8; wire::SIMILARITY_SIGNATURE_BYTES]>, - )>, -> { - let mut stream = service - .open_stream(owner, SIMILARITY_ALPN) - .await - .map_err(|error| anyhow::anyhow!("cannot reach similarity peer: {error}"))?; +) -> Result { + let owner = peer.owner; + let mut stream = match peer.ticket { + Some(ticket) => service.open_stream_to(&ticket, SIMILARITY_ALPN).await, + None => service.open_stream(owner, SIMILARITY_ALPN).await, + } + .map_err(|error| anyhow::anyhow!("cannot reach similarity peer: {error}"))?; super::record_stream_transport(&transport, "similarity", "outbound", "open", &stream); let response = wire::exchange(&mut stream, request).await?; super::record_stream_transport(&transport, "similarity", "outbound", "done", &stream); diff --git a/src/music/mod.rs b/src/music/mod.rs index f2123c2..c0962df 100644 --- a/src/music/mod.rs +++ b/src/music/mod.rs @@ -2541,6 +2541,34 @@ pub mod db_migrations { &[Operation::custom(create_similarity_embeddings).build()]; } + #[cot::db::migrations::migration_op] + async fn add_similarity_routing_signature( + ctx: migrations::MigrationContext<'_>, + ) -> cot::db::Result<()> { + ctx.db + .raw( + "ALTER TABLE furumusic__track_embedding + ADD COLUMN IF NOT EXISTS routing_signature BYTEA", + ) + .await?; + Ok(()) + } + + #[derive(Debug, Copy, Clone)] + pub struct M0044AddSimilarityRoutingSignature; + + impl migrations::Migration for M0044AddSimilarityRoutingSignature { + const APP_NAME: &'static str = "furumusic"; + const MIGRATION_NAME: &'static str = "m_0044_add_similarity_routing_signature"; + const DEPENDENCIES: &'static [migrations::MigrationDependency] = + &[migrations::MigrationDependency::migration( + "furumusic", + "m_0043_create_similarity_embeddings", + )]; + const OPERATIONS: &'static [Operation] = + &[Operation::custom(add_similarity_routing_signature).build()]; + } + pub const MIGRATIONS: &[&SyncDynMigration] = &[ &M0006CreateMediaFile, &M0007CreateArtist, @@ -2575,5 +2603,6 @@ pub mod db_migrations { &M0041CreateSyncedListenHistory, &M0042RepairLegacyListenQualification, &M0043CreateSimilarityEmbeddings, + &M0044AddSimilarityRoutingSignature, ]; } diff --git a/src/similarity.rs b/src/similarity.rs index 6419129..4bf166f 100644 --- a/src/similarity.rs +++ b/src/similarity.rs @@ -283,6 +283,90 @@ impl Manager { lock(&self.status).clone() } + /// Loads compact routing signatures for every current visible embedding. + /// Embeddings created before DHT routing existed are upgraded in place; + /// the CPU-heavy projection runs outside the async runtime. + pub async fn routing_signatures(&self, profile_id: &str) -> Result> { + let pool = self.pool().await?; + let missing = sqlx::query( + "SELECT e.track_id, e.dimensions, e.vector + FROM furumusic__track_embedding e + JOIN furumusic__track t ON t.id = e.track_id + JOIN furumusic__release r ON r.id = t.release_id + JOIN furumusic__media_file m ON m.id = t.audio_file_id + WHERE e.profile_id = $1 AND e.source_sha256 = m.sha256_hash + AND t.is_hidden = FALSE AND r.is_hidden = FALSE + AND (e.routing_signature IS NULL + OR octet_length(e.routing_signature) != 32) + ORDER BY e.track_id", + ) + .bind(profile_id) + .fetch_all(&pool) + .await? + .into_iter() + .map(|row| { + ( + row.get::(0), + row.get::(1), + row.get::, _>(2), + ) + }) + .collect::>(); + + let computed = tokio::task::spawn_blocking(move || { + missing + .into_iter() + .map(|(track_id, dimensions, bytes)| { + let vector = embedding_from_bytes(dimensions, &bytes)?; + let signature = music_dht::similarity_lsh::routing_signature(&vector)?; + Ok::<_, anyhow::Error>((track_id, signature)) + }) + .collect::>>() + }) + .await + .context("similarity routing backfill task failed")??; + + if !computed.is_empty() { + let mut transaction = pool.begin().await?; + for (track_id, signature) in computed { + sqlx::query( + "UPDATE furumusic__track_embedding + SET routing_signature = $3 + WHERE track_id = $1 AND profile_id = $2 + AND (routing_signature IS NULL + OR octet_length(routing_signature) != 32)", + ) + .bind(track_id) + .bind(profile_id) + .bind(signature.as_slice()) + .execute(&mut *transaction) + .await?; + } + transaction.commit().await?; + } + + let stored = sqlx::query_scalar::<_, Vec>( + "SELECT e.routing_signature + FROM furumusic__track_embedding e + JOIN furumusic__track t ON t.id = e.track_id + JOIN furumusic__release r ON r.id = t.release_id + JOIN furumusic__media_file m ON m.id = t.audio_file_id + WHERE e.profile_id = $1 AND e.source_sha256 = m.sha256_hash + AND t.is_hidden = FALSE AND r.is_hidden = FALSE + ORDER BY e.track_id", + ) + .bind(profile_id) + .fetch_all(&pool) + .await?; + stored + .into_iter() + .map(|signature| { + <[u8; 32]>::try_from(signature) + .map_err(|_| anyhow::anyhow!("invalid similarity routing signature length")) + }) + .collect() + } + pub fn start(self: &Arc) { let generation = self.generation.fetch_add(1, Ordering::AcqRel) + 1; let manager = Arc::clone(self); @@ -846,14 +930,16 @@ async fn store_embedding( vector.iter().all(|value| value.is_finite()), "embedding contains a non-finite value" ); + let routing_signature = music_dht::similarity_lsh::routing_signature(vector)?; sqlx::query( "INSERT INTO furumusic__track_embedding - (track_id, profile_id, dimensions, vector, source_sha256, - source_content_id, computed_at) - VALUES ($1, $2, $3, $4, $5, $6, $7) + (track_id, profile_id, dimensions, vector, routing_signature, + source_sha256, source_content_id, computed_at) + VALUES ($1, $2, $3, $4, $5, $6, $7, $8) ON CONFLICT (track_id, profile_id) DO UPDATE SET dimensions = EXCLUDED.dimensions, vector = EXCLUDED.vector, + routing_signature = EXCLUDED.routing_signature, source_sha256 = EXCLUDED.source_sha256, source_content_id = EXCLUDED.source_content_id, computed_at = EXCLUDED.computed_at", @@ -862,6 +948,7 @@ async fn store_embedding( .bind(profile_id) .bind(vector.len() as i32) .bind(embedding_to_bytes(vector)) + .bind(routing_signature.as_slice()) .bind(&track.source_sha256) .bind(&track.source_content_id) .bind(now_iso()) diff --git a/templates/admin/v2.html b/templates/admin/v2.html index c606911..62e3851 100644 --- a/templates/admin/v2.html +++ b/templates/admin/v2.html @@ -2210,7 +2210,7 @@ tbody tr:hover { -
Downloads the selected model and processes every visible local track. When federation is also enabled, this instance sends anonymized query embeddings to peers and answers their searches.
+
Downloads the selected model and processes every visible local track. With federation enabled, signed anonymous LSH summaries discover likely peers; full query embeddings are sent only to those peers, and this instance answers their searches.