Compare commits
3
Commits
e95d2e7fe1
...
federation
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
079d87a831 | ||
|
|
add764e51d | ||
|
|
87cb7fe74c |
@@ -123,6 +123,22 @@ The DHT is the distributed index. Nodes publish compact searchable
|
||||
descriptions of their local library and query the network without contacting a
|
||||
central search service.
|
||||
|
||||
Similarity discovery uses a separate schema-independent DHT overlay. A node
|
||||
derives a deterministic 256-bit routing signature from every local embedding,
|
||||
groups them into fixed two-level LSH buckets, and publishes compact summaries
|
||||
containing only fine-bucket representatives, its peer identity, and the ticket
|
||||
needed to dial a previously unknown owner. The owner signs every summary with
|
||||
its existing transport key, so storage peers can relay and cache it but cannot
|
||||
impersonate or modify it. Summaries expire with the ordinary library-record TTL
|
||||
and are replaceable federation cache, never local library authority.
|
||||
|
||||
A similarity search first performs bounded multi-probe LSH lookups to rank
|
||||
likely owners, then sends the existing normalized-vector request directly to
|
||||
at most 16 peers initially and 48 on fallback. Known peers remain a rollout
|
||||
fallback. The model, preprocessing, durable embeddings, exact cosine search,
|
||||
and consent policy remain client-owned; `music-dht` owns only compatible
|
||||
routing math, signed records, replication, and wire bounds.
|
||||
|
||||
Once a peer is known, communication moves to direct P2P streams provided by
|
||||
iroh through `music-dht`. Furumi defines separate application protocols for
|
||||
catalog requests, audio transfer, and trusted-device synchronization. This
|
||||
|
||||
+7
-1
@@ -20,12 +20,18 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
||||
- Track-seeded similarity search from the track-information popup, including
|
||||
bounded federated queries to compatible known peers.
|
||||
- The `furumi-fd/similarity/1` protocol in the visible protocol-version status.
|
||||
- Decentralized `similarity_dht` routing with signed anonymous two-level LSH
|
||||
summaries, multi-probe lookup beyond the locally known peer set, and known-
|
||||
peer fallback during gradual network upgrades.
|
||||
|
||||
### Changed
|
||||
|
||||
- Similarity wire types, bounds, validation, and stream framing now come from
|
||||
the shared `music-dht 0.3.1` API so native, web, and future clients can
|
||||
the shared `music-dht 0.4.0` API so native, web, and future clients can
|
||||
interoperate without sharing an embedding implementation.
|
||||
- Existing SQLite embeddings are backfilled once with compact 256-bit routing
|
||||
signatures; new embeddings store them immediately without changing exact
|
||||
local cosine search.
|
||||
- A similarity result page keeps the source track first as query context while
|
||||
excluding it from the actual nearest-neighbor ranking, labels the mode as
|
||||
`Search similar to`, and suppresses near-identical embeddings across releases
|
||||
|
||||
Generated
+7
-3
@@ -1520,7 +1520,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 = "c3e690b370c505d153bef214b21a8f2aa55d667367ac1e16bde8bc0de88963c2"
|
||||
dependencies = [
|
||||
"blake3",
|
||||
"data-encoding",
|
||||
@@ -1647,7 +1649,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "furumi_tui"
|
||||
version = "0.2.5"
|
||||
version = "0.2.6"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"blake3",
|
||||
@@ -3138,7 +3140,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 = "0c5b429b90a8f1b0980b3a35a6fa5445d7a275c737eb04db18db4d7f14c81478"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"blake3",
|
||||
|
||||
+1
-1
@@ -21,7 +21,7 @@ image = { version = "0.25.10", default-features = false, features = ["jpeg", "pn
|
||||
lofty = "0.22"
|
||||
# P2P federation: library index in a shared DHT + audio streaming between
|
||||
# peers (same protocol as furumi-fd).
|
||||
music-dht = "0.3.1"
|
||||
music-dht = "0.4.0"
|
||||
ratatui = "0.30.1"
|
||||
reqwest = { version = "0.12.28", default-features = false, features = ["rustls-tls", "stream"] }
|
||||
rhai = { version = "1", features = ["sync"] }
|
||||
|
||||
@@ -72,6 +72,9 @@ without rebuilding the player.
|
||||
Optional similarity search calculates versioned embeddings for local tracks
|
||||
in the background and keeps them in SQLite. It works offline; after a separate
|
||||
privacy consent it can also ask a bounded set of federation peers for matches.
|
||||
Compatible peers are selected through signed, anonymous LSH summaries in a
|
||||
decentralized DHT; no central recommendation index or shared calibration file
|
||||
is required.
|
||||
The first selectable model is downloaded on demand and is licensed separately
|
||||
by MTG under CC BY-NC-SA 4.0 (a proprietary license is also available from
|
||||
MTG); Furumi itself remains WTFPL.
|
||||
|
||||
+9
-15
@@ -8,7 +8,6 @@ use crate::app::Runtime;
|
||||
use crate::app::command::{self, Command, Parsed};
|
||||
use crate::app::event::AppEvent;
|
||||
use crate::app::state::{AppState, GlobalView, SearchState, Tab};
|
||||
use crate::library::models::SearchResults;
|
||||
|
||||
const SEARCH_DEBOUNCE: Duration = Duration::from_millis(180);
|
||||
const SEARCH_LIMIT: i64 = 12;
|
||||
@@ -99,6 +98,10 @@ fn set_view_cursor_zero(state: &mut AppState) {
|
||||
/// and the receiver drops responses that arrive out of date.
|
||||
pub(super) fn schedule_search(state: &mut AppState, runtime: &Runtime) {
|
||||
state.search.similarity_source = None;
|
||||
state.search.similarity_source_track = None;
|
||||
state.search.similarity_tracks.clear();
|
||||
state.search.similarity_stats = None;
|
||||
state.search.similarity_error = None;
|
||||
let seq = runtime.search_seq.fetch_add(1, Ordering::SeqCst) + 1;
|
||||
let query = state.search.query.clone();
|
||||
if query.is_empty() {
|
||||
@@ -160,6 +163,10 @@ pub(super) fn schedule_similarity_search(
|
||||
format!("{} — {artist}", track.title)
|
||||
};
|
||||
state.search.similarity_source = Some(track.id);
|
||||
state.search.similarity_source_track = Some(track.clone());
|
||||
state.search.similarity_tracks.clear();
|
||||
state.search.similarity_stats = None;
|
||||
state.search.similarity_error = None;
|
||||
state.search.loading = true;
|
||||
state.search.results = None;
|
||||
state.search.fed_tracks.clear();
|
||||
@@ -178,23 +185,10 @@ pub(super) fn schedule_similarity_search(
|
||||
let similarity = Arc::clone(&runtime.similarity);
|
||||
let tx = runtime.event_tx.clone();
|
||||
let track_id = track.id;
|
||||
let source_track = track.clone();
|
||||
tokio::task::spawn_blocking(move || {
|
||||
let result = similarity
|
||||
.search_track(track_id, 49)
|
||||
.map(|(matches, query)| {
|
||||
let mut tracks = Vec::with_capacity(1 + matches.len());
|
||||
tracks.push(source_track);
|
||||
tracks.extend(matches.into_iter().map(|found| found.track));
|
||||
(
|
||||
SearchResults {
|
||||
artists: Vec::new(),
|
||||
releases: Vec::new(),
|
||||
tracks,
|
||||
},
|
||||
query,
|
||||
)
|
||||
});
|
||||
.map(|(matches, query)| (matches, query));
|
||||
let (result, query) = match result {
|
||||
Ok((results, query)) => (Ok(results), Some(query)),
|
||||
Err(err) => (Err(format!("{err:#}")), None),
|
||||
|
||||
+6
-1
@@ -36,9 +36,14 @@ pub enum AppEvent {
|
||||
/// text search so stale pages cannot overwrite a newer request.
|
||||
SimilaritySearchLoaded {
|
||||
seq: u64,
|
||||
result: Result<SearchResults, String>,
|
||||
result: Result<Vec<crate::similarity::SimilarTrack>, String>,
|
||||
query: Option<crate::similarity::QueryVector>,
|
||||
},
|
||||
/// Federated candidates and diagnostics for a track-seeded search.
|
||||
FedSimilaritySearchLoaded {
|
||||
seq: u64,
|
||||
result: Result<crate::federation::FedSimilaritySearchResults, String>,
|
||||
},
|
||||
SimilarityStatus(crate::similarity::SimilarityStatus),
|
||||
/// `None` is emitted after clearing every stored embedding.
|
||||
SimilarityProfileActivated(Option<String>),
|
||||
|
||||
+225
-13
@@ -44,6 +44,10 @@ pub struct Runtime {
|
||||
pub similarity: Arc<crate::similarity::Manager>,
|
||||
/// When the last Federation-tab status snapshot was requested.
|
||||
pub fed_status_at: Option<std::time::Instant>,
|
||||
/// Keeps the last successful local-data snapshot visible while a newer
|
||||
/// one is calculated and collapses bursts of library-change events.
|
||||
pub local_library_stats_refreshing: Arc<std::sync::atomic::AtomicBool>,
|
||||
pub local_library_stats_refresh_requested: Arc<std::sync::atomic::AtomicBool>,
|
||||
pub library_network_refresh_at: Option<std::time::Instant>,
|
||||
pub library_network_refreshing: Arc<std::sync::atomic::AtomicBool>,
|
||||
pub library_network_cursors:
|
||||
@@ -102,11 +106,35 @@ fn refresh_local_content_ids(runtime: &Runtime) {
|
||||
}
|
||||
|
||||
fn refresh_local_library_stats(runtime: &Runtime) {
|
||||
runtime
|
||||
.local_library_stats_refresh_requested
|
||||
.store(true, std::sync::atomic::Ordering::Release);
|
||||
if runtime
|
||||
.local_library_stats_refreshing
|
||||
.swap(true, std::sync::atomic::Ordering::AcqRel)
|
||||
{
|
||||
return;
|
||||
}
|
||||
let library = Arc::clone(&runtime.library);
|
||||
let tx = runtime.event_tx.clone();
|
||||
let refreshing = Arc::clone(&runtime.local_library_stats_refreshing);
|
||||
let requested = Arc::clone(&runtime.local_library_stats_refresh_requested);
|
||||
tokio::task::spawn_blocking(move || {
|
||||
let result = library.local_stats().map_err(err_string);
|
||||
let _ = tx.send(AppEvent::LocalLibraryStatsLoaded(result));
|
||||
loop {
|
||||
requested.store(false, std::sync::atomic::Ordering::Release);
|
||||
let result = library.local_stats().map_err(err_string);
|
||||
let _ = tx.send(AppEvent::LocalLibraryStatsLoaded(result));
|
||||
if requested.load(std::sync::atomic::Ordering::Acquire) {
|
||||
continue;
|
||||
}
|
||||
refreshing.store(false, std::sync::atomic::Ordering::Release);
|
||||
if requested.swap(false, std::sync::atomic::Ordering::AcqRel)
|
||||
&& !refreshing.swap(true, std::sync::atomic::Ordering::AcqRel)
|
||||
{
|
||||
continue;
|
||||
}
|
||||
break;
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
@@ -303,6 +331,8 @@ pub async fn run(
|
||||
federation,
|
||||
similarity,
|
||||
fed_status_at: None,
|
||||
local_library_stats_refreshing: Arc::new(std::sync::atomic::AtomicBool::new(false)),
|
||||
local_library_stats_refresh_requested: Arc::new(std::sync::atomic::AtomicBool::new(false)),
|
||||
library_network_refresh_at: None,
|
||||
library_network_refreshing: Arc::new(std::sync::atomic::AtomicBool::new(false)),
|
||||
library_network_cursors: Arc::new(std::sync::Mutex::new(std::collections::HashMap::new())),
|
||||
@@ -2746,7 +2776,12 @@ fn on_library_changed(state: &mut AppState, runtime: &mut Runtime) {
|
||||
// until then.
|
||||
state.likes_loaded = false;
|
||||
state.local_content_ids_loaded = false;
|
||||
state.local_library_stats = None;
|
||||
// Refresh in place: status cards keep the last successful snapshot
|
||||
// instead of flashing `loading` for every background library event.
|
||||
if state.local_library_stats.is_none() {
|
||||
state.local_library_stats = Some(state::Loadable::Loading);
|
||||
}
|
||||
refresh_local_library_stats(runtime);
|
||||
|
||||
// Fresh copies of whatever sits in the queue. Federated placeholders
|
||||
// and ephemeral tracks (negative ids) are not library rows and keep
|
||||
@@ -3313,6 +3348,39 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent
|
||||
Err(message) => tracing::warn!(%message, "federated search failed"),
|
||||
}
|
||||
}
|
||||
AppEvent::FedSimilaritySearchLoaded { seq, result } => {
|
||||
if runtime.search_seq.load(std::sync::atomic::Ordering::SeqCst) != seq {
|
||||
return;
|
||||
}
|
||||
state.search.fed_loading = false;
|
||||
let selected_key = state.global.stack.last().and_then(|view| match view {
|
||||
state::GlobalView::Search { cursor } => state.search.similarity_key(*cursor),
|
||||
_ => None,
|
||||
});
|
||||
match result {
|
||||
Ok(results) => {
|
||||
let remote = results.tracks.into_iter().map(|hit| {
|
||||
state::SimilaritySearchHit::Federated {
|
||||
track: hit.track,
|
||||
score: hit.score,
|
||||
embedding_signature: hit.embedding_signature,
|
||||
}
|
||||
});
|
||||
state.search.similarity_tracks.extend(remote);
|
||||
rank_similarity_search_tracks(
|
||||
&mut state.search.similarity_tracks,
|
||||
state.similarity.settings.max_tracks_per_artist,
|
||||
);
|
||||
state.search.similarity_stats = Some(results.stats);
|
||||
state.search.similarity_error = None;
|
||||
}
|
||||
Err(message) => {
|
||||
tracing::warn!(%message, "federated similarity search failed");
|
||||
state.search.similarity_error = Some(message);
|
||||
}
|
||||
}
|
||||
restore_similarity_cursor(state, selected_key.as_deref());
|
||||
}
|
||||
AppEvent::FedTrackResolved {
|
||||
placeholder_id,
|
||||
resolve_key,
|
||||
@@ -3661,7 +3729,20 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent
|
||||
}
|
||||
state.search.loading = false;
|
||||
match result {
|
||||
Ok(results) => state.search.results = Some(results),
|
||||
Ok(results) => {
|
||||
state.search.similarity_tracks = results
|
||||
.into_iter()
|
||||
.map(|hit| state::SimilaritySearchHit::Local {
|
||||
track: hit.track,
|
||||
score: hit.score,
|
||||
embedding_signature: hit.embedding_signature,
|
||||
})
|
||||
.collect();
|
||||
rank_similarity_search_tracks(
|
||||
&mut state.search.similarity_tracks,
|
||||
state.similarity.settings.max_tracks_per_artist,
|
||||
);
|
||||
}
|
||||
Err(message) => {
|
||||
state.status_message = Some(format!("similarity search failed: {message}"));
|
||||
return;
|
||||
@@ -3679,7 +3760,7 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent
|
||||
.search_similar(query, 50)
|
||||
.await
|
||||
.map_err(|err| format!("{err:#}"));
|
||||
let _ = tx.send(AppEvent::FedSearchLoaded { seq, result });
|
||||
let _ = tx.send(AppEvent::FedSimilaritySearchLoaded { seq, result });
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -3792,15 +3873,17 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent
|
||||
tracing::warn!(%message, "local content id load failed");
|
||||
}
|
||||
},
|
||||
AppEvent::LocalLibraryStatsLoaded(result) => {
|
||||
state.local_library_stats = Some(match result {
|
||||
Ok(stats) => state::Loadable::Ready(stats),
|
||||
Err(message) => {
|
||||
tracing::warn!(%message, "local library stats load failed");
|
||||
state::Loadable::Failed(message)
|
||||
AppEvent::LocalLibraryStatsLoaded(result) => match result {
|
||||
Ok(stats) => {
|
||||
state.local_library_stats = Some(state::Loadable::Ready(stats));
|
||||
}
|
||||
Err(message) => {
|
||||
tracing::warn!(%message, "local library stats load failed");
|
||||
if !matches!(state.local_library_stats, Some(state::Loadable::Ready(_))) {
|
||||
state.local_library_stats = Some(state::Loadable::Failed(message));
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
},
|
||||
AppEvent::LocalContentAvailable { content_id } => {
|
||||
if let Some(content_id) = music_dht::normalize_content_id(&content_id) {
|
||||
state.local_content_ids.insert(content_id);
|
||||
@@ -3970,6 +4053,135 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent
|
||||
}
|
||||
}
|
||||
|
||||
fn rank_similarity_search_tracks(
|
||||
tracks: &mut Vec<state::SimilaritySearchHit>,
|
||||
max_tracks_per_artist: usize,
|
||||
) {
|
||||
const RESULT_LIMIT: usize = 49;
|
||||
const MAX_NEAR_DUPLICATE_SIGNATURE_DISTANCE: u32 = 8;
|
||||
|
||||
tracks.sort_by(|left, right| right.score().total_cmp(&left.score()));
|
||||
let candidates = std::mem::take(tracks);
|
||||
let mut content = std::collections::HashSet::new();
|
||||
let mut signatures = Vec::new();
|
||||
let mut artist_counts: std::collections::HashMap<String, usize> = Default::default();
|
||||
for hit in candidates {
|
||||
if !content.insert(hit.content_key()) {
|
||||
continue;
|
||||
}
|
||||
if hit.embedding_signature().is_some_and(|candidate| {
|
||||
signatures.iter().any(|existing| {
|
||||
music_dht::similarity::signature_distance(&candidate, existing)
|
||||
<= MAX_NEAR_DUPLICATE_SIGNATURE_DISTANCE
|
||||
})
|
||||
}) {
|
||||
continue;
|
||||
}
|
||||
let artist = hit.primary_artist_key();
|
||||
let count = artist_counts.entry(artist.clone()).or_default();
|
||||
if !artist.is_empty() && *count >= max_tracks_per_artist.clamp(1, RESULT_LIMIT) {
|
||||
continue;
|
||||
}
|
||||
*count += 1;
|
||||
if let Some(signature) = hit.embedding_signature() {
|
||||
signatures.push(signature);
|
||||
}
|
||||
tracks.push(hit);
|
||||
if tracks.len() >= RESULT_LIMIT {
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn restore_similarity_cursor(state: &mut AppState, selected_key: Option<&str>) {
|
||||
let selected_index = selected_key.and_then(|key| state.search.similarity_index_for_key(key));
|
||||
let len = state.search.similarity_len();
|
||||
if let Some(state::GlobalView::Search { cursor }) = state.global.stack.last_mut() {
|
||||
*cursor = selected_index.unwrap_or(*cursor).min(len.saturating_sub(1));
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod similarity_search_tests {
|
||||
use super::*;
|
||||
use crate::library::models::{ArtistRef, TrackItem};
|
||||
|
||||
fn local_hit(id: i64, artist: &str, score: f32, signature: u8) -> state::SimilaritySearchHit {
|
||||
state::SimilaritySearchHit::Local {
|
||||
track: TrackItem {
|
||||
id,
|
||||
title: format!("local {id}"),
|
||||
track_number: None,
|
||||
disc_number: None,
|
||||
duration_seconds: 1.0,
|
||||
artists: vec![ArtistRef {
|
||||
id,
|
||||
name: artist.to_string(),
|
||||
}],
|
||||
featured_artists: Vec::new(),
|
||||
release_id: id,
|
||||
release_title: "release".to_string(),
|
||||
release_year: None,
|
||||
file_path: format!("/music/{id}"),
|
||||
content_id: Some(format!("local-{id}")),
|
||||
cover_path: None,
|
||||
audio_format: None,
|
||||
audio_bitrate: None,
|
||||
audio_sample_rate: None,
|
||||
audio_bit_depth: None,
|
||||
file_size_bytes: None,
|
||||
play_count: 0,
|
||||
fed: None,
|
||||
},
|
||||
score,
|
||||
embedding_signature: [signature; music_dht::similarity::SIMILARITY_SIGNATURE_BYTES],
|
||||
}
|
||||
}
|
||||
|
||||
fn remote_hit(artist: &str, score: f32, signature: u8) -> state::SimilaritySearchHit {
|
||||
state::SimilaritySearchHit::Federated {
|
||||
track: crate::federation::FedTrack {
|
||||
item_id: format!("remote-{signature}"),
|
||||
owner: "peer".to_string(),
|
||||
own: false,
|
||||
title: format!("remote {signature}"),
|
||||
artist_names: vec![artist.to_string()],
|
||||
featured_artist_names: Vec::new(),
|
||||
year: None,
|
||||
duration_seconds: Some(1),
|
||||
content_id: Some(format!("remote-{signature}")),
|
||||
release_title: None,
|
||||
track_number: None,
|
||||
disc_number: None,
|
||||
},
|
||||
score,
|
||||
embedding_signature: Some(
|
||||
[signature; music_dht::similarity::SIMILARITY_SIGNATURE_BYTES],
|
||||
),
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn similarity_results_rank_local_and_remote_together_with_one_artist_cap() {
|
||||
let mut tracks = vec![
|
||||
local_hit(1, "same artist", 0.70, 1),
|
||||
remote_hit("other artist", 0.90, 2),
|
||||
remote_hit("same artist", 0.80, 3),
|
||||
local_hit(2, "same artist", 0.60, 4),
|
||||
];
|
||||
|
||||
rank_similarity_search_tracks(&mut tracks, 1);
|
||||
|
||||
assert_eq!(tracks.len(), 2);
|
||||
assert_eq!(tracks[0].score(), 0.90);
|
||||
assert_eq!(tracks[1].score(), 0.80);
|
||||
assert!(matches!(
|
||||
tracks[0],
|
||||
state::SimilaritySearchHit::Federated { .. }
|
||||
));
|
||||
}
|
||||
}
|
||||
|
||||
/// Mirror the playback state to the OS now-playing surface. `force` skips
|
||||
/// the position throttle (track switches, pauses).
|
||||
fn push_media_update(state: &AppState, runtime: &mut Runtime, force: bool) {
|
||||
|
||||
@@ -592,6 +592,36 @@ fn handle_fed_input(
|
||||
state.popup = Some(Popup::FedInput { field, input });
|
||||
}
|
||||
},
|
||||
FedInputField::SimilarityMinimumScore => match value.parse::<f32>() {
|
||||
Ok(score) if score.is_finite() && (0.0..=1.0).contains(&score) => {
|
||||
state.similarity.settings.minimum_score = score;
|
||||
super::perform_effect(
|
||||
state,
|
||||
runtime,
|
||||
crate::app::update::Effect::SimilarityApplySettings,
|
||||
);
|
||||
}
|
||||
_ => {
|
||||
state.status_message =
|
||||
Some("minimum similarity must be a number from 0.00 to 1.00".into());
|
||||
state.popup = Some(Popup::FedInput { field, input });
|
||||
}
|
||||
},
|
||||
FedInputField::SimilarityMaxTracksPerArtist => match value.parse::<usize>() {
|
||||
Ok(limit @ 1..=50) => {
|
||||
state.similarity.settings.max_tracks_per_artist = limit;
|
||||
super::perform_effect(
|
||||
state,
|
||||
runtime,
|
||||
crate::app::update::Effect::SimilarityApplySettings,
|
||||
);
|
||||
}
|
||||
_ => {
|
||||
state.status_message =
|
||||
Some("tracks per artist must be a number from 1 to 50".into());
|
||||
state.popup = Some(Popup::FedInput { field, input });
|
||||
}
|
||||
},
|
||||
FedInputField::ConnectTicket => {
|
||||
if value.is_empty() {
|
||||
state.status_message = Some("ticket is empty".into());
|
||||
|
||||
+149
-1
@@ -509,6 +509,8 @@ pub enum TrackSelectionScope {
|
||||
Release(i64),
|
||||
Playlist(i64),
|
||||
Queue,
|
||||
/// The unified local + federated similar-track result list.
|
||||
SimilaritySearch,
|
||||
/// The federated section of the search results (its tracks).
|
||||
FedSearch,
|
||||
/// The tracklist of the open federated release view.
|
||||
@@ -817,6 +819,8 @@ impl StatusDetailFocus {
|
||||
pub enum FedInputField {
|
||||
MusicDirectory,
|
||||
SimilarityWorkers,
|
||||
SimilarityMinimumScore,
|
||||
SimilarityMaxTracksPerArtist,
|
||||
NetworkId,
|
||||
ConnectTicket,
|
||||
DeviceName,
|
||||
@@ -829,6 +833,8 @@ impl FedInputField {
|
||||
match self {
|
||||
FedInputField::MusicDirectory => "Music save directory",
|
||||
FedInputField::SimilarityWorkers => "Similarity background workers",
|
||||
FedInputField::SimilarityMinimumScore => "Minimum similarity score",
|
||||
FedInputField::SimilarityMaxTracksPerArtist => "Tracks per artist",
|
||||
FedInputField::NetworkId => "Network ID",
|
||||
FedInputField::ConnectTicket => "Connect to peer (paste ticket)",
|
||||
FedInputField::DeviceName => "Device name",
|
||||
@@ -845,6 +851,12 @@ impl FedInputField {
|
||||
FedInputField::SimilarityWorkers => {
|
||||
"Enter the maximum number of tracks processed in parallel, from 1 to 16. The change takes effect immediately."
|
||||
}
|
||||
FedInputField::SimilarityMinimumScore => {
|
||||
"Enter the minimum cosine similarity from 0.00 to 1.00. Lower values show broader matches; higher values hide weak matches. Embeddings are not recalculated."
|
||||
}
|
||||
FedInputField::SimilarityMaxTracksPerArtist => {
|
||||
"Enter how many tracks by one primary artist may appear in similarity results, from 1 to 50. Embeddings are not recalculated."
|
||||
}
|
||||
FedInputField::NetworkId => {
|
||||
"A unique network id. It must match exactly on every client that should see and connect to the same peers."
|
||||
}
|
||||
@@ -880,15 +892,19 @@ pub enum SimilarityRow {
|
||||
Toggle,
|
||||
Model,
|
||||
Profile,
|
||||
MinimumScore,
|
||||
MaxTracksPerArtist,
|
||||
Workers,
|
||||
Clear,
|
||||
}
|
||||
|
||||
impl SimilarityRow {
|
||||
pub const ALL: [SimilarityRow; 5] = [
|
||||
pub const ALL: [SimilarityRow; 7] = [
|
||||
SimilarityRow::Toggle,
|
||||
SimilarityRow::Model,
|
||||
SimilarityRow::Profile,
|
||||
SimilarityRow::MinimumScore,
|
||||
SimilarityRow::MaxTracksPerArtist,
|
||||
SimilarityRow::Workers,
|
||||
SimilarityRow::Clear,
|
||||
];
|
||||
@@ -1211,6 +1227,85 @@ mod cmdline_history_tests {
|
||||
}
|
||||
|
||||
/// Live search state driven by the `:/query` command.
|
||||
#[derive(Debug, Clone)]
|
||||
pub enum SimilaritySearchHit {
|
||||
Local {
|
||||
track: TrackItem,
|
||||
score: f32,
|
||||
embedding_signature: [u8; music_dht::similarity::SIMILARITY_SIGNATURE_BYTES],
|
||||
},
|
||||
Federated {
|
||||
track: crate::federation::FedTrack,
|
||||
score: f32,
|
||||
embedding_signature: Option<[u8; music_dht::similarity::SIMILARITY_SIGNATURE_BYTES]>,
|
||||
},
|
||||
}
|
||||
|
||||
impl SimilaritySearchHit {
|
||||
pub fn score(&self) -> f32 {
|
||||
match self {
|
||||
Self::Local { score, .. } | Self::Federated { score, .. } => *score,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn embedding_signature(
|
||||
&self,
|
||||
) -> Option<[u8; music_dht::similarity::SIMILARITY_SIGNATURE_BYTES]> {
|
||||
match self {
|
||||
Self::Local {
|
||||
embedding_signature,
|
||||
..
|
||||
} => Some(*embedding_signature),
|
||||
Self::Federated {
|
||||
embedding_signature,
|
||||
..
|
||||
} => *embedding_signature,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn track_item(&self) -> TrackItem {
|
||||
match self {
|
||||
Self::Local { track, .. } => track.clone(),
|
||||
Self::Federated { track, .. } => crate::federation::pending_track(track),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn federated_track(&self) -> Option<&crate::federation::FedTrack> {
|
||||
match self {
|
||||
Self::Federated { track, .. } => Some(track),
|
||||
Self::Local { .. } => None,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn primary_artist_key(&self) -> String {
|
||||
match self {
|
||||
Self::Local { track, .. } => track
|
||||
.artists
|
||||
.first()
|
||||
.map(|artist| music_dht::normalize_name(&artist.name))
|
||||
.unwrap_or_default(),
|
||||
Self::Federated { track, .. } => track
|
||||
.artist_names
|
||||
.first()
|
||||
.map(|artist| music_dht::normalize_name(artist))
|
||||
.unwrap_or_default(),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn content_key(&self) -> String {
|
||||
match self {
|
||||
Self::Local { track, .. } => track
|
||||
.content_id
|
||||
.clone()
|
||||
.unwrap_or_else(|| format!("local:{}", track.id)),
|
||||
Self::Federated { track, .. } => track
|
||||
.content_id
|
||||
.clone()
|
||||
.unwrap_or_else(|| format!("remote:{}:{}", track.owner, track.item_id)),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Default)]
|
||||
pub struct SearchState {
|
||||
pub query: String,
|
||||
@@ -1226,6 +1321,59 @@ pub struct SearchState {
|
||||
/// Present only for a track-seeded search; text-search refreshes must not
|
||||
/// replace this page with a title query.
|
||||
pub similarity_source: Option<i64>,
|
||||
/// The source is pinned at row zero; candidates below it are globally
|
||||
/// ranked across the local library and federation.
|
||||
pub similarity_source_track: Option<TrackItem>,
|
||||
pub similarity_tracks: Vec<SimilaritySearchHit>,
|
||||
pub similarity_stats: Option<crate::federation::SimilaritySearchStats>,
|
||||
pub similarity_error: Option<String>,
|
||||
}
|
||||
|
||||
impl SearchState {
|
||||
pub fn similarity_len(&self) -> usize {
|
||||
usize::from(self.similarity_source_track.is_some()) + self.similarity_tracks.len()
|
||||
}
|
||||
|
||||
pub fn similarity_track(&self, index: usize) -> Option<TrackItem> {
|
||||
if index == 0 {
|
||||
return self.similarity_source_track.clone();
|
||||
}
|
||||
self.similarity_tracks
|
||||
.get(index.checked_sub(1)?)
|
||||
.map(SimilaritySearchHit::track_item)
|
||||
}
|
||||
|
||||
pub fn similarity_fed_track(&self, index: usize) -> Option<&crate::federation::FedTrack> {
|
||||
self.similarity_tracks
|
||||
.get(index.checked_sub(1)?)?
|
||||
.federated_track()
|
||||
}
|
||||
|
||||
pub fn similarity_key(&self, index: usize) -> Option<String> {
|
||||
if index == 0 {
|
||||
return self
|
||||
.similarity_source_track
|
||||
.as_ref()
|
||||
.map(|track| format!("source:{}", track.id));
|
||||
}
|
||||
self.similarity_tracks
|
||||
.get(index.checked_sub(1)?)
|
||||
.map(SimilaritySearchHit::content_key)
|
||||
}
|
||||
|
||||
pub fn similarity_index_for_key(&self, key: &str) -> Option<usize> {
|
||||
if self
|
||||
.similarity_source_track
|
||||
.as_ref()
|
||||
.is_some_and(|track| key == format!("source:{}", track.id))
|
||||
{
|
||||
return Some(0);
|
||||
}
|
||||
self.similarity_tracks
|
||||
.iter()
|
||||
.position(|hit| hit.content_key() == key)
|
||||
.map(|index| index + 1)
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
|
||||
|
||||
+113
-29
@@ -905,6 +905,9 @@ fn set_track_scope_cursor(state: &mut AppState, scope: &TrackSelectionScope, val
|
||||
TrackSelectionScope::Queue => {
|
||||
state.queue_tab.cursor = value;
|
||||
}
|
||||
TrackSelectionScope::SimilaritySearch => {
|
||||
set_view_cursor(state, value);
|
||||
}
|
||||
TrackSelectionScope::FedSearch => {
|
||||
let base = state.search.results.as_ref().map_or(0, |r| r.len())
|
||||
+ state.search.fed_artists.len();
|
||||
@@ -952,6 +955,14 @@ fn current_track_list_context(state: &AppState) -> Option<(TrackSelectionScope,
|
||||
_ => None,
|
||||
},
|
||||
GlobalView::Search { cursor } => {
|
||||
if state.search.similarity_source.is_some() {
|
||||
let len = state.search.similarity_len();
|
||||
return (*cursor < len).then_some((
|
||||
TrackSelectionScope::SimilaritySearch,
|
||||
*cursor,
|
||||
len,
|
||||
));
|
||||
}
|
||||
// Only the federated tracks section is selectable here.
|
||||
let base = state.search.results.as_ref().map_or(0, |r| r.len())
|
||||
+ state.search.fed_artists.len();
|
||||
@@ -1044,6 +1055,20 @@ fn current_track_list(state: &AppState) -> Option<(TrackSelectionScope, usize, V
|
||||
}
|
||||
|
||||
pub fn selected_tracks(state: &AppState) -> Vec<TrackItem> {
|
||||
if state.active_tab == Tab::Global
|
||||
&& state.search.similarity_source.is_some()
|
||||
&& let Some(GlobalView::Search { cursor }) = state.global.stack.last()
|
||||
{
|
||||
let len = state.search.similarity_len();
|
||||
let indices = state
|
||||
.track_selection
|
||||
.indices(&TrackSelectionScope::SimilaritySearch, len)
|
||||
.unwrap_or_else(|| vec![(*cursor).min(len.saturating_sub(1))]);
|
||||
return indices
|
||||
.into_iter()
|
||||
.filter_map(|index| state.search.similarity_track(index))
|
||||
.collect();
|
||||
}
|
||||
// Federated contexts produce queueable placeholders that behave like
|
||||
// regular tracks (queue, info, playback-on-demand).
|
||||
{
|
||||
@@ -1276,6 +1301,9 @@ pub fn selected_track(state: &AppState) -> Option<TrackItem> {
|
||||
_ => None,
|
||||
},
|
||||
GlobalView::Search { cursor } => {
|
||||
if state.search.similarity_source.is_some() {
|
||||
return state.search.similarity_track(*cursor);
|
||||
}
|
||||
let results = state.search.results.as_ref()?;
|
||||
let offset = cursor.checked_sub(results.artists.len() + results.releases.len())?;
|
||||
match results.tracks.get(offset) {
|
||||
@@ -1877,10 +1905,13 @@ fn move_selection(state: &mut AppState, dx: isize, dy: isize) {
|
||||
refresh_track_selection_cursor(state);
|
||||
}
|
||||
Some(GlobalView::Search { cursor }) => {
|
||||
// Local results plus the federated section below them.
|
||||
let total = (state.search.results.as_ref().map_or(0, |r| r.len())
|
||||
+ state.search.fed_artists.len()
|
||||
+ state.search.fed_tracks.len()) as isize;
|
||||
let total = if state.search.similarity_source.is_some() {
|
||||
state.search.similarity_len()
|
||||
} else {
|
||||
state.search.results.as_ref().map_or(0, |r| r.len())
|
||||
+ state.search.fed_artists.len()
|
||||
+ state.search.fed_tracks.len()
|
||||
} as isize;
|
||||
if total == 0 {
|
||||
return;
|
||||
}
|
||||
@@ -2000,9 +2031,13 @@ fn current_view_len(state: &AppState) -> usize {
|
||||
_ => 0,
|
||||
},
|
||||
Some(GlobalView::Search { .. }) => {
|
||||
state.search.results.as_ref().map_or(0, |r| r.len())
|
||||
+ state.search.fed_artists.len()
|
||||
+ state.search.fed_tracks.len()
|
||||
if state.search.similarity_source.is_some() {
|
||||
state.search.similarity_len()
|
||||
} else {
|
||||
state.search.results.as_ref().map_or(0, |r| r.len())
|
||||
+ state.search.fed_artists.len()
|
||||
+ state.search.fed_tracks.len()
|
||||
}
|
||||
}
|
||||
Some(GlobalView::FedArtist { .. }) => fed_card_len(state),
|
||||
Some(GlobalView::FedRelease { index, .. }) => fed_card_release(state, *index)
|
||||
@@ -2292,31 +2327,47 @@ fn select_current(state: &mut AppState) -> Option<Effect> {
|
||||
},
|
||||
_ => Outcome::Nothing,
|
||||
},
|
||||
Some(GlobalView::Search { cursor }) => match &state.search.results {
|
||||
Some(results) => {
|
||||
let artists = results.artists.len();
|
||||
let releases = results.releases.len();
|
||||
if cursor < artists {
|
||||
Outcome::Push(GlobalView::Artist {
|
||||
id: results.artists[cursor].id,
|
||||
cursor: 0,
|
||||
})
|
||||
} else if cursor < artists + releases {
|
||||
Outcome::Push(GlobalView::Release {
|
||||
id: results.releases[cursor - artists].id,
|
||||
cursor: 0,
|
||||
})
|
||||
} else if results.tracks.get(cursor - artists - releases).is_some() {
|
||||
Outcome::Play {
|
||||
tracks: results.tracks.clone(),
|
||||
start: cursor - artists - releases,
|
||||
}
|
||||
Some(GlobalView::Search { cursor }) => {
|
||||
if state.search.similarity_source.is_some() {
|
||||
let tracks = (0..state.search.similarity_len())
|
||||
.filter_map(|index| state.search.similarity_track(index))
|
||||
.collect::<Vec<_>>();
|
||||
if tracks.is_empty() {
|
||||
Outcome::Nothing
|
||||
} else {
|
||||
fed_outcome(state, cursor - artists - releases - results.tracks.len())
|
||||
Outcome::Play {
|
||||
start: cursor.min(tracks.len() - 1),
|
||||
tracks,
|
||||
}
|
||||
}
|
||||
} else {
|
||||
match &state.search.results {
|
||||
Some(results) => {
|
||||
let artists = results.artists.len();
|
||||
let releases = results.releases.len();
|
||||
if cursor < artists {
|
||||
Outcome::Push(GlobalView::Artist {
|
||||
id: results.artists[cursor].id,
|
||||
cursor: 0,
|
||||
})
|
||||
} else if cursor < artists + releases {
|
||||
Outcome::Push(GlobalView::Release {
|
||||
id: results.releases[cursor - artists].id,
|
||||
cursor: 0,
|
||||
})
|
||||
} else if results.tracks.get(cursor - artists - releases).is_some() {
|
||||
Outcome::Play {
|
||||
tracks: results.tracks.clone(),
|
||||
start: cursor - artists - releases,
|
||||
}
|
||||
} else {
|
||||
fed_outcome(state, cursor - artists - releases - results.tracks.len())
|
||||
}
|
||||
}
|
||||
None => fed_outcome(state, cursor),
|
||||
}
|
||||
}
|
||||
None => fed_outcome(state, cursor),
|
||||
},
|
||||
}
|
||||
Some(GlobalView::FedArtist { cursor }) => match &state.fed_artist_view {
|
||||
Some((_, Loadable::Ready(card))) => {
|
||||
let release_indices = fed_artist_visible_release_indices(state, card);
|
||||
@@ -2492,6 +2543,20 @@ pub(crate) fn selected_fed_tracks(state: &AppState) -> Vec<crate::federation::Fe
|
||||
if state.active_tab != Tab::Global {
|
||||
return Vec::new();
|
||||
}
|
||||
if state.search.similarity_source.is_some()
|
||||
&& let Some(GlobalView::Search { cursor }) = state.global.stack.last()
|
||||
{
|
||||
let scope = TrackSelectionScope::SimilaritySearch;
|
||||
let len = state.search.similarity_len();
|
||||
let indices = state
|
||||
.track_selection
|
||||
.indices(&scope, len)
|
||||
.unwrap_or_else(|| vec![(*cursor).min(len.saturating_sub(1))]);
|
||||
return indices
|
||||
.into_iter()
|
||||
.filter_map(|index| state.search.similarity_fed_track(index).cloned())
|
||||
.collect();
|
||||
}
|
||||
// An active Shift-V range in a federated scope.
|
||||
if let Some(scope) = state.track_selection.scope.clone() {
|
||||
match scope {
|
||||
@@ -2796,6 +2861,25 @@ fn federation_select(state: &mut AppState) -> Option<Effect> {
|
||||
});
|
||||
None
|
||||
}
|
||||
SettingsRow::Similarity(SimilarityRow::MinimumScore) => {
|
||||
state.popup = Some(Popup::FedInput {
|
||||
field: FedInputField::SimilarityMinimumScore,
|
||||
input: crate::app::input::LineEdit::new(format!(
|
||||
"{:.2}",
|
||||
state.similarity.settings.minimum_score
|
||||
)),
|
||||
});
|
||||
None
|
||||
}
|
||||
SettingsRow::Similarity(SimilarityRow::MaxTracksPerArtist) => {
|
||||
state.popup = Some(Popup::FedInput {
|
||||
field: FedInputField::SimilarityMaxTracksPerArtist,
|
||||
input: crate::app::input::LineEdit::new(
|
||||
state.similarity.settings.max_tracks_per_artist.to_string(),
|
||||
),
|
||||
});
|
||||
None
|
||||
}
|
||||
SettingsRow::Similarity(SimilarityRow::Workers) => {
|
||||
state.popup = Some(Popup::FedInput {
|
||||
field: FedInputField::SimilarityWorkers,
|
||||
|
||||
+41
-2
@@ -45,7 +45,7 @@ pub struct LibraryFilters {
|
||||
pub source_mode: LibrarySourceMode,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct SimilaritySettings {
|
||||
/// Local embedding/search master switch. Network participation follows
|
||||
/// federation and additionally requires the explicit privacy consent.
|
||||
@@ -57,6 +57,14 @@ pub struct SimilaritySettings {
|
||||
pub profile: String,
|
||||
#[serde(default = "default_similarity_workers")]
|
||||
pub workers: usize,
|
||||
/// Requester-side cosine score floor. This is search policy, not part of
|
||||
/// the embedding profile, so changing it never invalidates vectors.
|
||||
#[serde(default = "default_similarity_minimum_score")]
|
||||
pub minimum_score: f32,
|
||||
/// Requester-side diversity cap applied independently to local and
|
||||
/// federated candidates.
|
||||
#[serde(default = "default_similarity_max_tracks_per_artist")]
|
||||
pub max_tracks_per_artist: usize,
|
||||
#[serde(default)]
|
||||
pub federation_consent: bool,
|
||||
/// Exact fingerprint of the last fully usable profile. Keeping this
|
||||
@@ -73,13 +81,15 @@ impl Default for SimilaritySettings {
|
||||
model: default_similarity_model(),
|
||||
profile: default_similarity_profile(),
|
||||
workers: default_similarity_workers(),
|
||||
minimum_score: default_similarity_minimum_score(),
|
||||
max_tracks_per_artist: default_similarity_max_tracks_per_artist(),
|
||||
federation_consent: false,
|
||||
active_profile: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct AppSettings {
|
||||
#[serde(default = "default_volume")]
|
||||
pub volume: u8,
|
||||
@@ -116,6 +126,11 @@ impl AppSettings {
|
||||
self.similarity.profile = default_similarity_profile();
|
||||
}
|
||||
self.similarity.workers = self.similarity.workers.clamp(1, 16);
|
||||
if !self.similarity.minimum_score.is_finite() {
|
||||
self.similarity.minimum_score = default_similarity_minimum_score();
|
||||
}
|
||||
self.similarity.minimum_score = self.similarity.minimum_score.clamp(0.0, 1.0);
|
||||
self.similarity.max_tracks_per_artist = self.similarity.max_tracks_per_artist.clamp(1, 50);
|
||||
self
|
||||
}
|
||||
}
|
||||
@@ -138,6 +153,14 @@ fn default_similarity_workers() -> usize {
|
||||
.unwrap_or(1)
|
||||
}
|
||||
|
||||
fn default_similarity_minimum_score() -> f32 {
|
||||
0.70
|
||||
}
|
||||
|
||||
fn default_similarity_max_tracks_per_artist() -> usize {
|
||||
5
|
||||
}
|
||||
|
||||
/// The historical permanent-download location, kept as the default for
|
||||
/// backward compatibility with existing installations.
|
||||
pub fn default_music_dir() -> PathBuf {
|
||||
@@ -204,5 +227,21 @@ hide_featured_only = true
|
||||
assert!(settings.library.hide_featured_only);
|
||||
assert_eq!(settings.library.source_mode, LibrarySourceMode::Global);
|
||||
assert_eq!(settings.music_dir, default_music_dir());
|
||||
assert_eq!(settings.similarity.minimum_score, 0.70);
|
||||
assert_eq!(settings.similarity.max_tracks_per_artist, 5);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn similarity_search_policy_is_normalized_without_changing_the_profile() {
|
||||
let mut settings = AppSettings::default();
|
||||
let profile = settings.similarity.profile.clone();
|
||||
settings.similarity.minimum_score = f32::NAN;
|
||||
settings.similarity.max_tracks_per_artist = 0;
|
||||
|
||||
let settings = settings.normalized();
|
||||
|
||||
assert_eq!(settings.similarity.minimum_score, 0.70);
|
||||
assert_eq!(settings.similarity.max_tracks_per_artist, 1);
|
||||
assert_eq!(settings.similarity.profile, profile);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -171,6 +171,7 @@ mod tests {
|
||||
"catalog",
|
||||
"audio",
|
||||
"similarity",
|
||||
"similarity_dht",
|
||||
"device_sync",
|
||||
"jam",
|
||||
] {
|
||||
|
||||
+149
-5
@@ -25,6 +25,8 @@ use std::sync::atomic::{AtomicI64, Ordering};
|
||||
use std::time::Duration;
|
||||
|
||||
use anyhow::{Context, Result};
|
||||
use music_dht::similarity_dht::SimilarityDht;
|
||||
use music_dht::similarity_lsh::SIMILARITY_DHT_ALPN;
|
||||
use music_dht::{
|
||||
ByteStream, ByteStreamConnectionStats, EndpointId, ItemKind, ItemSpec, LibraryItem,
|
||||
MusicDhtConfig, MusicDhtService, NetworkId, PeerTicket, PublishStats, RendezvousConfig,
|
||||
@@ -336,6 +338,27 @@ pub struct FedSearchResults {
|
||||
pub tracks: Vec<FedTrack>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct ScoredFedTrack {
|
||||
pub track: FedTrack,
|
||||
pub score: f32,
|
||||
pub embedding_signature: Option<[u8; music_dht::similarity::SIMILARITY_SIGNATURE_BYTES]>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Default, PartialEq, Eq)]
|
||||
pub struct SimilaritySearchStats {
|
||||
pub tracks: usize,
|
||||
pub artists: usize,
|
||||
pub peers_queried: usize,
|
||||
pub elapsed_ms: u64,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Default)]
|
||||
pub struct FedSimilaritySearchResults {
|
||||
pub tracks: Vec<ScoredFedTrack>,
|
||||
pub stats: SimilaritySearchStats,
|
||||
}
|
||||
|
||||
/// A track found through federated search.
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub struct FedTrack {
|
||||
@@ -413,6 +436,7 @@ pub struct NetworkLibrarySource {
|
||||
|
||||
struct Running {
|
||||
service: Arc<MusicDhtService>,
|
||||
similarity_dht: Arc<SimilarityDht>,
|
||||
network_name: String,
|
||||
network_id: NetworkId,
|
||||
tasks: Vec<tokio::task::JoinHandle<()>>,
|
||||
@@ -697,12 +721,15 @@ impl Federation {
|
||||
.stream_protocol(AUDIO_ALPN)
|
||||
// ...and browse each other's per-artist catalogs over this one.
|
||||
.stream_protocol(CATALOG_ALPN)
|
||||
// Anonymous, bounded direct embedding queries.
|
||||
.stream_protocol(SIMILARITY_ALPN)
|
||||
// Anonymous, bounded direct embedding queries have their own
|
||||
// versioned contract and survive catalog-schema upgrades.
|
||||
.schema_independent_stream_protocol(SIMILARITY_ALPN)
|
||||
// Personal-device sync (likes, playlists, trusted devices).
|
||||
.stream_protocol(crate::devices::SYNC_ALPN)
|
||||
// Capability-scoped shared playback control.
|
||||
.stream_protocol(crate::jam::JAM_ALPN)
|
||||
// Signed LSH summaries form their own upgrade-safe DHT overlay.
|
||||
.schema_independent_stream_protocol(SIMILARITY_DHT_ALPN)
|
||||
// Informational application/protocol versions.
|
||||
.schema_independent_stream_protocol(capabilities::CAPABILITIES_ALPN)
|
||||
.build()
|
||||
@@ -717,6 +744,25 @@ impl Federation {
|
||||
"federation started"
|
||||
);
|
||||
|
||||
let similarity_dht = SimilarityDht::open(
|
||||
Arc::clone(&service),
|
||||
self.data_dir.join("similarity-routing.sqlite3"),
|
||||
)
|
||||
.await
|
||||
.map_err(|err| anyhow::anyhow!("failed to start the similarity DHT: {err}"))?;
|
||||
let similarity_dht_acceptor = service
|
||||
.stream_acceptor(SIMILARITY_DHT_ALPN)
|
||||
.map_err(|err| anyhow::anyhow!("failed to take similarity DHT acceptor: {err}"))?;
|
||||
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_dht_sync_task = tokio::spawn(similarity_route_sync_loop(
|
||||
Arc::clone(&similarity_dht),
|
||||
Arc::clone(&self.similarity),
|
||||
Arc::clone(&self.library),
|
||||
));
|
||||
|
||||
// Drain DHT events into the log; the channel is bounded.
|
||||
let event_task = tokio::spawn(async move {
|
||||
while let Some(event) = events.recv().await {
|
||||
@@ -797,6 +843,7 @@ impl Federation {
|
||||
|
||||
*guard = Some(Running {
|
||||
service,
|
||||
similarity_dht,
|
||||
network_name,
|
||||
network_id,
|
||||
tasks: vec![
|
||||
@@ -811,6 +858,9 @@ impl Federation {
|
||||
jam_poll_task,
|
||||
capabilities_serve_task,
|
||||
capabilities_probe_task,
|
||||
similarity_dht_serve_task,
|
||||
similarity_dht_maintenance_task,
|
||||
similarity_dht_sync_task,
|
||||
],
|
||||
});
|
||||
self.set_error(None);
|
||||
@@ -835,6 +885,20 @@ impl Federation {
|
||||
.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")
|
||||
}
|
||||
|
||||
fn ensure_connected_devices_enabled(&self) -> Result<()> {
|
||||
let settings = self.settings();
|
||||
anyhow::ensure!(
|
||||
@@ -1170,13 +1234,22 @@ impl Federation {
|
||||
&self,
|
||||
query: crate::similarity::QueryVector,
|
||||
limit: usize,
|
||||
) -> Result<FedSearchResults> {
|
||||
) -> Result<FedSimilaritySearchResults> {
|
||||
anyhow::ensure!(
|
||||
self.similarity.network_allowed(),
|
||||
"similarity federation has no consent"
|
||||
);
|
||||
let service = self.service().await?;
|
||||
similarity::search(service, query, limit, Arc::clone(&self.transport_stats)).await
|
||||
let (service, similarity_dht) = self.similarity_services().await?;
|
||||
let settings = self.similarity.settings();
|
||||
similarity::search(
|
||||
service,
|
||||
similarity_dht,
|
||||
query,
|
||||
limit,
|
||||
settings.minimum_score,
|
||||
Arc::clone(&self.transport_stats),
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
/// Resolves a share-link content id to one playable federated track.
|
||||
@@ -2565,6 +2638,77 @@ fn sort_fed_appearances(appearances: &mut [FedAppearsOn]) {
|
||||
});
|
||||
}
|
||||
|
||||
async fn similarity_route_sync_loop(
|
||||
routing: Arc<SimilarityDht>,
|
||||
similarity: Arc<crate::similarity::Manager>,
|
||||
library: Arc<Library>,
|
||||
) {
|
||||
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 !similarity.network_allowed() {
|
||||
if published_marker.take().is_some() {
|
||||
routing.clear_local_signatures();
|
||||
tracing::info!("local similarity DHT publication disabled");
|
||||
}
|
||||
continue;
|
||||
}
|
||||
let status = similarity.status();
|
||||
let Some(profile_id) = status.active_profile else {
|
||||
continue;
|
||||
};
|
||||
if status.phase != crate::similarity::Phase::Ready {
|
||||
continue;
|
||||
}
|
||||
let library = Arc::clone(&library);
|
||||
let profile_for_task = profile_id.clone();
|
||||
let loaded = tokio::task::spawn_blocking(move || {
|
||||
let signatures = library.similarity_routing_signatures(&profile_for_task)?;
|
||||
let mut hasher = blake3::Hasher::new();
|
||||
for signature in &signatures {
|
||||
hasher.update(signature);
|
||||
}
|
||||
Ok::<_, anyhow::Error>((signatures, hasher.finalize()))
|
||||
})
|
||||
.await;
|
||||
let (signatures, fingerprint) = match loaded {
|
||||
Ok(Ok(loaded)) => loaded,
|
||||
Ok(Err(error)) => {
|
||||
tracing::warn!(%error, %profile_id, "similarity routing signatures unavailable");
|
||||
continue;
|
||||
}
|
||||
Err(error) => {
|
||||
tracing::warn!(%error, "similarity routing signature task failed");
|
||||
continue;
|
||||
}
|
||||
};
|
||||
let marker = (profile_id.clone(), fingerprint);
|
||||
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 stop_running(running: Option<Running>) {
|
||||
let Some(running) = running else { return };
|
||||
for task in &running.tasks {
|
||||
|
||||
+146
-56
@@ -4,24 +4,30 @@
|
||||
//! module owns application policy: consent, peer fan-out, local index access,
|
||||
//! result conversion, deduplication, and ranking limits.
|
||||
|
||||
use std::collections::{HashMap, HashSet};
|
||||
use std::collections::HashSet;
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
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::federation::{FedSearchResults, FedTrack, TransportStats};
|
||||
use crate::federation::{
|
||||
FedSimilaritySearchResults, FedTrack, ScoredFedTrack, SimilaritySearchStats, TransportStats,
|
||||
};
|
||||
use crate::similarity::{Manager, QueryVector};
|
||||
|
||||
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 MAX_PER_ARTIST: usize = 3;
|
||||
const ROUTING_TIMEOUT: Duration = Duration::from_secs(5);
|
||||
const MAX_NEAR_DUPLICATE_SIGNATURE_DISTANCE: u32 = 8;
|
||||
|
||||
pub async fn serve_peers(
|
||||
@@ -57,7 +63,7 @@ async fn serve_one(
|
||||
let vector = request.vector;
|
||||
let limit = request.limit;
|
||||
let matches = tokio::task::spawn_blocking(move || {
|
||||
similarity.search_vector(&profile, &vector, None, None, limit)
|
||||
similarity.search_vector_for_peer(&profile, &vector, limit)
|
||||
})
|
||||
.await
|
||||
.context("local similarity task failed")
|
||||
@@ -122,20 +128,51 @@ async fn serve_one(
|
||||
|
||||
pub async fn search(
|
||||
service: Arc<MusicDhtService>,
|
||||
routing: Arc<SimilarityDht>,
|
||||
query: QueryVector,
|
||||
limit: usize,
|
||||
minimum_score: f32,
|
||||
transport: Arc<TransportStats>,
|
||||
) -> Result<FedSearchResults> {
|
||||
) -> Result<FedSimilaritySearchResults> {
|
||||
let started = Instant::now();
|
||||
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<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
|
||||
.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;
|
||||
@@ -148,36 +185,50 @@ 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::<Vec<_>>()
|
||||
.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;
|
||||
let mut peers_queried = initial;
|
||||
for response in responses {
|
||||
match response {
|
||||
Ok(peer_hits) => hits.extend(peer_hits),
|
||||
Ok(peer_hits) => {
|
||||
successful += 1;
|
||||
hits.extend(peer_hits);
|
||||
}
|
||||
Err(err) => tracing::debug!(%err, "similarity peer query skipped"),
|
||||
}
|
||||
}
|
||||
if initial < peers.len() && (hits.len() < limit || successful < initial.min(4)) {
|
||||
peers_queried = peers.len();
|
||||
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(err) => tracing::debug!(%err, "fallback similarity peer query skipped"),
|
||||
}
|
||||
}
|
||||
}
|
||||
hits.sort_by(|left, right| right.1.total_cmp(&left.1));
|
||||
let mut dedup = HashSet::new();
|
||||
let mut embedding_signatures = vec![query_signature];
|
||||
let mut artist_counts: HashMap<String, usize> = HashMap::new();
|
||||
let mut tracks = Vec::new();
|
||||
for (track, _, embedding_signature) in hits {
|
||||
for (track, score, embedding_signature) in hits {
|
||||
if score < minimum_score {
|
||||
break;
|
||||
}
|
||||
if query
|
||||
.source_content_id
|
||||
.as_deref()
|
||||
@@ -200,46 +251,85 @@ pub async fn search(
|
||||
}) {
|
||||
continue;
|
||||
}
|
||||
let artist = track
|
||||
.artist_names
|
||||
.first()
|
||||
.map(|name| music_dht::normalize_name(name))
|
||||
.unwrap_or_default();
|
||||
let count = artist_counts.entry(artist.clone()).or_default();
|
||||
if !artist.is_empty() && *count >= MAX_PER_ARTIST {
|
||||
continue;
|
||||
}
|
||||
*count += 1;
|
||||
if let Some(signature) = embedding_signature {
|
||||
embedding_signatures.push(signature);
|
||||
}
|
||||
tracks.push(track);
|
||||
tracks.push(ScoredFedTrack {
|
||||
track,
|
||||
score,
|
||||
embedding_signature,
|
||||
});
|
||||
if tracks.len() >= limit.min(wire::MAX_SIMILARITY_RESULTS) {
|
||||
break;
|
||||
}
|
||||
}
|
||||
Ok(FedSearchResults {
|
||||
artists: Vec::new(),
|
||||
let artists = tracks
|
||||
.iter()
|
||||
.filter_map(|hit| hit.track.artist_names.first())
|
||||
.map(|name| music_dht::normalize_name(name))
|
||||
.filter(|name| !name.is_empty())
|
||||
.collect::<HashSet<_>>()
|
||||
.len();
|
||||
let elapsed_ms = started.elapsed().as_millis().min(u128::from(u64::MAX)) as u64;
|
||||
Ok(FedSimilaritySearchResults {
|
||||
stats: SimilaritySearchStats {
|
||||
tracks: tracks.len(),
|
||||
artists,
|
||||
peers_queried,
|
||||
elapsed_ms,
|
||||
},
|
||||
tracks,
|
||||
})
|
||||
}
|
||||
|
||||
type PeerHits = Vec<(
|
||||
FedTrack,
|
||||
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(
|
||||
service: Arc<MusicDhtService>,
|
||||
owner: EndpointId,
|
||||
peer: QueryPeer,
|
||||
request: &SimilarityRequest,
|
||||
transport: Arc<TransportStats>,
|
||||
) -> Result<
|
||||
Vec<(
|
||||
FedTrack,
|
||||
f32,
|
||||
Option<[u8; wire::SIMILARITY_SIGNATURE_BYTES]>,
|
||||
)>,
|
||||
> {
|
||||
let mut stream = service
|
||||
.open_stream(owner, SIMILARITY_ALPN)
|
||||
.await
|
||||
.map_err(|err| anyhow::anyhow!("cannot reach similarity peer: {err}"))?;
|
||||
) -> Result<PeerHits> {
|
||||
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(|err| anyhow::anyhow!("cannot reach similarity peer: {err}"))?;
|
||||
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);
|
||||
|
||||
+75
-2
@@ -207,6 +207,7 @@ CREATE TABLE IF NOT EXISTS track_embeddings (
|
||||
profile_id TEXT NOT NULL REFERENCES similarity_profiles(profile_id) ON DELETE CASCADE,
|
||||
dimensions INTEGER NOT NULL,
|
||||
vector BLOB NOT NULL,
|
||||
routing_signature BLOB,
|
||||
source_content_id TEXT,
|
||||
computed_at_ms INTEGER NOT NULL,
|
||||
PRIMARY KEY (track_id, profile_id)
|
||||
@@ -1551,15 +1552,17 @@ impl Library {
|
||||
"embedding contains a non-finite value"
|
||||
);
|
||||
let bytes = embedding_to_bytes(vector);
|
||||
let routing_signature = music_dht::similarity_lsh::routing_signature(vector)?;
|
||||
let conn = self.lock();
|
||||
conn.execute(
|
||||
"INSERT INTO track_embeddings (
|
||||
track_id, profile_id, dimensions, vector,
|
||||
track_id, profile_id, dimensions, vector, routing_signature,
|
||||
source_content_id, computed_at_ms
|
||||
) VALUES (?1, ?2, ?3, ?4, ?5, ?6)
|
||||
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
|
||||
ON CONFLICT(track_id, profile_id) DO UPDATE SET
|
||||
dimensions = excluded.dimensions,
|
||||
vector = excluded.vector,
|
||||
routing_signature = excluded.routing_signature,
|
||||
source_content_id = excluded.source_content_id,
|
||||
computed_at_ms = excluded.computed_at_ms",
|
||||
params![
|
||||
@@ -1567,6 +1570,7 @@ impl Library {
|
||||
profile_id,
|
||||
vector.len() as i64,
|
||||
bytes,
|
||||
routing_signature.as_slice(),
|
||||
track.content_id.as_deref(),
|
||||
now_ms_i64(),
|
||||
],
|
||||
@@ -1633,6 +1637,65 @@ impl Library {
|
||||
Ok(embeddings)
|
||||
}
|
||||
|
||||
/// Loads the compact DHT-routing signatures for every current local
|
||||
/// embedding. Rows created before similarity routing existed are
|
||||
/// backfilled in place from their durable vectors.
|
||||
pub fn similarity_routing_signatures(&self, profile_id: &str) -> Result<Vec<[u8; 32]>> {
|
||||
let mut conn = self.lock();
|
||||
let transaction = conn.transaction()?;
|
||||
let missing = {
|
||||
let mut statement = transaction.prepare(
|
||||
"SELECT e.track_id, e.dimensions, e.vector
|
||||
FROM track_embeddings e
|
||||
JOIN tracks t ON t.id = e.track_id
|
||||
WHERE e.profile_id = ?1
|
||||
AND e.source_content_id IS t.content_id
|
||||
AND (e.routing_signature IS NULL OR length(e.routing_signature) != 32)
|
||||
ORDER BY e.track_id",
|
||||
)?;
|
||||
statement
|
||||
.query_map([profile_id], |row| {
|
||||
Ok((
|
||||
row.get::<_, i64>(0)?,
|
||||
row.get::<_, i64>(1)?,
|
||||
row.get::<_, Vec<u8>>(2)?,
|
||||
))
|
||||
})?
|
||||
.collect::<rusqlite::Result<Vec<_>>>()?
|
||||
};
|
||||
for (track_id, dimensions, vector_bytes) in missing {
|
||||
let vector = embedding_from_bytes(dimensions, &vector_bytes)?;
|
||||
let signature = music_dht::similarity_lsh::routing_signature(&vector)?;
|
||||
transaction.execute(
|
||||
"UPDATE track_embeddings
|
||||
SET routing_signature = ?3
|
||||
WHERE track_id = ?1 AND profile_id = ?2",
|
||||
params![track_id, profile_id, signature.as_slice()],
|
||||
)?;
|
||||
}
|
||||
transaction.commit()?;
|
||||
|
||||
let mut statement = conn.prepare(
|
||||
"SELECT e.routing_signature
|
||||
FROM track_embeddings e
|
||||
JOIN tracks t ON t.id = e.track_id
|
||||
WHERE e.profile_id = ?1
|
||||
AND e.source_content_id IS t.content_id
|
||||
ORDER BY e.track_id",
|
||||
)?;
|
||||
let stored = statement
|
||||
.query_map([profile_id], |row| row.get::<_, Vec<u8>>(0))?
|
||||
.collect::<rusqlite::Result<Vec<_>>>()?;
|
||||
let signatures = stored
|
||||
.into_iter()
|
||||
.map(|signature| {
|
||||
<[u8; 32]>::try_from(signature)
|
||||
.map_err(|_| anyhow::anyhow!("invalid similarity routing signature length"))
|
||||
})
|
||||
.collect::<Result<Vec<_>>>()?;
|
||||
Ok(signatures)
|
||||
}
|
||||
|
||||
pub fn similarity_storage_stats(&self, profile_id: &str) -> Result<SimilarityStorageStats> {
|
||||
let conn = self.lock();
|
||||
let total_tracks = conn.query_row("SELECT COUNT(*) FROM tracks", [], |row| {
|
||||
@@ -3444,6 +3507,16 @@ fn ensure_schema_migrations(conn: &Connection) -> Result<()> {
|
||||
[],
|
||||
)?;
|
||||
}
|
||||
let embedding_columns = table_columns(conn, "track_embeddings")?;
|
||||
if !embedding_columns
|
||||
.iter()
|
||||
.any(|column| column == "routing_signature")
|
||||
{
|
||||
conn.execute(
|
||||
"ALTER TABLE track_embeddings ADD COLUMN routing_signature BLOB",
|
||||
[],
|
||||
)?;
|
||||
}
|
||||
let mut rows = conn.prepare("SELECT id, title FROM playlists WHERE sync_id IS NULL")?;
|
||||
let missing = rows
|
||||
.query_map([], |row| {
|
||||
|
||||
+27
-4
@@ -701,24 +701,47 @@ fn similarity_embeddings_round_trip_and_keep_profiles_separate() {
|
||||
.into_iter()
|
||||
.find(|track| track.id == track_id)
|
||||
.unwrap();
|
||||
lib.store_similarity_embedding(&track, "profile-a", &[0.1, 0.2, 0.3])
|
||||
let first = [0.26726124, 0.5345225, 0.8017837];
|
||||
let second = [0.8017837, 0.5345225, 0.26726124];
|
||||
lib.store_similarity_embedding(&track, "profile-a", &first)
|
||||
.unwrap();
|
||||
lib.store_similarity_embedding(&track, "profile-b", &[0.3, 0.2, 0.1])
|
||||
lib.store_similarity_embedding(&track, "profile-b", &second)
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(
|
||||
lib.similarity_embedding(track_id, "profile-a").unwrap(),
|
||||
Some(vec![0.1, 0.2, 0.3])
|
||||
Some(first.to_vec())
|
||||
);
|
||||
assert_eq!(
|
||||
lib.similarity_embedding(track_id, "profile-b").unwrap(),
|
||||
Some(vec![0.3, 0.2, 0.1])
|
||||
Some(second.to_vec())
|
||||
);
|
||||
let stats = lib.similarity_storage_stats("profile-a").unwrap();
|
||||
assert_eq!(stats.total_tracks, 1);
|
||||
assert_eq!(stats.embedded_tracks, 1);
|
||||
assert_eq!(stats.stored_vectors, 2);
|
||||
assert_eq!(stats.stored_bytes, 24);
|
||||
lib.lock()
|
||||
.execute(
|
||||
"UPDATE track_embeddings SET routing_signature = NULL WHERE profile_id = 'profile-a'",
|
||||
[],
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
lib.similarity_routing_signatures("profile-a")
|
||||
.unwrap()
|
||||
.len(),
|
||||
1
|
||||
);
|
||||
let stored_signature_bytes: i64 = lib
|
||||
.lock()
|
||||
.query_row(
|
||||
"SELECT length(routing_signature) FROM track_embeddings WHERE profile_id = 'profile-a'",
|
||||
[],
|
||||
|row| row.get(0),
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(stored_signature_bytes, 32);
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
+180
-38
@@ -7,7 +7,7 @@
|
||||
use std::collections::{HashMap, VecDeque};
|
||||
use std::fs::File;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
|
||||
use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
|
||||
use std::sync::{Arc, Mutex, RwLock};
|
||||
use std::time::Instant;
|
||||
|
||||
@@ -36,7 +36,7 @@ const EMBEDDING_DIMENSIONS: usize = 1280;
|
||||
const MODEL_BATCH: usize = 8;
|
||||
const MAX_MODEL_BYTES: usize = 64 * 1024 * 1024;
|
||||
const RESULT_LIMIT: usize = 50;
|
||||
const MAX_PER_ARTIST: usize = 3;
|
||||
const PEER_CANDIDATE_MAX_PER_ARTIST: usize = 10;
|
||||
const NEAR_DUPLICATE_COSINE: f32 = 0.995;
|
||||
const FULL_TRACK_MAX_SECONDS: u32 = 5 * 60;
|
||||
const LONG_TRACK_WINDOW_SECONDS: u32 = 60;
|
||||
@@ -101,7 +101,7 @@ impl Phase {
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Default)]
|
||||
#[derive(Debug, Clone, Default, PartialEq, Eq)]
|
||||
pub struct SimilarityStatus {
|
||||
pub phase: Phase,
|
||||
pub active_profile: Option<String>,
|
||||
@@ -144,6 +144,8 @@ pub struct Manager {
|
||||
settings: Mutex<SimilaritySettings>,
|
||||
workers: AtomicUsize,
|
||||
generation: AtomicU64,
|
||||
pipeline_running: AtomicBool,
|
||||
rescan_requested: AtomicBool,
|
||||
status: Mutex<SimilarityStatus>,
|
||||
index: RwLock<Index>,
|
||||
model: Mutex<Option<(String, RunnableModel)>>,
|
||||
@@ -169,13 +171,20 @@ impl Manager {
|
||||
Err(err) => tracing::warn!(%err, "similarity index restore failed"),
|
||||
}
|
||||
}
|
||||
let target_profile = model_by_id(&settings.model)
|
||||
.filter(|_| profile_by_id(&settings.profile).is_some())
|
||||
.map(|model| profile_fingerprint(model, &settings.profile));
|
||||
let restored_profile_is_current = index.profile_id == target_profile;
|
||||
let status = SimilarityStatus {
|
||||
phase: if settings.enabled {
|
||||
Phase::Loading
|
||||
} else {
|
||||
phase: if !settings.enabled {
|
||||
Phase::Disabled
|
||||
} else if restored_profile_is_current {
|
||||
Phase::Ready
|
||||
} else {
|
||||
Phase::Loading
|
||||
},
|
||||
active_profile: index.profile_id.clone(),
|
||||
target_profile,
|
||||
model: settings.model.clone(),
|
||||
..SimilarityStatus::default()
|
||||
};
|
||||
@@ -184,6 +193,8 @@ impl Manager {
|
||||
event_tx,
|
||||
workers: AtomicUsize::new(settings.workers.clamp(1, 16)),
|
||||
generation: AtomicU64::new(0),
|
||||
pipeline_running: AtomicBool::new(false),
|
||||
rescan_requested: AtomicBool::new(false),
|
||||
settings: Mutex::new(settings),
|
||||
status: Mutex::new(status),
|
||||
index: RwLock::new(index),
|
||||
@@ -223,23 +234,51 @@ impl Manager {
|
||||
|| previous.model != settings.model
|
||||
|| previous.profile != settings.profile
|
||||
{
|
||||
self.generation.fetch_add(1, Ordering::AcqRel);
|
||||
self.start();
|
||||
}
|
||||
}
|
||||
|
||||
/// Requests a scan without cancelling useful work already in progress.
|
||||
/// Bursts of library-change notifications collapse into one follow-up
|
||||
/// pass, so metadata refreshes cannot repeatedly restart the model.
|
||||
pub fn start(self: &Arc<Self>) {
|
||||
let generation = self.generation.fetch_add(1, Ordering::AcqRel) + 1;
|
||||
self.rescan_requested.store(true, Ordering::Release);
|
||||
if self.pipeline_running.swap(true, Ordering::AcqRel) {
|
||||
return;
|
||||
}
|
||||
let this = Arc::clone(self);
|
||||
tokio::spawn(async move {
|
||||
if let Err(err) = this.run_pipeline(generation).await
|
||||
&& this.generation.load(Ordering::Acquire) == generation
|
||||
{
|
||||
tracing::error!(%err, "similarity pipeline failed");
|
||||
this.update_status(|status| {
|
||||
status.phase = Phase::Error;
|
||||
status.current_track = None;
|
||||
status.last_error = Some(format!("{err:#}"));
|
||||
});
|
||||
loop {
|
||||
// This pass covers every notification received before it
|
||||
// starts. A notification during the pass requests one more.
|
||||
this.rescan_requested.store(false, Ordering::Release);
|
||||
let generation = this.generation.load(Ordering::Acquire);
|
||||
if let Err(err) = this.run_pipeline(generation).await
|
||||
&& this.generation.load(Ordering::Acquire) == generation
|
||||
{
|
||||
tracing::error!(%err, "similarity pipeline failed");
|
||||
this.update_status(|status| {
|
||||
status.phase = Phase::Error;
|
||||
status.current_track = None;
|
||||
status.last_error = Some(format!("{err:#}"));
|
||||
});
|
||||
}
|
||||
|
||||
if this.rescan_requested.load(Ordering::Acquire) {
|
||||
continue;
|
||||
}
|
||||
|
||||
this.pipeline_running.store(false, Ordering::Release);
|
||||
// Close the small race between checking the request flag and
|
||||
// releasing ownership of the worker. If another worker has
|
||||
// already claimed it, that worker owns the pending pass.
|
||||
if this.rescan_requested.swap(false, Ordering::AcqRel)
|
||||
&& !this.pipeline_running.swap(true, Ordering::AcqRel)
|
||||
{
|
||||
continue;
|
||||
}
|
||||
break;
|
||||
}
|
||||
});
|
||||
}
|
||||
@@ -329,6 +368,52 @@ impl Manager {
|
||||
exclude_track_id: Option<i64>,
|
||||
exclude_content_id: Option<&str>,
|
||||
limit: usize,
|
||||
) -> Result<Vec<SimilarTrack>> {
|
||||
let settings = lock(&self.settings);
|
||||
let minimum_score = settings.minimum_score;
|
||||
let max_tracks_per_artist = settings.max_tracks_per_artist;
|
||||
drop(settings);
|
||||
self.search_vector_with_policy(
|
||||
profile_id,
|
||||
vector,
|
||||
exclude_track_id,
|
||||
exclude_content_id,
|
||||
limit,
|
||||
minimum_score,
|
||||
max_tracks_per_artist,
|
||||
)
|
||||
}
|
||||
|
||||
/// Returns a wider, policy-neutral candidate set to a remote requester.
|
||||
/// The requester applies its own score threshold and artist diversity
|
||||
/// limit; neither value is part of embedding compatibility.
|
||||
pub(crate) fn search_vector_for_peer(
|
||||
&self,
|
||||
profile_id: &str,
|
||||
vector: &[f32],
|
||||
limit: usize,
|
||||
) -> Result<Vec<SimilarTrack>> {
|
||||
self.search_vector_with_policy(
|
||||
profile_id,
|
||||
vector,
|
||||
None,
|
||||
None,
|
||||
limit,
|
||||
-1.0,
|
||||
PEER_CANDIDATE_MAX_PER_ARTIST,
|
||||
)
|
||||
}
|
||||
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
fn search_vector_with_policy(
|
||||
&self,
|
||||
profile_id: &str,
|
||||
vector: &[f32],
|
||||
exclude_track_id: Option<i64>,
|
||||
exclude_content_id: Option<&str>,
|
||||
limit: usize,
|
||||
minimum_score: f32,
|
||||
max_tracks_per_artist: usize,
|
||||
) -> Result<Vec<SimilarTrack>> {
|
||||
anyhow::ensure!(
|
||||
!vector.is_empty() && vector.len() <= 4096,
|
||||
@@ -357,7 +442,7 @@ impl Manager {
|
||||
entry.vector.as_slice(),
|
||||
)
|
||||
})
|
||||
.filter(|(_, score, _, _)| score.is_finite())
|
||||
.filter(|(_, score, _, _)| score.is_finite() && *score >= minimum_score)
|
||||
.collect();
|
||||
scores.sort_by(|left, right| right.1.total_cmp(&left.1));
|
||||
|
||||
@@ -371,7 +456,7 @@ impl Manager {
|
||||
continue;
|
||||
}
|
||||
let count = artist_counts.entry(artist.to_string()).or_default();
|
||||
if !artist.is_empty() && *count >= MAX_PER_ARTIST {
|
||||
if !artist.is_empty() && *count >= max_tracks_per_artist {
|
||||
continue;
|
||||
}
|
||||
*count += 1;
|
||||
@@ -431,7 +516,6 @@ impl Manager {
|
||||
)?;
|
||||
let stats = self.library.similarity_storage_stats(&profile_id)?;
|
||||
self.update_status(|status| {
|
||||
status.phase = Phase::Downloading;
|
||||
status.target_profile = Some(profile_id.clone());
|
||||
status.model = spec.id.to_string();
|
||||
status.total_tracks = stats.total_tracks;
|
||||
@@ -442,21 +526,18 @@ impl Manager {
|
||||
status.current_track = None;
|
||||
status.last_error = None;
|
||||
});
|
||||
let model_path = self.ensure_model(spec, generation).await?;
|
||||
self.ensure_generation(generation)?;
|
||||
self.update_status(|status| status.phase = Phase::Loading);
|
||||
let model = self.load_model(&profile_id, &model_path).await?;
|
||||
self.ensure_generation(generation)?;
|
||||
|
||||
let mut pending: VecDeque<_> = self.library.pending_similarity_tracks(&profile_id)?.into();
|
||||
let pending_total = pending.len();
|
||||
self.update_status(|status| {
|
||||
status.phase = if pending_total == 0 {
|
||||
Phase::Loading
|
||||
} else {
|
||||
Phase::Processing
|
||||
};
|
||||
});
|
||||
if pending.is_empty() {
|
||||
self.ensure_generation(generation)?;
|
||||
return self.activate_profile(profile_id);
|
||||
}
|
||||
|
||||
let model_path = self.ensure_model(spec, generation).await?;
|
||||
self.ensure_generation(generation)?;
|
||||
let model = self.load_model(&profile_id, &model_path).await?;
|
||||
self.ensure_generation(generation)?;
|
||||
self.update_status(|status| status.phase = Phase::Processing);
|
||||
|
||||
let mut jobs = tokio::task::JoinSet::new();
|
||||
while !pending.is_empty() || !jobs.is_empty() {
|
||||
@@ -506,6 +587,10 @@ impl Manager {
|
||||
}
|
||||
}
|
||||
self.ensure_generation(generation)?;
|
||||
self.activate_profile(profile_id)
|
||||
}
|
||||
|
||||
fn activate_profile(&self, profile_id: String) -> Result<()> {
|
||||
let entries = self.library.load_similarity_index(&profile_id)?;
|
||||
let total_tracks = self
|
||||
.library
|
||||
@@ -519,7 +604,12 @@ impl Manager {
|
||||
profile_id: Some(profile_id.clone()),
|
||||
entries,
|
||||
};
|
||||
lock(&self.settings).active_profile = Some(profile_id.clone());
|
||||
let profile_changed = {
|
||||
let mut settings = lock(&self.settings);
|
||||
let changed = settings.active_profile.as_deref() != Some(&profile_id);
|
||||
settings.active_profile = Some(profile_id.clone());
|
||||
changed
|
||||
};
|
||||
let stats = self.library.similarity_storage_stats(&profile_id)?;
|
||||
self.update_status(|status| {
|
||||
status.phase = Phase::Ready;
|
||||
@@ -531,9 +621,11 @@ impl Manager {
|
||||
status.stored_bytes = stats.stored_bytes;
|
||||
status.current_track = None;
|
||||
});
|
||||
let _ = self
|
||||
.event_tx
|
||||
.send(AppEvent::SimilarityProfileActivated(Some(profile_id)));
|
||||
if profile_changed {
|
||||
let _ = self
|
||||
.event_tx
|
||||
.send(AppEvent::SimilarityProfileActivated(Some(profile_id)));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -561,6 +653,7 @@ impl Manager {
|
||||
tokio::fs::remove_file(&path).await?;
|
||||
}
|
||||
|
||||
self.update_status(|status| status.phase = Phase::Downloading);
|
||||
let response = reqwest::get(spec.url).await?.error_for_status()?;
|
||||
let tmp = path.with_extension(format!("part-{}-{generation}", std::process::id()));
|
||||
let mut file = tokio::fs::File::create(&tmp).await?;
|
||||
@@ -603,6 +696,7 @@ impl Manager {
|
||||
{
|
||||
return Ok(Arc::clone(model));
|
||||
}
|
||||
self.update_status(|status| status.phase = Phase::Loading);
|
||||
let path = path.to_path_buf();
|
||||
let model = tokio::task::spawn_blocking(move || load_onnx(&path))
|
||||
.await
|
||||
@@ -614,10 +708,13 @@ impl Manager {
|
||||
fn update_status(&self, update: impl FnOnce(&mut SimilarityStatus)) {
|
||||
let snapshot = {
|
||||
let mut status = lock(&self.status);
|
||||
let previous = status.clone();
|
||||
update(&mut status);
|
||||
status.clone()
|
||||
(*status != previous).then(|| status.clone())
|
||||
};
|
||||
let _ = self.event_tx.send(AppEvent::SimilarityStatus(snapshot));
|
||||
if let Some(snapshot) = snapshot {
|
||||
let _ = self.event_tx.send(AppEvent::SimilarityStatus(snapshot));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -977,6 +1074,14 @@ fn write<T>(lock: &RwLock<T>) -> std::sync::RwLockWriteGuard<'_, T> {
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn unique_test_dir(label: &str) -> PathBuf {
|
||||
let unique = std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.unwrap()
|
||||
.as_nanos();
|
||||
std::env::temp_dir().join(format!("furumi-{label}-{}-{unique}", std::process::id()))
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn profile_fingerprint_changes_with_contract() {
|
||||
let model = &MODELS[0];
|
||||
@@ -1010,6 +1115,43 @@ mod tests {
|
||||
assert!(!is_near_duplicate(&distinct, &[&query]));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn repeated_rescans_keep_an_up_to_date_profile_ready() {
|
||||
let directory = unique_test_dir("similarity-stable-status");
|
||||
let library = Arc::new(Library::open(&directory.join("library.db")).unwrap());
|
||||
let profile_id = profile_fingerprint(&MODELS[0], DEFAULT_PROFILE_ID);
|
||||
let settings = SimilaritySettings {
|
||||
enabled: true,
|
||||
active_profile: Some(profile_id),
|
||||
..SimilaritySettings::default()
|
||||
};
|
||||
let (event_tx, mut event_rx) = tokio::sync::mpsc::unbounded_channel();
|
||||
let manager = Manager::new(Arc::clone(&library), event_tx, settings);
|
||||
|
||||
assert_eq!(manager.status().phase, Phase::Ready);
|
||||
for _ in 0..32 {
|
||||
manager.start();
|
||||
}
|
||||
tokio::time::timeout(std::time::Duration::from_secs(2), async {
|
||||
while manager.pipeline_running.load(Ordering::Acquire) {
|
||||
tokio::task::yield_now().await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(manager.status().phase, Phase::Ready);
|
||||
while let Ok(event) = event_rx.try_recv() {
|
||||
if let AppEvent::SimilarityStatus(status) = event {
|
||||
assert_eq!(status.phase, Phase::Ready);
|
||||
}
|
||||
}
|
||||
|
||||
drop(manager);
|
||||
drop(library);
|
||||
std::fs::remove_dir_all(directory).unwrap();
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn resampling_keeps_a_constant_signal() {
|
||||
let output = resample_sinc(&vec![0.25; 441], 44_100, 16_000);
|
||||
|
||||
+126
-105
@@ -103,6 +103,14 @@ fn draw_settings_rows(frame: &mut Frame, area: Rect, state: &AppState) {
|
||||
"Preprocessing profile",
|
||||
format!("{} (enter for details)", similarity.profile),
|
||||
),
|
||||
SimilarityRow::MinimumScore => (
|
||||
"Minimum similarity",
|
||||
format!("{:.2}", similarity.minimum_score),
|
||||
),
|
||||
SimilarityRow::MaxTracksPerArtist => (
|
||||
"Tracks per artist",
|
||||
similarity.max_tracks_per_artist.to_string(),
|
||||
),
|
||||
SimilarityRow::Workers => ("Background workers", similarity.workers.to_string()),
|
||||
SimilarityRow::Clear => ("Clear all stored embeddings", "↵".to_string()),
|
||||
};
|
||||
@@ -400,6 +408,7 @@ fn protocol_label(id: &str) -> &str {
|
||||
"catalog" => "Catalog",
|
||||
"audio" => "Audio transfer",
|
||||
"similarity" => "Similarity search",
|
||||
"similarity_dht" => "Similarity DHT",
|
||||
"device_sync" => "Device sync",
|
||||
"jam" => "Jam",
|
||||
other => other,
|
||||
@@ -540,21 +549,20 @@ fn short_id(id: &str) -> String {
|
||||
}
|
||||
|
||||
fn draw_status_column(frame: &mut Frame, area: Rect, state: &AppState) {
|
||||
if area.height < 12 {
|
||||
draw_status(frame, area, state);
|
||||
return;
|
||||
}
|
||||
let [similarity_area, _, federation_area] = Layout::vertical([
|
||||
Constraint::Length(8),
|
||||
Constraint::Length(1),
|
||||
Constraint::Min(0),
|
||||
])
|
||||
.areas(area);
|
||||
draw_similarity_status(frame, similarity_area, state);
|
||||
draw_status(frame, federation_area, state);
|
||||
draw_status(frame, area, state);
|
||||
}
|
||||
|
||||
fn draw_similarity_status(frame: &mut Frame, area: Rect, state: &AppState) {
|
||||
draw_summary_card(
|
||||
frame,
|
||||
area,
|
||||
state,
|
||||
" Similarity Processing ",
|
||||
similarity_summary_lines(state),
|
||||
);
|
||||
}
|
||||
|
||||
fn similarity_summary_lines(state: &AppState) -> Vec<Line<'static>> {
|
||||
let status = &state.similarity.status;
|
||||
let progress = if status.total_tracks == 0 {
|
||||
"0 / 0".to_string()
|
||||
@@ -571,34 +579,28 @@ fn draw_similarity_status(frame: &mut Frame, area: Rect, state: &AppState) {
|
||||
.as_deref()
|
||||
.map(short_id)
|
||||
.unwrap_or_else(|| "—".to_string());
|
||||
draw_summary_card(
|
||||
frame,
|
||||
area,
|
||||
state,
|
||||
" Similarity Processing ",
|
||||
vec![
|
||||
status_line("State", status.phase.label().to_string()),
|
||||
status_line("Progress", progress),
|
||||
status_line("Active", active),
|
||||
status_line("Processing", target),
|
||||
status_line(
|
||||
"Stored",
|
||||
format!(
|
||||
"{} vectors / {}",
|
||||
status.stored_vectors,
|
||||
short_bytes_label(status.stored_bytes)
|
||||
),
|
||||
vec![
|
||||
summary_line("State", status.phase.label().to_string()),
|
||||
summary_line("Progress", progress),
|
||||
summary_line("Active", active),
|
||||
summary_line("Processing", target),
|
||||
summary_line(
|
||||
"Stored",
|
||||
format!(
|
||||
"{} vectors / {}",
|
||||
status.stored_vectors,
|
||||
short_bytes_label(status.stored_bytes)
|
||||
),
|
||||
status_line(
|
||||
"Current / errors",
|
||||
status
|
||||
.current_track
|
||||
.clone()
|
||||
.or_else(|| status.last_error.clone())
|
||||
.unwrap_or_else(|| format!("{} errors", status.failed_tracks)),
|
||||
),
|
||||
],
|
||||
);
|
||||
),
|
||||
summary_line(
|
||||
"Current",
|
||||
status
|
||||
.current_track
|
||||
.clone()
|
||||
.or_else(|| status.last_error.clone())
|
||||
.unwrap_or_else(|| format!("{} errors", status.failed_tracks)),
|
||||
),
|
||||
]
|
||||
}
|
||||
|
||||
fn draw_status(frame: &mut Frame, area: Rect, state: &AppState) {
|
||||
@@ -615,69 +617,81 @@ fn draw_status(frame: &mut Frame, area: Rect, state: &AppState) {
|
||||
return;
|
||||
}
|
||||
|
||||
if area.width >= 60 && area.height >= 20 {
|
||||
let protocols_height =
|
||||
protocol_card_height(state, area.width.saturating_sub(2), area.height);
|
||||
let [top_area, _, bottom_area, _, protocols_area, _] = Layout::vertical([
|
||||
Constraint::Length(7),
|
||||
Constraint::Length(1),
|
||||
Constraint::Length(7),
|
||||
Constraint::Length(1),
|
||||
Constraint::Length(protocols_height),
|
||||
Constraint::Min(0),
|
||||
])
|
||||
.areas(area);
|
||||
let [status_area, _, local_area] = Layout::horizontal([
|
||||
Constraint::Percentage(50),
|
||||
Constraint::Length(1),
|
||||
Constraint::Percentage(50),
|
||||
])
|
||||
.areas(top_area);
|
||||
let [transport_area, _, devices_area] = Layout::horizontal([
|
||||
Constraint::Percentage(50),
|
||||
Constraint::Length(1),
|
||||
Constraint::Percentage(50),
|
||||
])
|
||||
.areas(bottom_area);
|
||||
draw_summary_card(
|
||||
frame,
|
||||
status_area,
|
||||
state,
|
||||
" Status ",
|
||||
node_summary_lines(state),
|
||||
);
|
||||
draw_summary_card(
|
||||
frame,
|
||||
local_area,
|
||||
state,
|
||||
" Local Data ",
|
||||
local_data_summary_lines(state),
|
||||
);
|
||||
draw_summary_card(
|
||||
frame,
|
||||
transport_area,
|
||||
state,
|
||||
" Iroh Transport ",
|
||||
transport_summary_lines(state),
|
||||
);
|
||||
draw_summary_card(
|
||||
frame,
|
||||
devices_area,
|
||||
state,
|
||||
" Connected Devices ",
|
||||
device_summary_lines(state),
|
||||
);
|
||||
draw_summary_card(
|
||||
frame,
|
||||
protocols_area,
|
||||
state,
|
||||
" Protocol Versions ",
|
||||
protocol_summary_lines(state, protocols_area.width.saturating_sub(2)),
|
||||
);
|
||||
return;
|
||||
if area.width >= 60 {
|
||||
let paired_width = area.width.saturating_sub(1) / 2;
|
||||
let paired_height =
|
||||
protocol_card_height(state, paired_width.saturating_sub(2), area.height).max(8);
|
||||
if area.height >= 16 + paired_height {
|
||||
let [top_area, _, middle_area, _, paired_area, _] = Layout::vertical([
|
||||
Constraint::Length(7),
|
||||
Constraint::Length(1),
|
||||
Constraint::Length(7),
|
||||
Constraint::Length(1),
|
||||
Constraint::Length(paired_height),
|
||||
Constraint::Min(0),
|
||||
])
|
||||
.areas(area);
|
||||
let [status_area, _, local_area] = Layout::horizontal([
|
||||
Constraint::Percentage(50),
|
||||
Constraint::Length(1),
|
||||
Constraint::Percentage(50),
|
||||
])
|
||||
.areas(top_area);
|
||||
let [transport_area, _, devices_area] = Layout::horizontal([
|
||||
Constraint::Percentage(50),
|
||||
Constraint::Length(1),
|
||||
Constraint::Percentage(50),
|
||||
])
|
||||
.areas(middle_area);
|
||||
let [similarity_area, _, protocols_area] = Layout::horizontal([
|
||||
Constraint::Percentage(50),
|
||||
Constraint::Length(1),
|
||||
Constraint::Percentage(50),
|
||||
])
|
||||
.areas(paired_area);
|
||||
draw_summary_card(
|
||||
frame,
|
||||
status_area,
|
||||
state,
|
||||
" Status ",
|
||||
node_summary_lines(state),
|
||||
);
|
||||
draw_summary_card(
|
||||
frame,
|
||||
local_area,
|
||||
state,
|
||||
" Local Data ",
|
||||
local_data_summary_lines(state),
|
||||
);
|
||||
draw_summary_card(
|
||||
frame,
|
||||
transport_area,
|
||||
state,
|
||||
" Iroh Transport ",
|
||||
transport_summary_lines(state),
|
||||
);
|
||||
draw_summary_card(
|
||||
frame,
|
||||
devices_area,
|
||||
state,
|
||||
" Connected Devices ",
|
||||
device_summary_lines(state),
|
||||
);
|
||||
draw_similarity_status(frame, similarity_area, state);
|
||||
draw_summary_card(
|
||||
frame,
|
||||
protocols_area,
|
||||
state,
|
||||
" Protocol Versions ",
|
||||
protocol_summary_lines(state, protocols_area.width.saturating_sub(2)),
|
||||
);
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
if area.height < 39 {
|
||||
let protocols_height = protocol_card_height(state, area.width.saturating_sub(2), area.height);
|
||||
let required_vertical_height = 41u16.saturating_add(protocols_height);
|
||||
if area.height < required_vertical_height {
|
||||
frame.render_widget(
|
||||
Paragraph::new(compact_status_lines(state))
|
||||
.wrap(ratatui::widgets::Wrap { trim: false }),
|
||||
@@ -695,6 +709,8 @@ fn draw_status(frame: &mut Frame, area: Rect, state: &AppState) {
|
||||
_,
|
||||
local_area,
|
||||
_,
|
||||
similarity_area,
|
||||
_,
|
||||
protocols_area,
|
||||
_,
|
||||
] = Layout::vertical([
|
||||
@@ -706,11 +722,9 @@ fn draw_status(frame: &mut Frame, area: Rect, state: &AppState) {
|
||||
Constraint::Length(1),
|
||||
Constraint::Length(7),
|
||||
Constraint::Length(1),
|
||||
Constraint::Length(protocol_card_height(
|
||||
state,
|
||||
area.width.saturating_sub(2),
|
||||
area.height,
|
||||
)),
|
||||
Constraint::Length(8),
|
||||
Constraint::Length(1),
|
||||
Constraint::Length(protocols_height),
|
||||
Constraint::Min(0),
|
||||
])
|
||||
.areas(area);
|
||||
@@ -743,6 +757,7 @@ fn draw_status(frame: &mut Frame, area: Rect, state: &AppState) {
|
||||
" Local Data ",
|
||||
local_data_summary_lines(state),
|
||||
);
|
||||
draw_similarity_status(frame, similarity_area, state);
|
||||
draw_summary_card(
|
||||
frame,
|
||||
protocols_area,
|
||||
@@ -766,6 +781,12 @@ fn compact_status_lines(state: &AppState) -> Vec<Line<'static>> {
|
||||
lines.push(Line::styled("Connected Devices", theme::header_for(state)));
|
||||
lines.extend(device_summary_lines(state).into_iter().take(2));
|
||||
lines.push(Line::default());
|
||||
lines.push(Line::styled(
|
||||
"Similarity Processing",
|
||||
theme::header_for(state),
|
||||
));
|
||||
lines.extend(similarity_summary_lines(state).into_iter().take(3));
|
||||
lines.push(Line::default());
|
||||
lines.push(Line::styled("Protocol Versions", theme::header_for(state)));
|
||||
lines.extend(protocol_summary_lines(state, 0).into_iter().take(3));
|
||||
lines
|
||||
|
||||
@@ -900,6 +900,9 @@ fn draw_search(frame: &mut Frame, area: Rect, state: &AppState, cursor: usize) {
|
||||
title.push_str("· searching… ");
|
||||
}
|
||||
let inner = bordered(frame, area, state, title);
|
||||
if search.similarity_source.is_some() {
|
||||
return draw_similarity_search(frame, inner, state, cursor);
|
||||
}
|
||||
|
||||
let empty_results = SearchResults::default();
|
||||
let results = match &search.results {
|
||||
@@ -1102,6 +1105,156 @@ fn draw_search(frame: &mut Frame, area: Rect, state: &AppState, cursor: usize) {
|
||||
}
|
||||
}
|
||||
|
||||
fn draw_similarity_search(frame: &mut Frame, area: Rect, state: &AppState, cursor: usize) {
|
||||
let search = &state.search;
|
||||
let status = if search.fed_loading {
|
||||
super::loading_line(state, "searching federation…")
|
||||
} else if let Some(stats) = &search.similarity_stats {
|
||||
let elapsed = if stats.elapsed_ms < 1_000 {
|
||||
format!("{} ms", stats.elapsed_ms)
|
||||
} else {
|
||||
format!("{:.2} s", stats.elapsed_ms as f64 / 1_000.0)
|
||||
};
|
||||
Line::from(vec![
|
||||
Span::styled("Federation · ", theme::accent_for(state)),
|
||||
Span::styled(
|
||||
format!(
|
||||
"{} tracks · {} artists · {} peers queried · {elapsed}",
|
||||
stats.tracks, stats.artists, stats.peers_queried
|
||||
),
|
||||
theme::dim(),
|
||||
),
|
||||
])
|
||||
} else if search.similarity_error.is_some() {
|
||||
Line::styled(
|
||||
"Federation search failed · showing local results",
|
||||
error_style(),
|
||||
)
|
||||
} else if search.loading {
|
||||
super::loading_line(state, "preparing local similarity search…")
|
||||
} else {
|
||||
Line::styled("Federation disabled · showing local results", theme::dim())
|
||||
};
|
||||
let [status_area, content] =
|
||||
Layout::vertical([Constraint::Length(1), Constraint::Min(0)]).areas(area);
|
||||
frame.render_widget(Paragraph::new(status), status_area);
|
||||
|
||||
let mut rows: Vec<(Line, Option<String>, Option<usize>)> = Vec::new();
|
||||
rows.push((Line::styled("Tracks", theme::header_for(state)), None, None));
|
||||
if let Some(track) = &search.similarity_source_track {
|
||||
let heart = if state.track_liked(track) {
|
||||
Span::styled("♥ ", theme::accent_for(state))
|
||||
} else {
|
||||
Span::raw(" ")
|
||||
};
|
||||
rows.push((
|
||||
Line::from(vec![
|
||||
heart,
|
||||
Span::raw(track.title.clone()),
|
||||
Span::styled(
|
||||
format!(" {} · {}", track.artist_line(), track.release_title),
|
||||
theme::dim(),
|
||||
),
|
||||
]),
|
||||
Some(super::track_meta_suffix(track, true)),
|
||||
Some(0),
|
||||
));
|
||||
}
|
||||
for (offset, hit) in search.similarity_tracks.iter().enumerate() {
|
||||
let index = offset + 1;
|
||||
match hit {
|
||||
crate::app::state::SimilaritySearchHit::Local { track, .. } => {
|
||||
let heart = if state.track_liked(track) {
|
||||
Span::styled("♥ ", theme::accent_for(state))
|
||||
} else {
|
||||
Span::raw(" ")
|
||||
};
|
||||
rows.push((
|
||||
Line::from(vec![
|
||||
heart,
|
||||
Span::raw(track.title.clone()),
|
||||
Span::styled(
|
||||
format!(" {} · {}", track.artist_line(), track.release_title),
|
||||
theme::dim(),
|
||||
),
|
||||
]),
|
||||
Some(super::track_meta_suffix(track, true)),
|
||||
Some(index),
|
||||
));
|
||||
}
|
||||
crate::app::state::SimilaritySearchHit::Federated { track, .. } => {
|
||||
let heart = if state.fed_track_liked(track) {
|
||||
Span::styled("♥ ", theme::accent_for(state))
|
||||
} else {
|
||||
Span::raw(" ")
|
||||
};
|
||||
let origin = if track.own {
|
||||
"your library".to_string()
|
||||
} else {
|
||||
format!("peer {}…", track.owner_short())
|
||||
};
|
||||
let mut meta = track.duration_label();
|
||||
if let Some(year) = track.year {
|
||||
if !meta.is_empty() {
|
||||
meta.push_str(" · ");
|
||||
}
|
||||
meta.push_str(&year.to_string());
|
||||
}
|
||||
rows.push((
|
||||
Line::from(vec![
|
||||
heart,
|
||||
fed_track_availability_prefix(state, track),
|
||||
Span::raw(track.title.clone()),
|
||||
Span::styled(
|
||||
format!(" {} · {origin}", track.artist_line()),
|
||||
theme::dim(),
|
||||
),
|
||||
]),
|
||||
Some(meta),
|
||||
Some(index),
|
||||
));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let scope = crate::app::state::TrackSelectionScope::SimilaritySearch;
|
||||
let mut selected = std::collections::HashSet::new();
|
||||
if state.track_selection.is_active_for(&scope)
|
||||
&& let Some(indices) = state
|
||||
.track_selection
|
||||
.indices(&scope, search.similarity_len())
|
||||
{
|
||||
selected.extend(indices);
|
||||
}
|
||||
let cursor_row = rows
|
||||
.iter()
|
||||
.position(|(_, _, row_cursor)| *row_cursor == Some(cursor))
|
||||
.unwrap_or(0);
|
||||
let visible = usize::from(content.height.max(1));
|
||||
let first = cursor_row
|
||||
.saturating_sub(visible / 2)
|
||||
.min(rows.len().saturating_sub(visible));
|
||||
for (offset, (line, right, row_cursor)) in
|
||||
rows.into_iter().enumerate().skip(first).take(visible)
|
||||
{
|
||||
let rect = Rect {
|
||||
x: content.x,
|
||||
y: content.y + (offset - first) as u16,
|
||||
width: content.width,
|
||||
height: 1,
|
||||
};
|
||||
if let Some(row_index) = row_cursor
|
||||
&& selected.contains(&row_index)
|
||||
&& row_index != cursor
|
||||
{
|
||||
frame
|
||||
.buffer_mut()
|
||||
.set_style(rect, theme::selection_for(state));
|
||||
}
|
||||
draw_row(frame, rect, state, line, right, row_cursor == Some(cursor));
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Federated artist card (assembled from peer catalogs)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
Reference in New Issue
Block a user