Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
6f337ee626 |
Generated
+5
-5
@@ -1793,9 +1793,9 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "federation-net"
|
name = "federation-net"
|
||||||
version = "0.2.0"
|
version = "0.3.0"
|
||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "15a8707baeccb46b5935138f9cb3df3c988c0730b807a2634d901f26b39250d6"
|
checksum = "c3e690b370c505d153bef214b21a8f2aa55d667367ac1e16bde8bc0de88963c2"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"blake3",
|
"blake3",
|
||||||
"data-encoding",
|
"data-encoding",
|
||||||
@@ -1938,7 +1938,7 @@ checksum = "e6d5a32815ae3f33302d95fdcb2ce17862f8c65363dcfd29360480ba1001fc9c"
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "furumusic"
|
name = "furumusic"
|
||||||
version = "0.9.8"
|
version = "0.10.1"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"async-stream",
|
"async-stream",
|
||||||
@@ -3775,9 +3775,9 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "music-dht"
|
name = "music-dht"
|
||||||
version = "0.3.1"
|
version = "0.4.0"
|
||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "91592b40de9f3c2158a39da105c821e9bf17f461fe142a56a8607fb0faf56a9c"
|
checksum = "0c5b429b90a8f1b0980b3a35a6fa5445d7a275c737eb04db18db4d7f14c81478"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"async-trait",
|
"async-trait",
|
||||||
"blake3",
|
"blake3",
|
||||||
|
|||||||
+2
-2
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "furumusic"
|
name = "furumusic"
|
||||||
version = "0.10.0"
|
version = "0.10.1"
|
||||||
edition = "2024"
|
edition = "2024"
|
||||||
description = "Reusable web-app boilerplate: auth, OIDC/SSO, admin panel, user management, i18n, PostgreSQL"
|
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"] }
|
librqbit = { version = "8.1.1", features = ["disable-upload"] }
|
||||||
# P2P federation: publishes the library into a shared DHT and serves audio /
|
# P2P federation: publishes the library into a shared DHT and serves audio /
|
||||||
# catalogs to furumi peers (TUI clients) over the frid stack.
|
# catalogs to furumi peers (TUI clients) over the frid stack.
|
||||||
music-dht = "0.3.1"
|
music-dht = "0.4.0"
|
||||||
|
|||||||
@@ -1,5 +1,12 @@
|
|||||||
# furumusic
|
# 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.
|
Reusable web-app boilerplate: auth, OIDC/SSO, admin panel, user management, i18n, PostgreSQL.
|
||||||
|
|
||||||
Built with Rust ([cot](https://cot.rs) framework).
|
Built with Rust ([cot](https://cot.rs) framework).
|
||||||
|
|||||||
@@ -69,6 +69,12 @@ mod tests {
|
|||||||
manifest.protocols.get(SIMILARITY_ID),
|
manifest.protocols.get(SIMILARITY_ID),
|
||||||
Some(&music_dht::similarity::SIMILARITY_PROTOCOL_VERSION)
|
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();
|
manifest.validate().unwrap();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+117
-5
@@ -29,6 +29,8 @@ use std::time::Duration;
|
|||||||
|
|
||||||
use anyhow::{Context, Result};
|
use anyhow::{Context, Result};
|
||||||
use music_dht::capabilities::CAPABILITIES_ALPN;
|
use music_dht::capabilities::CAPABILITIES_ALPN;
|
||||||
|
use music_dht::similarity_dht::SimilarityDht;
|
||||||
|
use music_dht::similarity_lsh::SIMILARITY_DHT_ALPN;
|
||||||
use music_dht::{
|
use music_dht::{
|
||||||
ByteStream, ByteStreamConnectionStats, ItemKind, ItemSpec, MusicDhtConfig, MusicDhtService,
|
ByteStream, ByteStreamConnectionStats, ItemKind, ItemSpec, MusicDhtConfig, MusicDhtService,
|
||||||
NetworkId, PeerTicket, PublishStats, RendezvousConfig, SyncStats,
|
NetworkId, PeerTicket, PublishStats, RendezvousConfig, SyncStats,
|
||||||
@@ -49,6 +51,7 @@ const TRANSPORT_SAMPLE_LIMIT: usize = 16;
|
|||||||
|
|
||||||
struct Running {
|
struct Running {
|
||||||
service: Arc<MusicDhtService>,
|
service: Arc<MusicDhtService>,
|
||||||
|
similarity_dht: Arc<SimilarityDht>,
|
||||||
network_name: String,
|
network_name: String,
|
||||||
tasks: Vec<tokio::task::JoinHandle<()>>,
|
tasks: Vec<tokio::task::JoinHandle<()>>,
|
||||||
}
|
}
|
||||||
@@ -221,7 +224,8 @@ pub fn record_stream_transport(
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub struct Federation {
|
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,
|
data_dir: PathBuf,
|
||||||
database_url: std::sync::Mutex<String>,
|
database_url: std::sync::Mutex<String>,
|
||||||
storage_dir: std::sync::Mutex<String>,
|
storage_dir: std::sync::Mutex<String>,
|
||||||
@@ -380,6 +384,9 @@ impl Federation {
|
|||||||
let dht_storage = Arc::new(PostgresFederationStorage::new(pool.clone()).await?);
|
let dht_storage = Arc::new(PostgresFederationStorage::new(pool.clone()).await?);
|
||||||
let secret_key = dht_storage.load_or_create_secret_key().await?;
|
let secret_key = dht_storage.load_or_create_secret_key().await?;
|
||||||
self.transport_stats.reset();
|
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()
|
let config = MusicDhtConfig::builder()
|
||||||
.data_dir(&self.data_dir)
|
.data_dir(&self.data_dir)
|
||||||
@@ -389,7 +396,8 @@ impl Federation {
|
|||||||
.stream_protocol(AUDIO_ALPN)
|
.stream_protocol(AUDIO_ALPN)
|
||||||
.stream_protocol(CATALOG_ALPN)
|
.stream_protocol(CATALOG_ALPN)
|
||||||
.stream_protocol(devices::SYNC_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)
|
.schema_independent_stream_protocol(CAPABILITIES_ALPN)
|
||||||
.build()
|
.build()
|
||||||
.map_err(|err| anyhow::anyhow!("invalid federation config: {err}"))?;
|
.map_err(|err| anyhow::anyhow!("invalid federation config: {err}"))?;
|
||||||
@@ -404,6 +412,25 @@ impl Federation {
|
|||||||
"federation started"
|
"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.
|
// Drain DHT events into the log; the channel is bounded.
|
||||||
let event_task = tokio::spawn(async move {
|
let event_task = tokio::spawn(async move {
|
||||||
while let Some(event) = events.recv().await {
|
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}"))?;
|
.map_err(|err| anyhow::anyhow!("failed to take the similarity acceptor: {err}"))?;
|
||||||
let similarity_task = tokio::spawn(similarity::serve_peers(
|
let similarity_task = tokio::spawn(similarity::serve_peers(
|
||||||
similarity_acceptor,
|
similarity_acceptor,
|
||||||
crate::similarity::handle(),
|
similarity_manager,
|
||||||
service.endpoint_id(),
|
service.endpoint_id(),
|
||||||
Arc::clone(&self.transport_stats),
|
Arc::clone(&self.transport_stats),
|
||||||
));
|
));
|
||||||
|
|
||||||
*guard = Some(Running {
|
*guard = Some(Running {
|
||||||
service,
|
service,
|
||||||
|
similarity_dht,
|
||||||
network_name,
|
network_name,
|
||||||
tasks: vec![
|
tasks: vec![
|
||||||
event_task,
|
event_task,
|
||||||
@@ -484,6 +512,9 @@ impl Federation {
|
|||||||
device_sync_task,
|
device_sync_task,
|
||||||
capabilities_task,
|
capabilities_task,
|
||||||
similarity_task,
|
similarity_task,
|
||||||
|
similarity_dht_serve_task,
|
||||||
|
similarity_dht_maintenance_task,
|
||||||
|
similarity_dht_sync_task,
|
||||||
],
|
],
|
||||||
});
|
});
|
||||||
self.set_error(None);
|
self.set_error(None);
|
||||||
@@ -507,6 +538,20 @@ impl Federation {
|
|||||||
.context("federation is not running")
|
.context("federation is not running")
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn similarity_services(&self) -> Result<(Arc<MusicDhtService>, Arc<SimilarityDht>)> {
|
||||||
|
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<Self>) {
|
async fn spawn_sync_soon(self: &Arc<Self>) {
|
||||||
if let Ok(service) = self.service().await {
|
if let Ok(service) = self.service().await {
|
||||||
let fed = Arc::clone(self);
|
let fed = Arc::clone(self);
|
||||||
@@ -810,6 +855,7 @@ impl Federation {
|
|||||||
"endpoint_id": service.endpoint_id().to_string(),
|
"endpoint_id": service.endpoint_id().to_string(),
|
||||||
"connected_peers": peers,
|
"connected_peers": peers,
|
||||||
"known_contacts": service.known_peers().len(),
|
"known_contacts": service.known_peers().len(),
|
||||||
|
"similarity_routing_peers": running.similarity_dht.known_peers(),
|
||||||
"published_items": published,
|
"published_items": published,
|
||||||
"transport": self.transport_stats.snapshot(),
|
"transport": self.transport_stats.snapshot(),
|
||||||
})
|
})
|
||||||
@@ -854,8 +900,15 @@ impl Federation {
|
|||||||
crate::similarity::handle().enabled(),
|
crate::similarity::handle().enabled(),
|
||||||
"similarity search is disabled"
|
"similarity search is disabled"
|
||||||
);
|
);
|
||||||
let service = self.service().await?;
|
let (service, similarity_dht) = self.similarity_services().await?;
|
||||||
similarity::search(service, query, limit, Arc::clone(&self.transport_stats)).await
|
similarity::search(
|
||||||
|
service,
|
||||||
|
similarity_dht,
|
||||||
|
query,
|
||||||
|
limit,
|
||||||
|
Arc::clone(&self.transport_stats),
|
||||||
|
)
|
||||||
|
.await
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn fed_device_status(
|
pub async fn fed_device_status(
|
||||||
@@ -983,6 +1036,65 @@ impl Federation {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn similarity_route_sync_loop(
|
||||||
|
routing: Arc<SimilarityDht>,
|
||||||
|
manager: Arc<crate::similarity::Manager>,
|
||||||
|
) {
|
||||||
|
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(
|
async fn persist_content_id(
|
||||||
pool: &PgPool,
|
pool: &PgPool,
|
||||||
media_file_id: i64,
|
media_file_id: i64,
|
||||||
|
|||||||
+111
-35
@@ -7,7 +7,10 @@ use std::time::Duration;
|
|||||||
use anyhow::{Context as _, Result};
|
use anyhow::{Context as _, Result};
|
||||||
use futures_util::stream::{self, StreamExt as _};
|
use futures_util::stream::{self, StreamExt as _};
|
||||||
use music_dht::similarity::{self as wire, SimilarityHit, SimilarityRequest, SimilarityResponse};
|
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};
|
use crate::similarity::{Manager, QueryVector};
|
||||||
|
|
||||||
@@ -15,9 +18,11 @@ use super::TransportStats;
|
|||||||
|
|
||||||
pub use music_dht::similarity::SIMILARITY_ALPN;
|
pub use music_dht::similarity::SIMILARITY_ALPN;
|
||||||
|
|
||||||
const MAX_QUERY_PEERS: usize = 16;
|
const INITIAL_QUERY_PEERS: usize = 16;
|
||||||
const QUERY_CONCURRENCY: usize = 6;
|
const MAX_QUERY_PEERS: usize = 48;
|
||||||
|
const QUERY_CONCURRENCY: usize = 8;
|
||||||
const QUERY_TIMEOUT: Duration = Duration::from_secs(5);
|
const QUERY_TIMEOUT: Duration = Duration::from_secs(5);
|
||||||
|
const ROUTING_TIMEOUT: Duration = Duration::from_secs(5);
|
||||||
const MAX_PER_ARTIST: usize = 3;
|
const MAX_PER_ARTIST: usize = 3;
|
||||||
const MAX_NEAR_DUPLICATE_SIGNATURE_DISTANCE: u32 = 8;
|
const MAX_NEAR_DUPLICATE_SIGNATURE_DISTANCE: u32 = 8;
|
||||||
|
|
||||||
@@ -142,20 +147,49 @@ async fn serve_one(
|
|||||||
|
|
||||||
pub async fn search(
|
pub async fn search(
|
||||||
service: Arc<MusicDhtService>,
|
service: Arc<MusicDhtService>,
|
||||||
|
routing: Arc<SimilarityDht>,
|
||||||
query: QueryVector,
|
query: QueryVector,
|
||||||
limit: usize,
|
limit: usize,
|
||||||
transport: Arc<TransportStats>,
|
transport: Arc<TransportStats>,
|
||||||
) -> Result<Vec<RemoteSimilarityTrack>> {
|
) -> Result<Vec<RemoteSimilarityTrack>> {
|
||||||
let own = service.endpoint_id();
|
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 seen = HashSet::new();
|
||||||
|
let mut peers: Vec<QueryPeer> = 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
|
for peer in service
|
||||||
.connected_peers()
|
.connected_peers()
|
||||||
.into_iter()
|
.into_iter()
|
||||||
.chain(service.known_peers().into_iter().map(|peer| peer.peer_id))
|
.chain(service.known_peers().into_iter().map(|peer| peer.peer_id))
|
||||||
{
|
{
|
||||||
if peer != own && seen.insert(peer) {
|
if peer != own && seen.insert(peer) {
|
||||||
peers.push(peer);
|
peers.push(QueryPeer {
|
||||||
|
owner: peer,
|
||||||
|
ticket: None,
|
||||||
|
});
|
||||||
}
|
}
|
||||||
if peers.len() >= MAX_QUERY_PEERS {
|
if peers.len() >= MAX_QUERY_PEERS {
|
||||||
break;
|
break;
|
||||||
@@ -168,30 +202,40 @@ pub async fn search(
|
|||||||
limit.clamp(1, wire::MAX_SIMILARITY_RESULTS),
|
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::<Vec<_>>()
|
|
||||||
.await;
|
|
||||||
|
|
||||||
let mut hits = Vec::new();
|
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 {
|
for response in responses {
|
||||||
match response {
|
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"),
|
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));
|
hits.sort_by(|left, right| right.1.total_cmp(&left.1));
|
||||||
let mut dedup = HashSet::new();
|
let mut dedup = HashSet::new();
|
||||||
let mut signatures = vec![query_signature];
|
let mut signatures = vec![query_signature];
|
||||||
@@ -241,22 +285,54 @@ pub async fn search(
|
|||||||
Ok(tracks)
|
Ok(tracks)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type PeerHits = Vec<(
|
||||||
|
RemoteSimilarityTrack,
|
||||||
|
f32,
|
||||||
|
Option<[u8; wire::SIMILARITY_SIGNATURE_BYTES]>,
|
||||||
|
)>;
|
||||||
|
|
||||||
|
#[derive(Clone)]
|
||||||
|
struct QueryPeer {
|
||||||
|
owner: EndpointId,
|
||||||
|
ticket: Option<PeerTicket>,
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn query_peers(
|
||||||
|
service: Arc<MusicDhtService>,
|
||||||
|
peers: &[QueryPeer],
|
||||||
|
request: Arc<SimilarityRequest>,
|
||||||
|
transport: Arc<TransportStats>,
|
||||||
|
) -> Vec<Result<PeerHits>> {
|
||||||
|
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(
|
async fn query_peer(
|
||||||
service: Arc<MusicDhtService>,
|
service: Arc<MusicDhtService>,
|
||||||
owner: EndpointId,
|
peer: QueryPeer,
|
||||||
request: &SimilarityRequest,
|
request: &SimilarityRequest,
|
||||||
transport: Arc<TransportStats>,
|
transport: Arc<TransportStats>,
|
||||||
) -> Result<
|
) -> Result<PeerHits> {
|
||||||
Vec<(
|
let owner = peer.owner;
|
||||||
RemoteSimilarityTrack,
|
let mut stream = match peer.ticket {
|
||||||
f32,
|
Some(ticket) => service.open_stream_to(&ticket, SIMILARITY_ALPN).await,
|
||||||
Option<[u8; wire::SIMILARITY_SIGNATURE_BYTES]>,
|
None => service.open_stream(owner, SIMILARITY_ALPN).await,
|
||||||
)>,
|
}
|
||||||
> {
|
.map_err(|error| anyhow::anyhow!("cannot reach similarity peer: {error}"))?;
|
||||||
let mut stream = 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);
|
super::record_stream_transport(&transport, "similarity", "outbound", "open", &stream);
|
||||||
let response = wire::exchange(&mut stream, request).await?;
|
let response = wire::exchange(&mut stream, request).await?;
|
||||||
super::record_stream_transport(&transport, "similarity", "outbound", "done", &stream);
|
super::record_stream_transport(&transport, "similarity", "outbound", "done", &stream);
|
||||||
|
|||||||
@@ -2541,6 +2541,34 @@ pub mod db_migrations {
|
|||||||
&[Operation::custom(create_similarity_embeddings).build()];
|
&[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] = &[
|
pub const MIGRATIONS: &[&SyncDynMigration] = &[
|
||||||
&M0006CreateMediaFile,
|
&M0006CreateMediaFile,
|
||||||
&M0007CreateArtist,
|
&M0007CreateArtist,
|
||||||
@@ -2575,5 +2603,6 @@ pub mod db_migrations {
|
|||||||
&M0041CreateSyncedListenHistory,
|
&M0041CreateSyncedListenHistory,
|
||||||
&M0042RepairLegacyListenQualification,
|
&M0042RepairLegacyListenQualification,
|
||||||
&M0043CreateSimilarityEmbeddings,
|
&M0043CreateSimilarityEmbeddings,
|
||||||
|
&M0044AddSimilarityRoutingSignature,
|
||||||
];
|
];
|
||||||
}
|
}
|
||||||
|
|||||||
+90
-3
@@ -283,6 +283,90 @@ impl Manager {
|
|||||||
lock(&self.status).clone()
|
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<Vec<[u8; 32]>> {
|
||||||
|
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::<i64, _>(0),
|
||||||
|
row.get::<i32, _>(1),
|
||||||
|
row.get::<Vec<u8>, _>(2),
|
||||||
|
)
|
||||||
|
})
|
||||||
|
.collect::<Vec<_>>();
|
||||||
|
|
||||||
|
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::<Result<Vec<_>>>()
|
||||||
|
})
|
||||||
|
.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<u8>>(
|
||||||
|
"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<Self>) {
|
pub fn start(self: &Arc<Self>) {
|
||||||
let generation = self.generation.fetch_add(1, Ordering::AcqRel) + 1;
|
let generation = self.generation.fetch_add(1, Ordering::AcqRel) + 1;
|
||||||
let manager = Arc::clone(self);
|
let manager = Arc::clone(self);
|
||||||
@@ -846,14 +930,16 @@ async fn store_embedding(
|
|||||||
vector.iter().all(|value| value.is_finite()),
|
vector.iter().all(|value| value.is_finite()),
|
||||||
"embedding contains a non-finite value"
|
"embedding contains a non-finite value"
|
||||||
);
|
);
|
||||||
|
let routing_signature = music_dht::similarity_lsh::routing_signature(vector)?;
|
||||||
sqlx::query(
|
sqlx::query(
|
||||||
"INSERT INTO furumusic__track_embedding
|
"INSERT INTO furumusic__track_embedding
|
||||||
(track_id, profile_id, dimensions, vector, source_sha256,
|
(track_id, profile_id, dimensions, vector, routing_signature,
|
||||||
source_content_id, computed_at)
|
source_sha256, source_content_id, computed_at)
|
||||||
VALUES ($1, $2, $3, $4, $5, $6, $7)
|
VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
|
||||||
ON CONFLICT (track_id, profile_id) DO UPDATE SET
|
ON CONFLICT (track_id, profile_id) DO UPDATE SET
|
||||||
dimensions = EXCLUDED.dimensions,
|
dimensions = EXCLUDED.dimensions,
|
||||||
vector = EXCLUDED.vector,
|
vector = EXCLUDED.vector,
|
||||||
|
routing_signature = EXCLUDED.routing_signature,
|
||||||
source_sha256 = EXCLUDED.source_sha256,
|
source_sha256 = EXCLUDED.source_sha256,
|
||||||
source_content_id = EXCLUDED.source_content_id,
|
source_content_id = EXCLUDED.source_content_id,
|
||||||
computed_at = EXCLUDED.computed_at",
|
computed_at = EXCLUDED.computed_at",
|
||||||
@@ -862,6 +948,7 @@ async fn store_embedding(
|
|||||||
.bind(profile_id)
|
.bind(profile_id)
|
||||||
.bind(vector.len() as i32)
|
.bind(vector.len() as i32)
|
||||||
.bind(embedding_to_bytes(vector))
|
.bind(embedding_to_bytes(vector))
|
||||||
|
.bind(routing_signature.as_slice())
|
||||||
.bind(&track.source_sha256)
|
.bind(&track.source_sha256)
|
||||||
.bind(&track.source_content_id)
|
.bind(&track.source_content_id)
|
||||||
.bind(now_iso())
|
.bind(now_iso())
|
||||||
|
|||||||
@@ -2210,7 +2210,7 @@ tbody tr:hover {
|
|||||||
<span x-text="settingsDraft.similarity_enabled ? 'Enabled for this instance' : 'Disabled'"></span>
|
<span x-text="settingsDraft.similarity_enabled ? 'Enabled for this instance' : 'Disabled'"></span>
|
||||||
<input type="checkbox" x-model="settingsDraft.similarity_enabled" />
|
<input type="checkbox" x-model="settingsDraft.similarity_enabled" />
|
||||||
</div>
|
</div>
|
||||||
<div class="setting-help">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.</div>
|
<div class="setting-help">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.</div>
|
||||||
</div>
|
</div>
|
||||||
<div class="setting-field settings-wide">
|
<div class="setting-field settings-wide">
|
||||||
<label>
|
<label>
|
||||||
@@ -2368,6 +2368,7 @@ tbody tr:hover {
|
|||||||
<div class="probe-row"><span>Network</span><strong x-text="(federationStatus.node && federationStatus.node.network) || '-'"></strong></div>
|
<div class="probe-row"><span>Network</span><strong x-text="(federationStatus.node && federationStatus.node.network) || '-'"></strong></div>
|
||||||
<div class="probe-row"><span>Connected peers</span><strong x-text="federationStatus.node && federationStatus.node.connected_peers ? federationStatus.node.connected_peers.length : 0"></strong></div>
|
<div class="probe-row"><span>Connected peers</span><strong x-text="federationStatus.node && federationStatus.node.connected_peers ? federationStatus.node.connected_peers.length : 0"></strong></div>
|
||||||
<div class="probe-row"><span>Known contacts</span><strong x-text="(federationStatus.node && federationStatus.node.known_contacts) ?? '-'"></strong></div>
|
<div class="probe-row"><span>Known contacts</span><strong x-text="(federationStatus.node && federationStatus.node.known_contacts) ?? '-'"></strong></div>
|
||||||
|
<div class="probe-row"><span>Similarity routing peers</span><strong x-text="(federationStatus.node && federationStatus.node.similarity_routing_peers) ?? '-'"></strong></div>
|
||||||
<div class="probe-row"><span>Published items</span><strong x-text="(federationStatus.node && federationStatus.node.published_items) ?? '-'"></strong></div>
|
<div class="probe-row"><span>Published items</span><strong x-text="(federationStatus.node && federationStatus.node.published_items) ?? '-'"></strong></div>
|
||||||
<div class="probe-row"><span>Last sync</span><strong x-text="federationStatus.last_sync || 'not yet'"></strong></div>
|
<div class="probe-row"><span>Last sync</span><strong x-text="federationStatus.last_sync || 'not yet'"></strong></div>
|
||||||
</div>
|
</div>
|
||||||
|
|||||||
Reference in New Issue
Block a user