Fixed UI bugs

This commit is contained in:
Aleksandr Bogomiakov
2026-08-14 02:14:46 +01:00
parent 2dc40c2c3f
commit 0d7c9b9f5e
24 changed files with 4351 additions and 158 deletions
+1
View File
@@ -216,6 +216,7 @@ impl Actor {
}
urgent_targets.push(previous);
}
self.finish_listen(music_dht::device_sync::ListenEndReason::Stopped);
self.audio.stop();
self.device_role = DevicePlaybackRole::Control;
self.active_device_id.clone_from(&target.id);
+253 -4
View File
@@ -12,10 +12,11 @@ use std::time::Duration;
use anyhow::{Context as _, Result};
use music_dht::device_sync::{
DEFAULT_INVITE_TTL_MS, DEVICE_SYNC_PROTOCOL_VERSION, DeviceProfileWire, InviteWire,
PlaybackCommand, PlaybackSnapshot, SnapshotLike, SnapshotLikeTombstone, SnapshotPlaylist,
SnapshotPlaylistItem, SnapshotPlaylistItemTombstone, SnapshotPlaylistTombstone, SyncOpPayload,
SyncOpWire, SyncSnapshot, SyncedFedTrack, WireMessage, encode_invite, finish_response,
finish_send, hash_secret, parse_invite, random_hex, read_msg, ticket_endpoint_id, write_msg,
ListenEndReason, ListenEvent, ListenTrackMetadata, PlaybackCommand, PlaybackSnapshot,
SnapshotLike, SnapshotLikeTombstone, SnapshotPlaylist, SnapshotPlaylistItem,
SnapshotPlaylistItemTombstone, SnapshotPlaylistTombstone, SyncOpPayload, SyncOpWire,
SyncSnapshot, SyncedFedTrack, WireMessage, encode_invite, finish_response, finish_send,
hash_secret, parse_invite, random_hex, read_msg, ticket_endpoint_id, write_msg,
};
use music_dht::{ByteStream, MusicDhtService, PeerTicket, StreamAcceptor};
use rusqlite::{Connection, OptionalExtension as _, params};
@@ -120,6 +121,10 @@ impl DeviceSync {
Ok((identity.device_id, identity.name))
}
pub fn new_listen_id() -> String {
format!("{}-{}", now_ms(), random_hex(12))
}
pub fn set_device_name(&self, name: &str, endpoint_ticket: Option<&str>) -> Result<()> {
let name = if name.trim().is_empty() {
"furumi".to_string()
@@ -367,6 +372,56 @@ impl DeviceSync {
})
}
pub fn record_listen(
&self,
listen_id: String,
track: &furumi_domain::Track,
started_at_ms: i64,
listened_ms: i64,
ended_reason: ListenEndReason,
) -> Result<()> {
let Some(content_id) = track
.key
.content_id()
.and_then(|id| music_dht::normalize_content_id(id.as_str()))
else {
return Ok(());
};
let mut artist_names = track
.artists
.iter()
.map(|artist| artist.name.clone())
.filter(|artist| !artist.trim().is_empty())
.collect::<Vec<_>>();
if artist_names.is_empty() && !track.artist.trim().is_empty() {
artist_names.push(track.artist.clone());
}
let event = ListenEvent {
listen_id,
content_id,
started_at_ms,
listened_ms: listened_ms.max(0),
track_duration_ms: (track.duration_seconds > 0.0).then_some(
crate::support::seconds_to_milliseconds(track.duration_seconds),
),
ended_reason,
track: ListenTrackMetadata {
title: track.title.clone(),
artist_names,
featured_artist_names: track
.featured_artists
.iter()
.map(|artist| artist.name.clone())
.collect(),
release_title: (!track.release.trim().is_empty()).then(|| track.release.clone()),
},
};
if event.should_record() {
self.record_op(SyncOpPayload::ListenRecorded { event })?;
}
Ok(())
}
pub fn record_playlist_created(&self, id: i64, title: &str) -> Result<()> {
let playlist_id = self.library.ensure_playlist_sync_id(id)?;
self.record_op(SyncOpPayload::PlaylistCreated {
@@ -698,6 +753,8 @@ async fn handle_pair_request(
finish_response(&mut stream, RESPONSE_DRAIN).await?;
return Ok(());
}
let requester_group_id = requester_group_id.filter(|group| !group.trim().is_empty());
let requester_group_active_devices = requester_group_active_devices.max(1);
let devices_json = serde_json::to_string(&requester_devices)?;
lock(&sync.conn).execute(
"INSERT OR IGNORE INTO sync_pending_pairing
@@ -1016,6 +1073,15 @@ fn pair_request_id(invite_id: &str, device_id: &str) -> String {
format!("pair_{}", &digest[..16])
}
fn requester_group_conflict(
local_group_id: &str,
requester_group_id: Option<&str>,
requester_group_active_devices: usize,
) -> bool {
requester_group_id.is_some_and(|group| !group.trim().is_empty() && group != local_group_id)
&& requester_group_active_devices > 1
}
fn valid_pair_request(
sync: &DeviceSync,
invite_id: &str,
@@ -1115,6 +1181,150 @@ enum PairAttempt {
mod tests {
use super::*;
static NEXT_TEST_DB: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(1);
fn with_test_sync(test: impl FnOnce(&DeviceSync)) {
let conn = Connection::open_in_memory().unwrap();
init_schema(&conn).unwrap();
let unique = NEXT_TEST_DB.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let library_path = std::env::temp_dir().join(format!(
"furumi-desktop-devices-test-{}-{}-{}.sqlite3",
std::process::id(),
now_ms(),
unique
));
let library = Arc::new(furumi_library::Library::open(&library_path).unwrap());
let (events, _event_rx) = tokio::sync::mpsc::channel(4);
let sync = DeviceSync {
conn: Arc::new(Mutex::new(conn)),
library,
events,
playback: Arc::new(Mutex::new(PlaybackStore::default())),
sync_requested: Arc::new(tokio::sync::Notify::new()),
};
sync.ensure_identity().unwrap();
test(&sync);
drop(sync);
for suffix in ["", "-wal", "-shm"] {
let mut path = library_path.as_os_str().to_os_string();
path.push(suffix);
let _ = std::fs::remove_file(path);
}
}
fn test_profile(device_id: &str) -> DeviceProfileWire {
DeviceProfileWire {
device_id: device_id.into(),
name: device_id.into(),
client_version: CLIENT_VERSION.into(),
protocol_version: DEVICE_SYNC_PROTOCOL_VERSION,
endpoint_id: String::new(),
endpoint_ticket: String::new(),
revoked: false,
revoke_cutoff_seq: None,
updated_at_ms: now_ms(),
}
}
#[test]
fn locally_finished_listen_enters_the_shared_history() {
with_test_sync(|sync| {
let content_id =
furumi_domain::ContentId::parse(format!("b3:{}", "a".repeat(64))).unwrap();
let track = furumi_domain::Track {
key: furumi_domain::TrackKey::remote(content_id.clone()),
title: "Shared listen".into(),
artist: "Artist".into(),
artists: vec![furumi_domain::ArtistRef {
key: furumi_domain::ArtistKey::Federation {
peer_id: "peer".into(),
id: "artist".into(),
},
name: "Artist".into(),
}],
featured_artists: Vec::new(),
release: "Release".into(),
release_id: furumi_domain::ReleaseKey::Federation {
peer_id: "peer".into(),
id: "release".into(),
},
duration_seconds: 180.0,
track_number: Some(1),
disc_number: Some(1),
cover_uri: None,
audio_format: None,
audio_bitrate_kbps: None,
audio_sample_rate_hz: None,
audio_bit_depth: None,
file_size_bytes: None,
liked: false,
audio_source: furumi_domain::AudioSource::Federation {
peer_id: "peer".into(),
content_id,
},
};
sync.record_listen(
"desktop-listen".into(),
&track,
now_ms(),
180_000,
ListenEndReason::Finished,
)
.unwrap();
let history = sync.library.listen_history(10).unwrap();
assert_eq!(history.len(), 1);
assert_eq!(history[0].listen_id, "desktop-listen");
assert_eq!(history[0].title, "Shared listen");
});
}
fn insert_pending_pairing(sync: &DeviceSync, requester_group_devices: &[DeviceProfileWire]) {
lock(&sync.conn)
.execute(
"INSERT INTO sync_pending_pairing
(request_id, device_id, name, client_version, endpoint_id,
endpoint_ticket, invite_id, created_at_ms, status,
requester_group_id, requester_group_active_devices,
requester_group_devices_json)
VALUES ('request', 'dev_requester', 'Requester', ?1, '', '',
'invite', ?2, 'pending', 'grp_remote', 2, ?3)",
params![
CLIENT_VERSION,
now_ms(),
serde_json::to_string(requester_group_devices).unwrap()
],
)
.unwrap();
}
fn device_known(sync: &DeviceSync, device_id: &str) -> bool {
lock(&sync.conn)
.query_row(
"SELECT 1 FROM sync_devices WHERE device_id = ?1",
[device_id],
|_| Ok(()),
)
.optional()
.unwrap()
.is_some()
}
fn device_trusted(sync: &DeviceSync, device_id: &str) -> bool {
lock(&sync.conn)
.query_row(
"SELECT trusted_at_ms IS NOT NULL FROM sync_devices WHERE device_id = ?1",
[device_id],
|row| row.get::<_, bool>(0),
)
.optional()
.unwrap()
.unwrap_or(false)
}
#[test]
fn pairing_request_ids_are_stable_and_scoped_to_the_device() {
assert_eq!(
@@ -1126,4 +1336,43 @@ mod tests {
pair_request_id("invite", "device-b")
);
}
#[test]
fn group_choice_is_only_required_for_an_existing_different_group() {
assert!(requester_group_conflict("grp_local", Some("grp_remote"), 2));
assert!(!requester_group_conflict(
"grp_local",
Some("grp_remote"),
1
));
assert!(!requester_group_conflict("grp_local", Some("grp_local"), 3));
assert!(!requester_group_conflict("grp_local", None, 3));
assert!(!requester_group_conflict("grp_local", Some(" "), 3));
}
#[test]
fn pairing_choice_either_joins_the_requester_group_or_keeps_the_local_group() {
let requester_peer = test_profile("dev_requester_peer");
with_test_sync(|sync| {
let local_group = sync.ensure_identity().unwrap().group_id;
insert_pending_pairing(sync, std::slice::from_ref(&requester_peer));
sync.answer_pairing("request", true, false).unwrap();
assert_eq!(sync.ensure_identity().unwrap().group_id, local_group);
assert!(device_trusted(sync, "dev_requester"));
assert!(!device_known(sync, "dev_requester_peer"));
});
with_test_sync(|sync| {
insert_pending_pairing(sync, std::slice::from_ref(&requester_peer));
sync.answer_pairing("request", true, true).unwrap();
assert_eq!(sync.ensure_identity().unwrap().group_id, "grp_remote");
assert!(device_trusted(sync, "dev_requester"));
assert!(device_known(sync, "dev_requester_peer"));
});
}
}
+21 -5
View File
@@ -108,6 +108,7 @@ impl DeviceSync {
.collect::<rusqlite::Result<Vec<_>>>()?
};
let pending = {
let local_group_id = identity.group_id.as_str();
let mut stmt = conn.prepare(
"SELECT request_id, device_id, name, client_version,
requester_group_id, requester_group_active_devices
@@ -115,14 +116,25 @@ impl DeviceSync {
ORDER BY created_at_ms",
)?;
stmt.query_map([], |row| {
let requester_group_id = row.get::<_, Option<String>>(4)?;
let requester_group_active_devices =
usize::try_from(row.get::<_, i64>(5)?.max(0)).unwrap_or(usize::MAX);
let group_conflict = requester_group_conflict(
local_group_id,
requester_group_id.as_deref(),
requester_group_active_devices,
);
Ok(PendingPairing {
request_id: row.get(0)?,
device_id: row.get(1)?,
name: row.get(2)?,
client_version: row.get(3)?,
requester_group_id: row.get(4)?,
requester_group_active_devices: usize::try_from(row.get::<_, i64>(5)?.max(0))
.unwrap_or(usize::MAX),
requester_group_id: group_conflict.then_some(requester_group_id).flatten(),
requester_group_active_devices: if group_conflict {
requester_group_active_devices
} else {
0
},
})
})?
.collect::<rusqlite::Result<Vec<_>>>()?
@@ -480,8 +492,12 @@ impl DeviceSync {
command,
} => self.apply_playback_command(target_device_id, command, &op.op_id)?,
SyncOpPayload::ListenRecorded { event } => {
self.library
.apply_listen_event(event, &op.origin_device_id)?;
if self
.library
.apply_listen_event(event, &op.origin_device_id)?
{
self.notify_library();
}
}
}
Ok(())
+71 -5
View File
@@ -14,6 +14,8 @@ use furumi_domain::{
use music_dht::catalog::{
CATALOG_ALPN, CatalogArtist, CatalogImageHeader, CatalogRequest, CatalogResponse,
};
use music_dht::similarity_dht::SimilarityDht;
use music_dht::similarity_lsh::SIMILARITY_DHT_ALPN;
use music_dht::{
EndpointId, ItemKind, ItemSpec, LibraryItem, MusicDhtConfig, MusicDhtService, NetworkId,
RendezvousConfig,
@@ -82,29 +84,60 @@ const IMAGE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(4);
pub struct Client {
service: Arc<MusicDhtService>,
similarity_dht: Arc<SimilarityDht>,
media_dir: PathBuf,
}
impl Client {
pub async fn start(data_dir: PathBuf, media_dir: PathBuf, network: &str) -> Result<Arc<Self>> {
pub async fn start(
data_dir: PathBuf,
media_dir: PathBuf,
network: &str,
similarity: Arc<crate::similarity::Manager>,
) -> Result<Arc<Self>> {
tokio::fs::create_dir_all(&data_dir).await?;
tokio::fs::create_dir_all(&media_dir).await?;
let similarity_routing_path = data_dir.join("similarity-routing.sqlite3");
let config = MusicDhtConfig::builder()
.data_dir(data_dir)
.data_dir(&data_dir)
.network_id(NetworkId::from_name(network))
.rendezvous(RendezvousConfig::default())
.stream_protocol(CATALOG_ALPN)
.stream_protocol(AUDIO_ALPN)
.stream_protocol(music_dht::device_sync::SYNC_ALPN_V1)
.stream_protocol(music_dht::device_sync::SYNC_ALPN_V2)
.schema_independent_stream_protocol(crate::federation_similarity::SIMILARITY_ALPN)
.schema_independent_stream_protocol(SIMILARITY_DHT_ALPN)
.build()
.context("invalid federation configuration")?;
let (service, mut events) = MusicDhtService::start(config)
.await
.context("starting federation node")?;
tokio::spawn(async move { while events.recv().await.is_some() {} });
let service = Arc::new(service);
let similarity_dht = SimilarityDht::open(Arc::clone(&service), similarity_routing_path)
.await
.context("starting similarity routing overlay")?;
let routing_acceptor = service
.stream_acceptor(SIMILARITY_DHT_ALPN)
.context("starting similarity routing listener")?;
tokio::spawn(Arc::clone(&similarity_dht).serve(routing_acceptor));
tokio::spawn(Arc::clone(&similarity_dht).maintenance());
tokio::spawn(crate::federation_similarity::sync_routes(
Arc::clone(&similarity_dht),
Arc::clone(&similarity),
));
let similarity_acceptor = service
.stream_acceptor(crate::federation_similarity::SIMILARITY_ALPN)
.context("starting similarity listener")?;
tokio::spawn(crate::federation_similarity::serve(
similarity_acceptor,
similarity,
service.endpoint_id(),
));
Ok(Arc::new(Self {
service: Arc::new(service),
service,
similarity_dht,
media_dir,
}))
}
@@ -171,6 +204,33 @@ impl Client {
Ok((results, stats))
}
pub async fn search_similar(
&self,
query: crate::similarity::QueryVector,
limit: usize,
minimum_score: f32,
max_tracks_per_artist: usize,
) -> Result<Vec<crate::federation_similarity::ScoredTrack>> {
let mut hits = crate::federation_similarity::search(
Arc::clone(&self.service),
Arc::clone(&self.similarity_dht),
query,
limit,
minimum_score,
max_tracks_per_artist,
)
.await?;
let mut results = SearchResults {
tracks: hits.iter().map(|hit| hit.track.clone()).collect(),
..SearchResults::default()
};
self.fetch_artwork(&mut results).await;
for (hit, with_artwork) in hits.iter_mut().zip(results.tracks) {
hit.track = with_artwork;
}
Ok(hits)
}
/// Resolves a portable queue entry when another connected device only
/// knows its stable audio content id.
pub async fn track_by_content_id(&self, content_id: &str) -> Result<Track> {
@@ -181,11 +241,17 @@ impl Client {
.into_iter()
.chain(outcome.network_results)
.collect::<Vec<_>>();
convert_items(&items, own)
let mut track = convert_items(&items, own)
.tracks
.into_iter()
.next()
.context("no federation peer currently publishes this track")
.context("no federation peer currently publishes this track")?;
if track.cover_uri.is_none()
&& let Some(path) = self.artwork_for_track(&track).await
{
track.cover_uri = Some(path.to_string_lossy().into_owned());
}
Ok(track)
}
pub async fn publish(&self, specs: Vec<ItemSpec>) -> Result<()> {
+378
View File
@@ -0,0 +1,378 @@
//! Local-index adapter and DHT-routed peer client for Furumi similarity.
use std::collections::{HashMap, HashSet};
use std::sync::Arc;
use std::time::Duration;
use anyhow::{Context as _, Result};
use furumi_domain::{ArtistKey, ArtistRef, AudioSource, ContentId, ReleaseKey, Track, TrackKey};
use futures_util::stream::{self, StreamExt as _};
use music_dht::similarity::{self as wire, SimilarityHit, SimilarityRequest, SimilarityResponse};
use music_dht::similarity_dht::SimilarityDht;
use music_dht::{EndpointId, ItemId, ItemKind, MusicDhtService, PeerTicket, StreamAcceptor};
use crate::similarity::{Manager, QueryVector};
pub use music_dht::similarity::SIMILARITY_ALPN;
const MAX_QUERY_PEERS: usize = 48;
const QUERY_CONCURRENCY: usize = 8;
const QUERY_TIMEOUT: Duration = Duration::from_secs(5);
const ROUTING_TIMEOUT: Duration = Duration::from_secs(5);
const ROUTE_SYNC_INTERVAL: Duration = Duration::from_secs(30);
const MAX_NEAR_DUPLICATE_SIGNATURE_DISTANCE: u32 = 8;
pub struct ScoredTrack {
pub track: Track,
pub score: f32,
pub embedding_signature: Option<[u8; wire::SIMILARITY_SIGNATURE_BYTES]>,
}
pub async fn serve(mut acceptor: StreamAcceptor, manager: Arc<Manager>, own: EndpointId) {
while let Some(stream) = acceptor.accept().await {
let manager = Arc::clone(&manager);
tokio::spawn(async move {
let _ = serve_one(stream, manager, own).await;
});
}
}
async fn serve_one(
mut stream: music_dht::ByteStream,
manager: Arc<Manager>,
own: EndpointId,
) -> Result<()> {
let request = wire::read_request(&mut stream).await?;
let response = if manager.network_allowed() {
let matches = tokio::task::spawn_blocking(move || {
manager.search_vector_for_peer(&request.profile_id, &request.vector, request.limit)
})
.await
.context("local similarity task failed")
.and_then(|result| result);
match matches {
Ok(matches) => SimilarityResponse::success(
matches
.into_iter()
.filter_map(|found| {
let track = found.track;
let hit = SimilarityHit {
score: found.score,
item_id: hex_encode(
ItemId::derive(
&own,
ItemKind::Track,
&format!("track:{}", track.id),
)
.as_bytes(),
),
title: track.title,
artist_names: track
.artists
.into_iter()
.map(|artist| artist.name)
.collect(),
featured_artist_names: track
.featured_artists
.into_iter()
.map(|artist| artist.name)
.collect(),
year: track.release_year,
duration_seconds: Some(
crate::support::seconds_to_milliseconds(track.duration_seconds)
/ 1_000,
),
content_id: track.content_id,
release_title: Some(track.release_title),
track_number: track.track_number,
disc_number: track.disc_number,
embedding_signature: Some(found.embedding_signature),
};
hit.validate().is_ok().then_some(hit)
})
.collect(),
)?,
Err(error) => {
SimilarityResponse::refused(format!("similarity query is unavailable: {error:#}"))?
}
}
} else {
SimilarityResponse::refused("similarity federation is disabled or has no privacy consent")?
};
wire::write_response(&mut stream, &response).await?;
stream.send.finish()?;
let _ = stream.send.stopped().await;
Ok(())
}
#[allow(
clippy::too_many_lines,
reason = "bounded peer fan-out and ranking policy"
)]
pub async fn search(
service: Arc<MusicDhtService>,
routing: Arc<SimilarityDht>,
query: QueryVector,
limit: usize,
minimum_score: f32,
max_tracks_per_artist: usize,
) -> Result<Vec<ScoredTrack>> {
let own = service.endpoint_id();
let routed = tokio::time::timeout(
ROUTING_TIMEOUT,
routing.find_peers(&query.profile_id, &query.vector, MAX_QUERY_PEERS),
)
.await
.ok()
.and_then(Result::ok)
.unwrap_or_default();
let mut seen = HashSet::new();
let mut peers = routed
.into_iter()
.filter_map(|ticket| {
let owner = ticket.endpoint_id();
(owner != own && seen.insert(owner)).then_some(QueryPeer {
owner,
ticket: Some(ticket),
})
})
.collect::<Vec<_>>();
for owner in service
.connected_peers()
.into_iter()
.chain(service.known_peers().into_iter().map(|peer| peer.peer_id))
{
if owner != own && seen.insert(owner) {
peers.push(QueryPeer {
owner,
ticket: None,
});
}
if peers.len() >= MAX_QUERY_PEERS {
break;
}
}
let query_signature = wire::embedding_signature(&query.vector)?;
let request = Arc::new(SimilarityRequest::new(
query.profile_id,
query.vector,
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);
async move {
tokio::time::timeout(QUERY_TIMEOUT, query_peer(service, peer, &request))
.await
.map_err(|_| anyhow::anyhow!("similarity peer timed out"))?
}
}))
.buffer_unordered(QUERY_CONCURRENCY)
.collect::<Vec<_>>()
.await;
let mut hits = responses
.into_iter()
.filter_map(Result::ok)
.flatten()
.collect::<Vec<_>>();
hits.sort_by(|left, right| right.1.total_cmp(&left.1));
let mut dedup = HashSet::new();
let mut signatures = vec![query_signature];
let mut artist_counts = HashMap::<String, usize>::new();
let mut results = Vec::new();
for (track, score, signature) in hits {
if score < minimum_score {
break;
}
if query.source_content_id.as_deref().is_some_and(|source| {
track
.key
.content_id()
.is_some_and(|id| id.as_str() == source)
}) {
continue;
}
let identity = track
.key
.content_id()
.map_or_else(|| format!("{:?}", track.key), |id| id.as_str().to_owned());
if !dedup.insert(identity) {
continue;
}
if signature.is_some_and(|candidate| {
signatures.iter().any(|existing| {
wire::signature_distance(&candidate, existing)
<= MAX_NEAR_DUPLICATE_SIGNATURE_DISTANCE
})
}) {
continue;
}
let artist = track
.artists
.first()
.map(|artist| music_dht::normalize_name(&artist.name))
.unwrap_or_default();
let count = artist_counts.entry(artist.clone()).or_default();
if !artist.is_empty() && *count >= max_tracks_per_artist {
continue;
}
*count += 1;
if let Some(signature) = signature {
signatures.push(signature);
}
results.push(ScoredTrack {
track,
score,
embedding_signature: signature,
});
if results.len() >= limit.min(wire::MAX_SIMILARITY_RESULTS) {
break;
}
}
Ok(results)
}
type PeerHits = Vec<(Track, f32, Option<[u8; wire::SIMILARITY_SIGNATURE_BYTES]>)>;
#[derive(Clone)]
struct QueryPeer {
owner: EndpointId,
ticket: Option<PeerTicket>,
}
async fn query_peer(
service: Arc<MusicDhtService>,
peer: QueryPeer,
request: &SimilarityRequest,
) -> 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,
}?;
let response = wire::exchange(&mut stream, request).await?;
anyhow::ensure!(
response.ok,
"peer refused similarity query: {}",
response.error.unwrap_or_default()
);
Ok(response
.hits
.into_iter()
.filter_map(|hit| hit_to_track(owner, hit))
.collect())
}
/// Keeps the routing overlay synchronized with the active durable index.
pub async fn sync_routes(routing: Arc<SimilarityDht>, manager: Arc<Manager>) {
let mut interval = tokio::time::interval(ROUTE_SYNC_INTERVAL);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
let mut published_marker: Option<(String, blake3::Hash)> = None;
loop {
interval.tick().await;
let manager = Arc::clone(&manager);
let loaded = tokio::task::spawn_blocking(move || manager.routing_signatures()).await;
let Ok(Ok(snapshot)) = loaded else {
continue;
};
let Some((profile_id, signatures)) = snapshot else {
if published_marker.take().is_some() {
routing.clear_local_signatures();
}
continue;
};
let mut hasher = blake3::Hasher::new();
for signature in &signatures {
hasher.update(signature);
}
let marker = (profile_id.clone(), hasher.finalize());
if published_marker.as_ref() == Some(&marker) {
continue;
}
if routing
.sync_local_signatures(profile_id, signatures)
.await
.is_ok()
{
published_marker = Some(marker);
}
}
}
fn hit_to_track(
owner: EndpointId,
hit: SimilarityHit,
) -> Option<(Track, f32, Option<[u8; wire::SIMILARITY_SIGNATURE_BYTES]>)> {
let peer_id = owner.to_string();
let content_id = hit
.content_id
.as_deref()
.and_then(|id| ContentId::parse(id).ok())?;
let refs = |names: Vec<String>| {
names
.into_iter()
.map(|name| ArtistRef {
key: ArtistKey::Federation {
peer_id: peer_id.clone(),
id: music_dht::normalize_name(&name),
},
name,
})
.collect::<Vec<_>>()
};
let artists = refs(hit.artist_names);
let featured_artists = refs(hit.featured_artist_names);
let artist = artists
.iter()
.map(|artist| artist.name.as_str())
.chain(featured_artists.iter().map(|artist| artist.name.as_str()))
.collect::<Vec<_>>()
.join(", ");
let release = hit.release_title.unwrap_or_default();
let score = hit.score;
let signature = hit.embedding_signature;
Some((
Track {
key: TrackKey::federation(peer_id.clone(), hit.item_id, Some(content_id.clone())),
title: hit.title,
artist,
artists,
featured_artists,
release: release.clone(),
release_id: ReleaseKey::Federation {
peer_id: peer_id.clone(),
id: format!("name:{}", music_dht::normalize_name(&release)),
},
duration_seconds: hit
.duration_seconds
.and_then(|value| u32::try_from(value).ok())
.map_or(0.0, f64::from),
track_number: hit.track_number.and_then(|value| u32::try_from(value).ok()),
disc_number: hit.disc_number.and_then(|value| u32::try_from(value).ok()),
cover_uri: None,
audio_format: None,
audio_bitrate_kbps: None,
audio_sample_rate_hz: None,
audio_bit_depth: None,
file_size_bytes: None,
liked: false,
audio_source: AudioSource::Federation {
peer_id,
content_id,
},
},
score,
signature,
))
}
fn hex_encode(bytes: &[u8]) -> String {
const HEX: &[u8; 16] = b"0123456789abcdef";
let mut output = String::with_capacity(bytes.len() * 2);
for byte in bytes {
output.push(char::from(HEX[usize::from(byte >> 4)]));
output.push(char::from(HEX[usize::from(byte & 0x0f)]));
}
output
}
+415 -18
View File
@@ -13,7 +13,7 @@ use furumi_backend_api::{
FederationActivitySnapshot, FederationDebugSnapshot, FederationOperation, LibrarySnapshot,
PendingPairingSnapshot, PlaybackRepeat, PlaybackStatus, PlaylistSnapshot, RemoteData,
RequestId, SearchResults, SearchSnapshot, SearchStats, SendCommandError, SettingsSnapshot,
VersionEntrySnapshot,
SimilarityStatusSnapshot, VersionEntrySnapshot,
};
use furumi_domain::{
Artist, ArtistId, ArtistKey, ArtistRef, Artwork, AudioSource, CatalogSource, ContentId,
@@ -26,18 +26,20 @@ mod actor_devices;
mod audio;
mod devices;
mod federation;
mod federation_similarity;
mod settings;
mod similarity;
mod streaming;
mod support;
use support::{
apply_federated_metadata, expand_tilde, extrapolated_control_position, federation_specs,
find_catalog_track, library_snapshot, library_track, local_search_results,
apply_federated_metadata, extrapolated_control_position, federated_audio_directory,
federation_specs, find_catalog_track, library_snapshot, library_track, local_search_results,
merge_release_preserving_local, merge_search_results, normalize_device_name,
playback_state_acknowledges_command, portable_playback_placeholder,
remote_snapshot_has_authority, runtime_build_info, sanitize_filename, selected_track_position,
spawn_settings_worker, track_is_liked, track_to_library_fed, track_to_playback_track,
track_to_synced_fed, unix_time_ms, volume_percent,
remote_snapshot_has_authority, runtime_build_info, sanitize_filename, seconds_to_milliseconds,
selected_track_position, spawn_settings_worker, track_is_liked, track_to_library_fed,
track_to_playback_track, track_to_synced_fed, unix_time_ms, volume_percent,
};
use settings::SettingsStore;
@@ -46,6 +48,14 @@ const COMMAND_CAPACITY: usize = 64;
const INTERNAL_CAPACITY: usize = 32;
const CONTROL_COMMAND_ACK_TIMEOUT: Duration = Duration::from_secs(12);
const CONTROL_POSITION_ACK_TOLERANCE_SECONDS: f64 = 3.0;
const DESKTOP_APPLICATION_ID: &str = "furumi-desktop";
const SIMILARITY_RESULT_LIMIT: usize = 50;
type SimilarityCandidate = (
Track,
f32,
Option<[u8; music_dht::similarity::SIMILARITY_SIGNATURE_BYTES]>,
);
#[derive(Clone)]
pub struct BackendHandle {
@@ -81,20 +91,26 @@ impl BackendHandle {
///
/// Returns an I/O error when the runtime or its owner thread cannot be created.
pub fn spawn_backend() -> Result<BackendHandle, BackendStartupError> {
let project_dirs = directories::ProjectDirs::from("cy", "hexor", "Furumi")
// The music library deliberately keeps using furumi_library::default_db_path()
// below. Everything else belongs to this client and must not reuse the TUI's
// device identity, settings, federation node, or caches.
let project_dirs = directories::ProjectDirs::from("cy", "hexor", DESKTOP_APPLICATION_ID)
.ok_or(BackendStartupError::DataDirectoryUnavailable)?;
let settings_path = project_dirs.data_local_dir().join("furumi-desktop.sqlite3");
let federation_data_dir = project_dirs.data_local_dir().join("federation");
let federation_media_dir = project_dirs.cache_dir().join("federation-media");
let similarity_model_dir = project_dirs.cache_dir().join("similarity-models");
let devices_db_path = project_dirs
.data_local_dir()
.join("devices")
.join("sync.sqlite3");
let settings_store = SettingsStore::open(&settings_path)?;
let library_db_path =
furumi_library::default_db_path().map_err(BackendStartupError::Library)?;
let default_library_path = library_db_path.with_file_name("federation-media");
let settings_store = SettingsStore::open(&settings_path, &default_library_path)?;
let loaded_settings = settings_store.load()?;
let library_path = furumi_library::default_db_path().map_err(BackendStartupError::Library)?;
let catalog = std::sync::Arc::new(
furumi_library::Library::open(&library_path).map_err(BackendStartupError::Library)?,
furumi_library::Library::open(&library_db_path).map_err(BackendStartupError::Library)?,
);
let loaded_library = library_snapshot(&catalog).map_err(BackendStartupError::Library)?;
let runtime = tokio::runtime::Builder::new_multi_thread()
@@ -132,6 +148,7 @@ pub fn spawn_backend() -> Result<BackendHandle, BackendStartupError> {
audio,
federation_data_dir,
federation_media_dir,
similarity_model_dir,
devices,
}));
})?;
@@ -218,6 +235,10 @@ enum InternalEvent {
keys: Vec<TrackKey>,
cover_uri: Option<String>,
},
HistoryTrackResolved {
content_id: String,
result: Result<Track, String>,
},
FederationDebugUpdated(FederationDebugSnapshot),
DevicesChanged,
DeviceLibraryChanged,
@@ -226,6 +247,12 @@ enum InternalEvent {
DeviceOperationFinished(Result<DeviceOperationResult, String>),
DevicePlaybackSnapshot(music_dht::device_sync::PlaybackSnapshot),
DevicePlaybackCommand(music_dht::device_sync::PlaybackCommand),
SimilarityStatus(similarity::SimilarityStatus),
SimilarityProfileActivated(Option<String>),
SimilaritySearchFinished {
source_title: String,
result: Result<Vec<Track>, String>,
},
}
enum DeviceOperationResult {
@@ -246,8 +273,10 @@ struct Actor {
federation: Option<std::sync::Arc<federation::Client>>,
federation_data_dir: std::path::PathBuf,
federation_media_dir: std::path::PathBuf,
similarity: std::sync::Arc<similarity::Manager>,
ephemeral_audio: Option<(TrackKey, std::path::PathBuf)>,
pending_queue_artwork: HashSet<String>,
pending_history_resolutions: HashSet<String>,
federation_debug_pending: bool,
devices: std::sync::Arc<devices::DeviceSync>,
device_service: Option<std::sync::Arc<music_dht::MusicDhtService>>,
@@ -256,6 +285,14 @@ struct Actor {
active_device_name: String,
control_anchor: Option<ControlPlaybackAnchor>,
pending_control: Option<PendingControlState>,
listen_session: Option<ListenSession>,
}
#[derive(Debug, Clone)]
struct ListenSession {
id: String,
track: Track,
started_at_ms: i64,
}
#[derive(Debug, Clone)]
@@ -285,6 +322,7 @@ struct ActorBootstrap {
audio: audio::Controller,
federation_data_dir: std::path::PathBuf,
federation_media_dir: std::path::PathBuf,
similarity_model_dir: std::path::PathBuf,
devices: std::sync::Arc<devices::DeviceSync>,
}
@@ -307,6 +345,7 @@ fn initialize_actor(bootstrap: ActorBootstrap) -> ActorRuntime {
audio,
federation_data_dir,
federation_media_dir,
similarity_model_dir,
devices,
} = bootstrap;
let (settings_tx, settings_rx) = std_mpsc::channel();
@@ -336,6 +375,12 @@ fn initialize_actor(bootstrap: ActorBootstrap) -> ActorRuntime {
settings_error,
..BackendSnapshot::default()
};
let similarity = similarity::Manager::new(
std::sync::Arc::clone(&catalog),
internal_tx.clone(),
initial_state.settings.similarity.clone(),
similarity_model_dir,
);
let mut actor = Actor {
state: initial_state,
library,
@@ -349,8 +394,10 @@ fn initialize_actor(bootstrap: ActorBootstrap) -> ActorRuntime {
federation: None,
federation_data_dir,
federation_media_dir,
similarity,
ephemeral_audio: None,
pending_queue_artwork: HashSet::new(),
pending_history_resolutions: HashSet::new(),
federation_debug_pending: false,
devices,
device_service: None,
@@ -359,7 +406,12 @@ fn initialize_actor(bootstrap: ActorBootstrap) -> ActorRuntime {
active_device_name,
control_anchor: None,
pending_control: None,
listen_session: None,
};
actor.state.similarity_status = similarity_status_snapshot(&actor.similarity.status());
if actor.state.settings.similarity.enabled {
actor.similarity.start();
}
actor.refresh_connected_devices();
ActorRuntime {
actor,
@@ -368,6 +420,22 @@ fn initialize_actor(bootstrap: ActorBootstrap) -> ActorRuntime {
}
}
fn similarity_status_snapshot(status: &similarity::SimilarityStatus) -> SimilarityStatusSnapshot {
SimilarityStatusSnapshot {
phase: status.phase.label().into(),
active_profile: status.active_profile.clone(),
target_profile: status.target_profile.clone(),
model: status.model.clone(),
total_tracks: status.total_tracks,
completed_tracks: status.completed_tracks,
failed_tracks: status.failed_tracks,
stored_vectors: status.stored_vectors,
stored_bytes: status.stored_bytes,
current_track: status.current_track.clone(),
error: status.last_error.clone(),
}
}
async fn run_actor(bootstrap: ActorBootstrap) {
let ActorRuntime {
mut actor,
@@ -400,6 +468,7 @@ async fn run_actor(bootstrap: ActorBootstrap) {
if let Some((_, token)) = actor.active_search.take() {
token.cancel();
}
actor.finish_listen(music_dht::device_sync::ListenEndReason::Stopped);
actor.audio.stop();
}
@@ -416,6 +485,7 @@ impl Actor {
let data_dir = self.federation_data_dir.clone();
let media_dir = self.federation_media_dir.clone();
let network = self.state.settings.network_id.trim().to_owned();
let similarity = std::sync::Arc::clone(&self.similarity);
let internal = self.internal.clone();
self.state.federation_activity = FederationActivitySnapshot {
operation: FederationOperation::Idle,
@@ -425,7 +495,7 @@ impl Actor {
};
self.publish();
tokio::spawn(async move {
let result = federation::Client::start(data_dir, media_dir, &network)
let result = federation::Client::start(data_dir, media_dir, &network, similarity)
.await
.map_err(|error| format!("federation: {error:#}"));
let _ = internal
@@ -659,6 +729,7 @@ impl Actor {
}
}
BackendCommand::Stop => {
self.finish_listen(music_dht::device_sync::ListenEndReason::Stopped);
if self.device_role == DevicePlaybackRole::Active {
self.audio.stop();
}
@@ -725,6 +796,36 @@ impl Actor {
self.play_current();
}
}
BackendCommand::MoveQueueItem {
item_id,
target_index,
} => {
if self.state.queue.move_item(item_id, target_index) {
self.send_control_state(false);
self.publish();
}
}
BackendCommand::RemoveQueueItem { item_id } => {
if let Some(removed_current) = self.state.queue.remove_item(item_id) {
if removed_current {
self.finish_listen(music_dht::device_sync::ListenEndReason::Stopped);
self.remove_ephemeral_audio();
if self.state.queue.current().is_some() {
self.play_current();
} else {
self.audio.stop();
self.state.playback.status = PlaybackStatus::Stopped;
self.state.playback.position_seconds = 0.0;
self.state.playback.duration_seconds = 0.0;
self.send_control_state(true);
self.publish();
}
} else {
self.send_control_state(false);
self.publish();
}
}
}
BackendCommand::PlayContext { tracks, selected } => {
let tracks = self.resolve_tracks(&tracks);
if !tracks.is_empty() {
@@ -845,6 +946,8 @@ impl Actor {
}
}
BackendCommand::Search { request_id, query } => self.start_search(request_id, query),
BackendCommand::SearchSimilar { track } => self.start_similarity_search(&track),
BackendCommand::ClearSimilarity => self.similarity.clear(),
BackendCommand::CancelSearch { request_id } => self.cancel_search(request_id),
BackendCommand::LoadArtist {
request_id,
@@ -857,11 +960,13 @@ impl Actor {
artist_key,
artist_name,
} => self.load_detail(request_id, artist_key, artist_name, Some(key)),
BackendCommand::UpdateSettings(settings) => {
BackendCommand::UpdateSettings(mut settings) => {
settings.similarity = settings.similarity.normalized();
let federation_changed = self.state.settings.federation_enabled
!= settings.federation_enabled
|| self.state.settings.network_id != settings.network_id;
let device_name_changed = self.state.settings.device_name != settings.device_name;
let similarity_changed = self.state.settings.similarity != settings.similarity;
let pending_device_name = settings.device_name.clone();
self.state.settings = settings.clone();
self.state.settings_error = None;
@@ -878,6 +983,9 @@ impl Actor {
self.apply_local_device_name(&pending_device_name);
self.schedule_device_name_publish(pending_device_name);
}
if similarity_changed {
self.similarity.apply(&self.state.settings.similarity);
}
}
BackendCommand::Shutdown => {}
}
@@ -1147,6 +1255,7 @@ impl Actor {
});
}
self.federation = Some(client);
self.resolve_history_tracks();
self.state.federation_activity.pending = false;
self.resolve_queue_artwork();
self.refresh_federation_debug();
@@ -1223,6 +1332,7 @@ impl Actor {
.current()
.is_some_and(|item| item.track.key.matches(&key))
{
self.finish_listen(music_dht::device_sync::ListenEndReason::Stopped);
self.state.playback.status = PlaybackStatus::Stopped;
self.state.playback_error = Some(message);
self.publish();
@@ -1270,6 +1380,7 @@ impl Actor {
}
}
Err(message) => {
self.finish_listen(music_dht::device_sync::ListenEndReason::Stopped);
self.state.playback.status = PlaybackStatus::Stopped;
self.state.playback_error = Some(message);
self.publish();
@@ -1299,6 +1410,27 @@ impl Actor {
self.publish();
}
}
InternalEvent::HistoryTrackResolved { content_id, result } => {
self.pending_history_resolutions.remove(&content_id);
if let Ok(mut resolved) = result {
let mut changed = false;
for track in &mut self.library.recently_played {
if track
.key
.content_id()
.is_some_and(|id| id.as_str() == content_id)
{
resolved.liked = track.liked;
*track = resolved.clone();
changed = true;
}
}
if changed {
self.state.library = RemoteData::Ready(self.library.clone());
self.publish();
}
}
}
InternalEvent::FederationDebugUpdated(debug) => {
self.federation_debug_pending = false;
if self.state.federation_debug != debug {
@@ -1311,6 +1443,8 @@ impl Actor {
self.library = library.clone();
self.state.library = RemoteData::Ready(library);
self.reconcile_likes();
self.similarity.start();
self.resolve_history_tracks();
}
self.refresh_connected_devices();
self.publish();
@@ -1338,9 +1472,129 @@ impl Actor {
InternalEvent::DevicePlaybackCommand(command) => {
self.apply_device_playback_command(command);
}
InternalEvent::SimilarityStatus(status) => {
self.state.similarity_status = similarity_status_snapshot(&status);
self.publish();
}
InternalEvent::SimilarityProfileActivated(active_profile) => {
self.state.settings.similarity.active_profile = active_profile;
let _ = self.settings.send(self.state.settings.clone());
self.publish();
}
InternalEvent::SimilaritySearchFinished {
source_title,
result,
} => {
self.state.similarity_search.pending = false;
self.state.similarity_search.source_title = source_title;
match result {
Ok(mut tracks) => {
let liked = self
.catalog
.liked_content_ids()
.unwrap_or_default()
.into_iter()
.chain(self.catalog.fed_like_ids().unwrap_or_default())
.collect::<HashSet<_>>();
for track in &mut tracks {
track.liked = track_is_liked(track, &liked);
}
self.state.similarity_search.results = tracks;
self.state.similarity_search.error = None;
}
Err(error) => {
self.state.similarity_search.results.clear();
self.state.similarity_search.error = Some(error);
}
}
self.publish();
}
}
}
fn start_similarity_search(&mut self, key: &TrackKey) {
let Some(source) = self.track(key).cloned() else {
self.state.similarity_search.error = Some("Track is unavailable".into());
self.publish();
return;
};
if !self.state.settings.similarity.enabled {
self.state.similarity_search.error =
Some("Enable Similarity search in Settings first".into());
self.state.similarity_search.results.clear();
self.publish();
return;
}
let Some(track_id) = source.key.local_id().map(LocalTrackId::get) else {
self.state.similarity_search.error =
Some("Similarity search currently starts from a local track".into());
self.state.similarity_search.results.clear();
self.publish();
return;
};
self.state
.similarity_search
.source_title
.clone_from(&source.title);
self.state.similarity_search.results.clear();
self.state.similarity_search.error = None;
self.state.similarity_search.pending = true;
self.publish();
let manager = std::sync::Arc::clone(&self.similarity);
let federation = self.federation.clone();
let settings = self.state.settings.similarity.clone();
let internal = self.internal.clone();
let source_title = source.title;
tokio::spawn(async move {
let local = tokio::task::spawn_blocking(move || manager.search_track(track_id, 50))
.await
.map_err(|error| format!("similarity worker failed: {error}"))
.and_then(|result| result.map_err(|error| format!("{error:#}")));
let result = match local {
Ok((local, query)) => {
let mut scored = local
.into_iter()
.map(|found| {
(
library_track(found.track, ""),
found.score,
Some(found.embedding_signature),
)
})
.collect::<Vec<_>>();
if settings.federation_consent
&& let Some(federation) = federation
&& let Ok(remote) = federation
.search_similar(
query,
50,
settings.minimum_score,
settings.max_tracks_per_artist,
)
.await
{
scored.extend(
remote
.into_iter()
.map(|hit| (hit.track, hit.score, hit.embedding_signature)),
);
}
Ok(rank_similarity_candidates(
scored,
settings.max_tracks_per_artist,
))
}
Err(error) => Err(error),
};
let _ = internal
.send(InternalEvent::SimilaritySearchFinished {
source_title,
result,
})
.await;
});
}
fn apply_remote_playback_snapshot(
&mut self,
snapshot: music_dht::device_sync::PlaybackSnapshot,
@@ -1378,6 +1632,9 @@ impl Actor {
self.pending_control = None;
}
if self.device_role == DevicePlaybackRole::Active {
self.finish_listen(music_dht::device_sync::ListenEndReason::Stopped);
}
self.device_role = DevicePlaybackRole::Control;
self.active_device_id.clone_from(&snapshot.device_id);
self.active_device_name.clone_from(&snapshot.device_name);
@@ -1445,6 +1702,7 @@ impl Actor {
self.publish();
}
audio::Event::Finished => {
self.finish_listen(music_dht::device_sync::ListenEndReason::Finished);
self.remove_ephemeral_audio();
if self.advance_queue(true) {
self.play_current();
@@ -1456,6 +1714,7 @@ impl Actor {
}
}
audio::Event::Failed(message) => {
self.finish_listen(music_dht::device_sync::ListenEndReason::Stopped);
self.state.playback.status = PlaybackStatus::Stopped;
self.state.playback_error = Some(message);
self.publish();
@@ -1466,10 +1725,12 @@ impl Actor {
fn play_current(&mut self) {
self.remove_ephemeral_audio();
let Some(track) = self.state.queue.current().map(|item| item.track.clone()) else {
self.finish_listen(music_dht::device_sync::ListenEndReason::Stopped);
self.state.playback.status = PlaybackStatus::Stopped;
self.publish();
return;
};
self.begin_listen(&track);
self.state.playback.position_seconds = 0.0;
self.state.playback.duration_seconds = track.duration_seconds.max(0.0);
self.state.playback_error = None;
@@ -1495,6 +1756,46 @@ impl Actor {
self.publish();
}
fn begin_listen(&mut self, track: &Track) {
if self.device_role != DevicePlaybackRole::Active {
return;
}
if let Some(session) = &mut self.listen_session
&& session.track.key.matches(&track.key)
{
session.track = track.clone();
return;
}
self.finish_listen(music_dht::device_sync::ListenEndReason::Replaced);
if track.key.content_id().is_some() {
self.listen_session = Some(ListenSession {
id: devices::DeviceSync::new_listen_id(),
track: track.clone(),
started_at_ms: unix_time_ms(),
});
}
}
fn finish_listen(&mut self, ended_reason: music_dht::device_sync::ListenEndReason) {
let Some(session) = self.listen_session.take() else {
return;
};
let listened_ms = if ended_reason == music_dht::device_sync::ListenEndReason::Finished {
seconds_to_milliseconds(session.track.duration_seconds)
} else {
seconds_to_milliseconds(self.state.playback.position_seconds)
};
if let Err(error) = self.devices.record_listen(
session.id,
&session.track,
session.started_at_ms,
listened_ms,
ended_reason,
) {
self.state.settings_error = Some(format!("listening history: {error:#}"));
}
}
fn shuffle_new_context(&mut self) {
if self.state.playback.shuffle {
self.state.queue.shuffle_upcoming();
@@ -1514,6 +1815,8 @@ impl Actor {
fn start_federated_playback(&mut self, track: Track, peer_id: String, content_id: ContentId) {
let Some(client) = self.federation.clone() else {
self.finish_listen(music_dht::device_sync::ListenEndReason::Stopped);
self.state.playback.status = PlaybackStatus::Stopped;
self.state.playback_error = Some("federation is still starting".into());
self.publish();
return;
@@ -1548,11 +1851,11 @@ impl Actor {
let key = track.key.clone();
let internal = self.internal.clone();
let keep = self.state.settings.save_federated_on_listen;
let directory = if keep {
expand_tilde(&self.state.settings.library_path)
} else {
self.federation_media_dir.join("stream-cache")
};
let directory = federated_audio_directory(
&self.state.settings.library_path,
keep,
&self.federation_media_dir,
);
let stem = format!("fed-{}", sanitize_filename(&track.title));
tokio::spawn(async move {
let (tx, mut rx) = mpsc::channel(4);
@@ -1673,7 +1976,20 @@ impl Actor {
}
fn track(&self, key: &TrackKey) -> Option<&Track> {
find_catalog_track(&self.library, &self.state.search.results, key)
self.state
.queue
.items()
.iter()
.map(|item| &item.track)
.find(|track| track.key.matches(key))
.or_else(|| find_catalog_track(&self.library, &self.state.search.results, key))
.or_else(|| {
self.state
.similarity_search
.results
.iter()
.find(|track| track.key.matches(key))
})
}
fn resolve_tracks(&self, keys: &[TrackKey]) -> Vec<Track> {
@@ -1688,6 +2004,7 @@ impl Actor {
self.library = library.clone();
self.state.library = RemoteData::Ready(library);
self.reconcile_likes();
self.resolve_history_tracks();
}
Err(error) => {
self.state.settings_error = Some(format!("library: {error:#}"));
@@ -1696,6 +2013,35 @@ impl Actor {
self.publish();
}
fn resolve_history_tracks(&mut self) {
let Some(client) = self.federation.clone() else {
return;
};
let unresolved = self
.library
.recently_played
.iter()
.filter(|track| track.key.local_id().is_none() && track.cover_uri.is_none())
.filter_map(|track| track.key.content_id().map(|id| id.as_str().to_owned()))
.collect::<HashSet<_>>();
for content_id in unresolved {
if !self.pending_history_resolutions.insert(content_id.clone()) {
continue;
}
let internal = self.internal.clone();
let client = client.clone();
tokio::spawn(async move {
let result = client
.track_by_content_id(&content_id)
.await
.map_err(|error| format!("history track lookup failed: {error:#}"));
let _ = internal
.send(InternalEvent::HistoryTrackResolved { content_id, result })
.await;
});
}
}
fn reconcile_likes(&mut self) {
let liked = self
.catalog
@@ -1714,11 +2060,15 @@ impl Actor {
) {
track.liked = track_is_liked(track, &liked);
}
for track in &mut self.state.similarity_search.results {
track.liked = track_is_liked(track, &liked);
}
let replacements = self
.library
.featured_releases
.iter()
.flat_map(|release| release.tracks.iter())
.chain(self.library.recently_played.iter())
.chain(
self.library
.playlists
@@ -1848,5 +2198,52 @@ impl Actor {
}
}
fn rank_similarity_candidates(
mut candidates: Vec<SimilarityCandidate>,
max_tracks_per_artist: usize,
) -> Vec<Track> {
const MAX_NEAR_DUPLICATE_SIGNATURE_DISTANCE: u32 = 8;
candidates.sort_by(|left, right| right.1.total_cmp(&left.1));
let mut content = HashSet::new();
let mut signatures = Vec::new();
let mut artist_counts = HashMap::<String, usize>::new();
let mut tracks = Vec::new();
for (track, _, signature) in candidates {
let identity = track
.key
.content_id()
.map_or_else(|| format!("{:?}", track.key), |id| id.as_str().to_owned());
if !content.insert(identity) {
continue;
}
if signature.is_some_and(|candidate| {
signatures.iter().any(|existing| {
music_dht::similarity::signature_distance(&candidate, existing)
<= MAX_NEAR_DUPLICATE_SIGNATURE_DISTANCE
})
}) {
continue;
}
let artist = track.artists.first().map_or_else(
|| music_dht::normalize_name(&track.artist),
|artist| music_dht::normalize_name(&artist.name),
);
let count = artist_counts.entry(artist.clone()).or_default();
if !artist.is_empty() && *count >= max_tracks_per_artist.clamp(1, SIMILARITY_RESULT_LIMIT) {
continue;
}
*count += 1;
if let Some(signature) = signature {
signatures.push(signature);
}
tracks.push(track);
if tracks.len() >= SIMILARITY_RESULT_LIMIT {
break;
}
}
tracks
}
#[cfg(test)]
mod tests;
+136 -8
View File
@@ -4,6 +4,9 @@ use std::path::Path;
use furumi_backend_api::SettingsSnapshot;
use rusqlite::{Connection, OptionalExtension, params};
const LEGACY_DEFAULT_LIBRARY_PATH: &str = "~/Music/Furumi";
const PLATFORM_LIBRARY_PATH_MIGRATION: i64 = 4;
const MIGRATIONS: &[(i64, &str)] = &[
(
1,
@@ -28,6 +31,19 @@ const MIGRATIONS: &[(i64, &str)] = &[
3,
"ALTER TABLE app_settings ADD COLUMN device_name TEXT NOT NULL DEFAULT '';",
),
(
5,
r"
ALTER TABLE app_settings ADD COLUMN similarity_enabled INTEGER NOT NULL DEFAULT 0 CHECK (similarity_enabled IN (0, 1));
ALTER TABLE app_settings ADD COLUMN similarity_model TEXT NOT NULL DEFAULT 'discogs-effnet-bsdynamic-1';
ALTER TABLE app_settings ADD COLUMN similarity_profile TEXT NOT NULL DEFAULT 'furumi-full-track-v1';
ALTER TABLE app_settings ADD COLUMN similarity_workers INTEGER NOT NULL DEFAULT 2;
ALTER TABLE app_settings ADD COLUMN similarity_minimum_score REAL NOT NULL DEFAULT 0.70;
ALTER TABLE app_settings ADD COLUMN similarity_max_tracks_per_artist INTEGER NOT NULL DEFAULT 5;
ALTER TABLE app_settings ADD COLUMN similarity_federation_consent INTEGER NOT NULL DEFAULT 0 CHECK (similarity_federation_consent IN (0, 1));
ALTER TABLE app_settings ADD COLUMN similarity_active_profile TEXT;
",
),
];
pub struct SettingsStore {
@@ -35,32 +51,47 @@ pub struct SettingsStore {
}
impl SettingsStore {
pub fn open(path: &Path) -> rusqlite::Result<Self> {
pub fn open(path: &Path, default_library_path: &Path) -> rusqlite::Result<Self> {
if let Some(parent) = path.parent() {
fs::create_dir_all(parent)
.map_err(|error| rusqlite::Error::ToSqlConversionFailure(Box::new(error)))?;
}
let connection = Connection::open(path)?;
let mut store = Self { connection };
store.migrate()?;
store.migrate(default_library_path)?;
Ok(store)
}
#[cfg(test)]
fn in_memory() -> rusqlite::Result<Self> {
fn in_memory(default_library_path: &Path) -> rusqlite::Result<Self> {
let connection = Connection::open_in_memory()?;
let mut store = Self { connection };
store.migrate()?;
store.migrate(default_library_path)?;
Ok(store)
}
pub fn load(&self) -> rusqlite::Result<SettingsSnapshot> {
self.connection.query_row(
"SELECT network_id, library_path, federation_enabled, language,
save_federated_on_listen, device_name
save_federated_on_listen, device_name,
similarity_enabled, similarity_model, similarity_profile,
similarity_workers, similarity_minimum_score,
similarity_max_tracks_per_artist, similarity_federation_consent,
similarity_active_profile
FROM app_settings WHERE singleton_id = 1",
[],
|row| {
let similarity = furumi_backend_api::SimilaritySettingsSnapshot {
enabled: row.get::<_, i64>(6)? != 0,
model: row.get(7)?,
profile: row.get(8)?,
workers: usize::try_from(row.get::<_, i64>(9)?).unwrap_or(1),
minimum_score: row.get(10)?,
max_tracks_per_artist: usize::try_from(row.get::<_, i64>(11)?).unwrap_or(1),
federation_consent: row.get::<_, i64>(12)? != 0,
active_profile: row.get(13)?,
}
.normalized();
Ok(SettingsSnapshot {
network_id: row.get(0)?,
library_path: row.get(1)?,
@@ -68,6 +99,7 @@ impl SettingsStore {
language: row.get(3)?,
save_federated_on_listen: row.get::<_, i64>(4)? != 0,
device_name: row.get(5)?,
similarity,
})
},
)
@@ -81,7 +113,15 @@ impl SettingsStore {
federation_enabled = ?3,
language = ?4,
save_federated_on_listen = ?5,
device_name = ?6
device_name = ?6,
similarity_enabled = ?7,
similarity_model = ?8,
similarity_profile = ?9,
similarity_workers = ?10,
similarity_minimum_score = ?11,
similarity_max_tracks_per_artist = ?12,
similarity_federation_consent = ?13,
similarity_active_profile = ?14
WHERE singleton_id = 1",
params![
settings.network_id,
@@ -90,12 +130,20 @@ impl SettingsStore {
settings.language,
i64::from(settings.save_federated_on_listen),
settings.device_name,
i64::from(settings.similarity.enabled),
settings.similarity.model,
settings.similarity.profile,
i64::try_from(settings.similarity.workers).unwrap_or(16),
settings.similarity.minimum_score,
i64::try_from(settings.similarity.max_tracks_per_artist).unwrap_or(50),
i64::from(settings.similarity.federation_consent),
settings.similarity.active_profile,
],
)?;
Ok(())
}
fn migrate(&mut self) -> rusqlite::Result<()> {
fn migrate(&mut self, default_library_path: &Path) -> rusqlite::Result<()> {
self.connection.execute_batch(
"PRAGMA foreign_keys = ON;
CREATE TABLE IF NOT EXISTS schema_migrations (
@@ -119,14 +167,55 @@ impl SettingsStore {
}
let transaction = self.connection.transaction()?;
transaction.execute_batch(sql)?;
if version == 5 {
transaction.execute(
"UPDATE app_settings SET similarity_workers = ?1 WHERE singleton_id = 1",
[i64::try_from(
furumi_backend_api::SimilaritySettingsSnapshot::default().workers,
)
.unwrap_or(1)],
)?;
}
transaction.execute(
"INSERT INTO schema_migrations (version) VALUES (?1)",
[version],
)?;
transaction.commit()?;
}
self.migrate_platform_library_path(default_library_path)?;
Ok(())
}
fn migrate_platform_library_path(
&mut self,
default_library_path: &Path,
) -> rusqlite::Result<()> {
let applied = self
.connection
.query_row(
"SELECT 1 FROM schema_migrations WHERE version = ?1",
[PLATFORM_LIBRARY_PATH_MIGRATION],
|_| Ok(()),
)
.optional()?
.is_some();
if applied {
return Ok(());
}
let default_library_path = default_library_path.to_string_lossy();
let transaction = self.connection.transaction()?;
transaction.execute(
"UPDATE app_settings
SET library_path = ?1
WHERE library_path = ?2 OR trim(library_path) = ''",
params![default_library_path.as_ref(), LEGACY_DEFAULT_LIBRARY_PATH],
)?;
transaction.execute(
"INSERT INTO schema_migrations (version) VALUES (?1)",
[PLATFORM_LIBRARY_PATH_MIGRATION],
)?;
transaction.commit()
}
}
#[cfg(test)]
@@ -135,19 +224,58 @@ mod tests {
#[test]
fn migration_creates_defaults_and_settings_round_trip() {
let store = SettingsStore::in_memory().unwrap();
let default_library_path = Path::new("/platform/furumi/federation-media");
let store = SettingsStore::in_memory(default_library_path).unwrap();
let mut settings = store.load().unwrap();
assert_eq!(settings.network_id, "furumi");
assert_eq!(
settings.library_path,
default_library_path.to_string_lossy()
);
assert!(settings.federation_enabled);
assert!(settings.save_federated_on_listen);
assert!(settings.device_name.is_empty());
assert!(!settings.similarity.enabled);
assert_eq!(
settings.similarity.workers,
furumi_backend_api::SimilaritySettingsSnapshot::default().workers
);
assert!((settings.similarity.minimum_score - 0.70).abs() < f32::EPSILON);
settings.network_id = "friends".into();
settings.device_name = "Studio Mac".into();
settings.library_path = "/music/library".into();
settings.federation_enabled = false;
settings.similarity.enabled = true;
settings.similarity.workers = 7;
settings.similarity.minimum_score = 0.82;
settings.similarity.max_tracks_per_artist = 9;
settings.similarity.federation_consent = true;
settings.similarity.active_profile = Some("sim1:test".into());
store.save(&settings).unwrap();
assert_eq!(store.load().unwrap(), settings);
}
#[test]
fn platform_default_migration_preserves_a_custom_library_path() {
let first_default = Path::new("/first/furumi/federation-media");
let mut store = SettingsStore::in_memory(first_default).unwrap();
let mut settings = store.load().unwrap();
settings.library_path = "/custom/music".into();
store.save(&settings).unwrap();
store
.connection
.execute(
"DELETE FROM schema_migrations WHERE version = ?1",
[PLATFORM_LIBRARY_PATH_MIGRATION],
)
.unwrap();
store
.migrate(Path::new("/second/furumi/federation-media"))
.unwrap();
assert_eq!(store.load().unwrap().library_path, "/custom/music");
}
}
File diff suppressed because it is too large Load Diff
+77 -5
View File
@@ -15,6 +15,18 @@ pub(super) fn expand_tilde(value: &str) -> std::path::PathBuf {
}
}
pub(super) fn federated_audio_directory(
library_path: &str,
keep: bool,
federation_cache_dir: &std::path::Path,
) -> std::path::PathBuf {
if keep {
expand_tilde(library_path)
} else {
federation_cache_dir.join("stream-cache")
}
}
pub(super) fn normalize_device_name(value: &str) -> String {
let value = value.trim();
if value.is_empty() {
@@ -31,10 +43,18 @@ pub(super) fn selected_track_position(tracks: &[Track], selected: &TrackKey) ->
.unwrap_or(0)
}
pub(super) fn seconds_to_milliseconds(seconds: f64) -> i64 {
if !seconds.is_finite() || seconds <= 0.0 {
return 0;
}
let duration = Duration::from_secs_f64(seconds.min(Duration::MAX.as_secs_f64()));
i64::try_from(duration.as_millis()).unwrap_or(i64::MAX)
}
pub(super) fn runtime_build_info() -> BuildInfoSnapshot {
use music_dht::capabilities::{
CATALOG_ID, CapabilityManifest, DEVICE_SYNC_ID, FEDERATION_NET_ID, MUSIC_DHT_ID,
RENDEZVOUS_ID, TICKET_ID,
RENDEZVOUS_ID, SIMILARITY_DHT_ID, SIMILARITY_ID, TICKET_ID,
};
let manifest = CapabilityManifest::frid("furumi-desktop", env!("CARGO_PKG_VERSION"));
@@ -70,6 +90,8 @@ pub(super) fn runtime_build_info() -> BuildInfoSnapshot {
protocol("Rendezvous", RENDEZVOUS_ID),
protocol("Music DHT", MUSIC_DHT_ID),
protocol("Catalog", CATALOG_ID),
protocol("Similarity search", SIMILARITY_ID),
protocol("Similarity routing", SIMILARITY_DHT_ID),
VersionEntrySnapshot {
name: "Audio transfer".into(),
version: federation::AUDIO_PROTOCOL_VERSION.to_string(),
@@ -98,6 +120,7 @@ pub(super) fn find_catalog_track<'a>(
.featured_releases
.iter()
.flat_map(|release| release.tracks.iter())
.chain(library.recently_played.iter())
.chain(
library
.playlists
@@ -785,10 +808,22 @@ pub(super) fn library_snapshot(
tracks,
});
}
let recently_played = releases
.iter()
.flat_map(|release| release.tracks.iter().cloned())
.take(12)
let recently_played = catalog
.listen_history(500)?
.into_iter()
.filter_map(|entry| {
let content_id = ContentId::parse(entry.content_id.clone()).ok()?;
let mut track = catalog
.track_by_content_id(content_id.as_str())
.ok()
.flatten()
.map_or_else(
|| history_placeholder(&entry, content_id),
|track| library_track(track, ""),
);
track.liked = track_is_liked(&track, &liked_ids);
Some(track)
})
.collect();
let mut playlists = Vec::new();
for card in catalog.playlists()? {
@@ -816,6 +851,43 @@ pub(super) fn library_snapshot(
})
}
fn history_placeholder(entry: &furumi_library::ListenHistoryEntry, content_id: ContentId) -> Track {
let peer_id = "history".to_owned();
let artist = ArtistRef {
key: ArtistKey::Federation {
peer_id: peer_id.clone(),
id: format!("name:{}", music_dht::normalize_name(&entry.artist)),
},
name: entry.artist.clone(),
};
Track {
key: TrackKey::remote(content_id.clone()),
title: entry.title.clone(),
artist: entry.artist.clone(),
artists: vec![artist],
featured_artists: Vec::new(),
release: String::new(),
release_id: ReleaseKey::Federation {
peer_id: peer_id.clone(),
id: format!("history:{}", entry.listen_id),
},
duration_seconds: 0.0,
track_number: None,
disc_number: None,
cover_uri: None,
audio_format: None,
audio_bitrate_kbps: None,
audio_sample_rate_hz: None,
audio_bit_depth: None,
file_size_bytes: None,
liked: false,
audio_source: AudioSource::Federation {
peer_id,
content_id,
},
}
}
pub(super) fn track_is_liked(track: &Track, liked_ids: &HashSet<String>) -> bool {
track
.key
+65
View File
@@ -1,5 +1,19 @@
use super::*;
#[test]
fn configured_library_path_controls_permanent_federated_audio_storage() {
let cache = std::path::Path::new("/cache/furumi-desktop/federation-media");
assert_eq!(
federated_audio_directory("/chosen/music", true, cache),
std::path::Path::new("/chosen/music")
);
assert_eq!(
federated_audio_directory("/chosen/music", false, cache),
cache.join("stream-cache")
);
}
fn merge_test_track(key: TrackKey, audio_source: AudioSource) -> Track {
Track {
key,
@@ -26,6 +40,57 @@ fn merge_test_track(key: TrackKey, audio_source: AudioSource) -> Track {
}
}
fn similarity_test_candidate(
id: i64,
artist: &str,
score: f32,
signature: u8,
) -> SimilarityCandidate {
let mut track = merge_test_track(
TrackKey::local(LocalTrackId::new(id)),
AudioSource::LocalFile(format!("{id}.flac").into()),
);
track.title = format!("Track {id}");
track.artist = artist.into();
track.artists[0].name = artist.into();
(
track,
score,
Some([signature; music_dht::similarity::SIMILARITY_SIGNATURE_BYTES]),
)
}
#[test]
fn similarity_results_apply_one_artist_cap_after_combining_sources() {
let tracks = rank_similarity_candidates(
vec![
similarity_test_candidate(1, "same artist", 0.70, 1),
similarity_test_candidate(2, "other artist", 0.90, 2),
similarity_test_candidate(3, "same artist", 0.80, 3),
similarity_test_candidate(4, "same artist", 0.60, 4),
],
1,
);
assert_eq!(tracks.len(), 2);
assert_eq!(tracks[0].title, "Track 2");
assert_eq!(tracks[1].title, "Track 3");
}
#[test]
fn similarity_results_drop_cross_source_near_duplicates() {
let tracks = rank_similarity_candidates(
vec![
similarity_test_candidate(1, "first", 0.90, 7),
similarity_test_candidate(2, "second", 0.80, 7),
],
5,
);
assert_eq!(tracks.len(), 1);
assert_eq!(tracks[0].title, "Track 1");
}
#[test]
fn merging_search_results_deduplicates_source_keys() {
let artist = Artist {