feat(library): cluster merge, resolve, and search commands (server cluster step 2)

Add library_cluster_list_tracks, library_cluster_resolve_candidates, and
library_search_cluster with priority-ordered dedup by cluster_key and
duration guard. Register Tauri commands; no frontend changes in this step.
This commit is contained in:
Maxim Isaev
2026-06-05 16:22:10 +03:00
parent dbccbea6cc
commit f3d4a06726
9 changed files with 800 additions and 2 deletions
@@ -20,7 +20,9 @@ use crate::cover_resolve::CoverEntryDto;
use crate::cross_server;
use crate::dto::{
count_local_tracks, local_tracks_max_updated_ms, track_index_nonempty, ArtifactInputDto,
FactInputDto, LibraryAdvancedSearchRequest, LibraryAdvancedSearchResponse,
FactInputDto, LibraryAdvancedSearchRequest, LibraryAdvancedSearchResponse,
LibraryClusterListTracksRequest, LibraryClusterResolveRequest,
LibraryClusterResolveResponse,
LibraryCrossServerSearchResponse, LibraryLiveSearchRequest, LibraryLiveSearchResponse, LibraryTrackDto,
LibraryTracksEnvelope, OfflinePathDto, PlaySessionDayDetailDto, PlaySessionHeatmapDayDto,
PlaySessionInputDto, PlaySessionRecentDayDto, PlaySessionYearBoundsDto, PlaySessionYearSummaryDto, PurgeReportDto, SyncJobDto, SyncStateDto,
@@ -546,6 +548,80 @@ pub async fn library_search_cross_server(
cross_server::run_cross_server_search(&runtime.store, &query, limit, servers.as_deref())
}
#[tauri::command]
pub async fn library_cluster_list_tracks(
runtime: State<'_, LibraryRuntime>,
request: LibraryClusterListTracksRequest,
) -> Result<LibraryTracksEnvelope, String> {
let store = Arc::clone(&runtime.store);
let servers_ordered = request.servers_ordered;
let limit = request.limit.unwrap_or(100);
let offset = request.offset.unwrap_or(0);
library_spawn_blocking(move || {
crate::server_cluster::list_merged_tracks(&store, &servers_ordered, limit, offset)
})
.await
}
#[tauri::command]
pub async fn library_cluster_resolve_candidates(
runtime: State<'_, LibraryRuntime>,
request: LibraryClusterResolveRequest,
) -> Result<LibraryClusterResolveResponse, String> {
let store = Arc::clone(&runtime.store);
library_spawn_blocking(move || {
if let Some(key) = request.cluster_key.filter(|k| !k.is_empty()) {
let candidates = crate::server_cluster::resolve_candidates_by_cluster_key(
&store,
&request.servers_ordered,
&key,
)?;
return Ok(LibraryClusterResolveResponse {
candidates,
cluster_key: Some(key),
});
}
let server_id = request
.server_id
.as_deref()
.filter(|s| !s.is_empty())
.ok_or_else(|| "cluster_key or (server_id, track_id) required".to_string())?;
let track_id = request
.track_id
.as_deref()
.filter(|s| !s.is_empty())
.ok_or_else(|| "cluster_key or (server_id, track_id) required".to_string())?;
let cluster_key =
crate::server_cluster::cluster_key_for_track(&store, server_id, track_id)?;
let candidates = crate::server_cluster::resolve_candidates_for_track(
&store,
&request.servers_ordered,
server_id,
track_id,
)?;
Ok(LibraryClusterResolveResponse {
candidates,
cluster_key,
})
})
.await
}
#[tauri::command]
pub async fn library_search_cluster(
runtime: State<'_, LibraryRuntime>,
query: String,
limit: Option<u32>,
servers_ordered: Vec<String>,
) -> Result<LibraryCrossServerSearchResponse, String> {
let store = Arc::clone(&runtime.store);
let limit = limit.unwrap_or(100);
library_spawn_blocking(move || {
crate::server_cluster::run_cluster_search(&store, &query, limit, &servers_ordered)
})
.await
}
// ── helpers ──────────────────────────────────────────────────────────
fn hydrate_refs(
@@ -658,6 +658,51 @@ pub struct LibraryCrossServerSearchResponse {
pub servers_searched: Vec<String>,
}
/// Cluster candidate row for playback / write fan-out resolution.
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(rename_all = "camelCase")]
pub struct LibraryClusterCandidateDto {
pub server_id: String,
pub track_id: String,
pub duration_sec: i64,
pub priority_rank: u32,
pub is_winner: bool,
}
/// `library_cluster_list_tracks` request.
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(rename_all = "camelCase")]
pub struct LibraryClusterListTracksRequest {
/// Ordered member server ids (index 0 = highest priority).
pub servers_ordered: Vec<String>,
#[serde(default)]
pub limit: Option<u32>,
#[serde(default)]
pub offset: Option<u32>,
}
/// `library_cluster_resolve_candidates` request — provide cluster_key OR seed track.
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(rename_all = "camelCase")]
pub struct LibraryClusterResolveRequest {
pub servers_ordered: Vec<String>,
#[serde(default)]
pub cluster_key: Option<String>,
#[serde(default)]
pub server_id: Option<String>,
#[serde(default)]
pub track_id: Option<String>,
}
/// `library_cluster_resolve_candidates` response.
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(rename_all = "camelCase")]
pub struct LibraryClusterResolveResponse {
pub candidates: Vec<LibraryClusterCandidateDto>,
#[serde(default)]
pub cluster_key: Option<String>,
}
/// Read `MAX(server_updated_at)` for non-deleted tracks on this server
/// — used by `SyncStateDto` so callers can show "tracks watermark" in
/// Settings without a separate column.
@@ -0,0 +1,199 @@
//! Merged track listing for cluster scope (spec §4 Tier 1).
use rusqlite::types::Value as SqlValue;
use crate::dto::{LibraryTrackDto, LibraryTracksEnvelope};
use crate::repos;
use crate::search::{aliased_track_columns, PAGE_LIMIT_MAX};
use crate::store::LibraryStore;
use super::db::ATTACH_ALIAS;
use super::merge::DURATION_TOLERANCE_SEC;
use super::priority::{in_list_sql, priority_case_sql};
/// List merged tracks — one row per `cluster_key` (priority winner), solo rows
/// for empty-key tracks and duration outliers.
pub fn list_merged_tracks(
store: &LibraryStore,
servers_ordered: &[String],
limit: u32,
offset: u32,
) -> Result<LibraryTracksEnvelope, String> {
if servers_ordered.is_empty() {
return Ok(LibraryTracksEnvelope {
tracks: vec![],
total: 0,
});
}
let limit = limit.clamp(1, PAGE_LIMIT_MAX);
let offset = offset.min(i32::MAX as u32) as i32;
let (in_placeholders, mut in_params) = in_list_sql(servers_ordered);
let (priority_sql, mut priority_params) = priority_case_sql("t.server_id", servers_ordered);
let cols = aliased_track_columns("t");
let sql = format!(
"WITH candidates AS (
SELECT
t.rowid AS tid,
t.server_id,
t.id AS track_id,
k.cluster_key,
COALESCE(k.duration_sec, t.duration_sec) AS dur,
({priority_sql}) AS priority_rank
FROM track t
LEFT JOIN {ATTACH_ALIAS}.track_cluster_key k
ON k.server_id = t.server_id AND k.track_id = t.id
WHERE t.deleted = 0 AND t.server_id IN ({in_placeholders})
),
refs AS (
SELECT cluster_key, MIN(priority_rank) AS best_rank
FROM candidates
WHERE cluster_key IS NOT NULL
GROUP BY cluster_key
),
ref_dur AS (
SELECT c.cluster_key, c.dur AS ref_dur
FROM candidates c
JOIN refs r ON c.cluster_key = r.cluster_key AND c.priority_rank = r.best_rank
),
partitioned AS (
SELECT c.tid,
CASE
WHEN c.cluster_key IS NULL THEN 'solo:' || c.server_id || ':' || c.track_id
WHEN ABS(c.dur - rd.ref_dur) <= {tol} THEN c.cluster_key
ELSE 'solo:' || c.server_id || ':' || c.track_id
END AS merge_key,
c.priority_rank
FROM candidates c
LEFT JOIN ref_dur rd ON c.cluster_key = rd.cluster_key
),
winners AS (
SELECT tid,
ROW_NUMBER() OVER (PARTITION BY merge_key ORDER BY priority_rank) AS rn
FROM partitioned
)
SELECT {cols}
FROM winners w
JOIN track t ON t.rowid = w.tid
WHERE w.rn = 1
ORDER BY t.title COLLATE NOCASE, t.server_id, t.id
LIMIT ? OFFSET ?",
tol = DURATION_TOLERANCE_SEC,
);
let mut params: Vec<SqlValue> = Vec::new();
params.append(&mut priority_params);
params.append(&mut in_params);
params.push(SqlValue::Integer(limit as i64));
params.push(SqlValue::Integer(offset as i64));
let tracks: Vec<LibraryTrackDto> = store.with_read_conn(|conn| {
let mut stmt = conn.prepare(&sql)?;
let rows = stmt.query_map(rusqlite::params_from_iter(params.iter()), |r| {
repos::row_to_track_row(r).map(|row| LibraryTrackDto::from_row(&row))
})?;
rows.collect::<rusqlite::Result<Vec<_>>>()
})?;
Ok(LibraryTracksEnvelope {
total: tracks.len() as u32,
tracks,
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::repos::{TrackRepository, TrackRow};
use crate::server_cluster::rebuild::rebuild_all_cluster_keys;
fn track(server: &str, id: &str, title: &str, artist: &str, album: &str, dur: i64) -> TrackRow {
TrackRow {
server_id: server.into(),
id: id.into(),
title: title.into(),
title_sort: None,
artist: Some(artist.into()),
artist_id: None,
album: album.into(),
album_id: None,
album_artist: Some(artist.into()),
duration_sec: dur,
track_number: Some(1),
disc_number: Some(1),
year: None,
genre: None,
suffix: None,
bit_rate: None,
size_bytes: None,
cover_art_id: None,
starred_at: None,
user_rating: None,
play_count: None,
played_at: None,
server_path: None,
library_id: None,
isrc: None,
mbid_recording: None,
bpm: None,
replay_gain_track_db: None,
replay_gain_album_db: None,
content_hash: None,
server_updated_at: None,
server_created_at: None,
deleted: false,
synced_at: 1,
raw_json: "{}".into(),
}
}
#[test]
fn merge_collapses_same_cluster_key_by_priority() {
let store = LibraryStore::open_in_memory();
TrackRepository::new(&store)
.upsert_batch(&[
track("s1", "t1", "Song", "Band", "LP", 200),
track("s2", "t2", "Song", "Band", "LP", 201),
])
.unwrap();
rebuild_all_cluster_keys(&store).unwrap();
let env = list_merged_tracks(&store, &["s1".into(), "s2".into()], 50, 0).unwrap();
assert_eq!(env.tracks.len(), 1);
assert_eq!(env.tracks[0].server_id, "s1");
}
#[test]
fn unavailable_priority_falls_through() {
let store = LibraryStore::open_in_memory();
TrackRepository::new(&store)
.upsert_batch(&[
track("s1", "t1", "Song", "Band", "LP", 200),
track("s2", "t2", "Song", "Band", "LP", 201),
])
.unwrap();
rebuild_all_cluster_keys(&store).unwrap();
let env = list_merged_tracks(&store, &["s2".into()], 50, 0).unwrap();
assert_eq!(env.tracks.len(), 1);
assert_eq!(env.tracks[0].server_id, "s2");
}
#[test]
fn empty_key_tracks_never_merge() {
let store = LibraryStore::open_in_memory();
store
.with_conn_mut("misc", |c| {
c.execute(
"INSERT INTO track (server_id, id, title, artist, album, duration_sec, synced_at, raw_json) \
VALUES ('s1', 't1', 'A', '', 'X', 1, 1, '{}'), \
('s2', 't2', 'A', '', 'X', 1, 1, '{}')",
[],
)
})
.unwrap();
rebuild_all_cluster_keys(&store).unwrap();
let env = list_merged_tracks(&store, &["s1".into(), "s2".into()], 50, 0).unwrap();
assert_eq!(env.tracks.len(), 2);
}
}
@@ -0,0 +1,64 @@
//! Duration guard and partition keys for cluster merge (spec §2.3).
pub const DURATION_TOLERANCE_SEC: i64 = 5;
/// Synthetic partition for tracks without a `cluster_key` row (never merged).
pub fn solo_partition_key(server_id: &str, track_id: &str) -> String {
format!("solo:{server_id}:{track_id}")
}
/// Within one `cluster_key` group, split rows that fall outside ± tolerance of
/// the reference (priority-1 available candidate duration). Returns partition
/// keys: merged survivors share `cluster_key`; outliers get solo keys.
pub fn duration_partitions(
cluster_key: &str,
rows: &[(String, String, i64, u32)],
) -> Vec<(String, String, String)> {
// (server_id, track_id, duration_sec, priority_rank)
if rows.is_empty() {
return Vec::new();
}
let mut sorted = rows.to_vec();
sorted.sort_by_key(|(_, _, _, rank)| *rank);
let reference_duration = sorted[0].2;
let mut merged: Vec<&(String, String, i64, u32)> = Vec::new();
let mut outliers: Vec<&(String, String, i64, u32)> = Vec::new();
for row in &sorted {
if (row.2 - reference_duration).abs() <= DURATION_TOLERANCE_SEC {
merged.push(row);
} else {
outliers.push(row);
}
}
let mut out = Vec::new();
if !merged.is_empty() {
merged.sort_by_key(|(_, _, _, rank)| *rank);
let (sid, tid, _, _) = merged[0];
out.push((cluster_key.to_string(), sid.clone(), tid.clone()));
}
for (sid, tid, _, _) in outliers {
out.push((solo_partition_key(sid, tid), sid.clone(), tid.clone()));
}
out
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn outlier_splits_to_solo_partition() {
let rows = vec![
("s1".into(), "t1".into(), 180, 0),
("s2".into(), "t2".into(), 182, 1),
("s3".into(), "t3".into(), 240, 2),
];
let parts = duration_partitions("ck1", &rows);
assert_eq!(parts.len(), 2);
assert_eq!(parts[0].0, "ck1");
assert_eq!(parts[0].1, "s1");
assert_eq!(parts[1].0, "solo:s3:t3");
}
}
@@ -4,15 +4,26 @@
mod db;
mod keys;
mod list;
mod merge;
mod norm;
mod priority;
mod rebuild;
mod resolve;
mod search;
pub use db::{
attach_cluster_database, attach_cluster_database_uri, cluster_db_path, ensure_cluster_schema,
init_cluster_meta, needs_norm_rebuild, ATTACH_ALIAS, CLUSTER_DB_FILENAME, NORM_VERSION,
};
pub use keys::{TrackClusterKeys, compute_track_cluster_keys};
pub use keys::{compute_track_cluster_keys, TrackClusterKeys};
pub use list::list_merged_tracks;
pub use merge::DURATION_TOLERANCE_SEC;
pub use norm::norm_field;
pub use rebuild::{
rebuild_all_cluster_keys, rebuild_cluster_keys_for_server, rebuild_if_norm_version_stale,
};
pub use resolve::{
cluster_key_for_track, resolve_candidates_by_cluster_key, resolve_candidates_for_track,
};
pub use search::run_cluster_search;
@@ -0,0 +1,44 @@
//! Priority rank SQL from an ordered server list (index 0 = highest).
use rusqlite::types::Value as SqlValue;
/// Build `CASE server_col WHEN ? THEN 0 … ELSE 9999 END` plus bind values.
pub fn priority_case_sql(server_col: &str, servers_ordered: &[String]) -> (String, Vec<SqlValue>) {
if servers_ordered.is_empty() {
return ("9999".to_string(), Vec::new());
}
let mut sql = format!("CASE {server_col}");
let mut params = Vec::with_capacity(servers_ordered.len());
for (rank, sid) in servers_ordered.iter().enumerate() {
sql.push_str(&format!(" WHEN ? THEN {rank}"));
params.push(SqlValue::Text(sid.clone()));
}
sql.push_str(" ELSE 9999 END");
(sql, params)
}
/// `server_id IN (?,?,…)` placeholders and bind values.
pub fn in_list_sql(servers: &[String]) -> (String, Vec<SqlValue>) {
if servers.is_empty() {
return ("0".to_string(), Vec::new());
}
let placeholders = vec!["?"; servers.len()].join(", ");
let params = servers
.iter()
.map(|s| SqlValue::Text(s.clone()))
.collect();
(placeholders, params)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn priority_case_orders_servers() {
let (sql, params) = priority_case_sql("t.server_id", &["a".into(), "b".into()]);
assert!(sql.contains("WHEN ? THEN 0"));
assert!(sql.contains("WHEN ? THEN 1"));
assert_eq!(params.len(), 2);
}
}
@@ -0,0 +1,198 @@
//! Resolve cluster candidates for playback / writes (spec §56).
use rusqlite::types::Value as SqlValue;
use rusqlite::OptionalExtension;
use crate::dto::LibraryClusterCandidateDto;
use crate::store::LibraryStore;
use super::db::ATTACH_ALIAS;
use super::merge::duration_partitions;
use super::priority::priority_case_sql;
/// All `(server_id, track_id)` rows sharing a `cluster_key`, ordered by priority.
pub fn resolve_candidates_by_cluster_key(
store: &LibraryStore,
servers_ordered: &[String],
cluster_key: &str,
) -> Result<Vec<LibraryClusterCandidateDto>, String> {
if servers_ordered.is_empty() || cluster_key.is_empty() {
return Ok(Vec::new());
}
let (priority_sql, mut priority_params) = priority_case_sql("t.server_id", servers_ordered);
let in_placeholders = vec!["?"; servers_ordered.len()].join(", ");
let mut in_params: Vec<SqlValue> = servers_ordered
.iter()
.map(|s| SqlValue::Text(s.clone()))
.collect();
let sql = format!(
"SELECT t.server_id, t.id, COALESCE(k.duration_sec, t.duration_sec), ({priority_sql})
FROM {ATTACH_ALIAS}.track_cluster_key k
JOIN track t ON t.server_id = k.server_id AND t.id = k.track_id
WHERE k.cluster_key = ? AND t.deleted = 0 AND t.server_id IN ({in_placeholders})
ORDER BY 4, t.server_id, t.id"
);
let mut params: Vec<SqlValue> = Vec::new();
params.append(&mut priority_params);
params.push(SqlValue::Text(cluster_key.to_string()));
params.append(&mut in_params);
let rows: Vec<(String, String, i64, u32)> = store.with_read_conn(|conn| {
let mut stmt = conn.prepare(&sql)?;
let collected = stmt.query_map(rusqlite::params_from_iter(params.iter()), |r| {
Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get::<_, i64>(3)? as u32))
})?;
collected.collect::<rusqlite::Result<Vec<_>>>()
})?;
let with_rank = rows;
let partitions = duration_partitions(cluster_key, &with_rank);
let winner = partitions.first();
Ok(with_rank
.into_iter()
.map(|(server_id, track_id, duration_sec, priority_rank)| {
let is_winner = winner
.map(|(_, ws, wt)| ws == &server_id && wt == &track_id)
.unwrap_or(false);
LibraryClusterCandidateDto {
server_id,
track_id,
duration_sec,
priority_rank,
is_winner,
}
})
.collect())
}
/// Resolve `cluster_key` from a seed track, then return candidates.
pub fn resolve_candidates_for_track(
store: &LibraryStore,
servers_ordered: &[String],
server_id: &str,
track_id: &str,
) -> Result<Vec<LibraryClusterCandidateDto>, String> {
let cluster_key: Option<String> = store.with_read_conn(|conn| {
conn.query_row(
&format!(
"SELECT cluster_key FROM {ATTACH_ALIAS}.track_cluster_key \
WHERE server_id = ?1 AND track_id = ?2"
),
rusqlite::params![server_id, track_id],
|r| r.get(0),
)
.optional()
})?;
let Some(cluster_key) = cluster_key else {
return Ok(vec![LibraryClusterCandidateDto {
server_id: server_id.to_string(),
track_id: track_id.to_string(),
duration_sec: track_duration(store, server_id, track_id)?,
priority_rank: servers_ordered
.iter()
.position(|s| s == server_id)
.unwrap_or(9999) as u32,
is_winner: true,
}]);
};
resolve_candidates_by_cluster_key(store, servers_ordered, &cluster_key)
}
fn track_duration(store: &LibraryStore, server_id: &str, track_id: &str) -> Result<i64, String> {
store
.with_read_conn(|conn| {
conn.query_row(
"SELECT duration_sec FROM track WHERE server_id = ?1 AND id = ?2 AND deleted = 0",
rusqlite::params![server_id, track_id],
|r| r.get(0),
)
})
.map_err(|e| e.to_string())
}
/// Lookup cluster key for a track (for search/seed mapping).
pub fn cluster_key_for_track(
store: &LibraryStore,
server_id: &str,
track_id: &str,
) -> Result<Option<String>, String> {
store
.with_read_conn(|conn| {
conn.query_row(
&format!(
"SELECT cluster_key FROM {ATTACH_ALIAS}.track_cluster_key \
WHERE server_id = ?1 AND track_id = ?2"
),
rusqlite::params![server_id, track_id],
|r| r.get(0),
)
.optional()
})
.map_err(|e| e.to_string())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::repos::{TrackRepository, TrackRow};
use crate::server_cluster::rebuild::rebuild_all_cluster_keys;
fn tr(server: &str, id: &str) -> TrackRow {
TrackRow {
server_id: server.into(),
id: id.into(),
title: "Song".into(),
title_sort: None,
artist: Some("Band".into()),
artist_id: None,
album: "LP".into(),
album_id: None,
album_artist: Some("Band".into()),
duration_sec: 180,
track_number: Some(1),
disc_number: Some(1),
year: None,
genre: None,
suffix: None,
bit_rate: None,
size_bytes: None,
cover_art_id: None,
starred_at: None,
user_rating: None,
play_count: None,
played_at: None,
server_path: None,
library_id: None,
isrc: None,
mbid_recording: None,
bpm: None,
replay_gain_track_db: None,
replay_gain_album_db: None,
content_hash: None,
server_updated_at: None,
server_created_at: None,
deleted: false,
synced_at: 1,
raw_json: "{}".into(),
}
}
#[test]
fn resolve_orders_by_priority() {
let store = LibraryStore::open_in_memory();
TrackRepository::new(&store)
.upsert_batch(&[tr("s1", "t1"), tr("s2", "t2")])
.unwrap();
rebuild_all_cluster_keys(&store).unwrap();
let key = cluster_key_for_track(&store, "s1", "t1").unwrap().unwrap();
let cands = resolve_candidates_by_cluster_key(&store, &["s1".into(), "s2".into()], &key).unwrap();
assert_eq!(cands.len(), 2);
assert!(cands[0].is_winner);
assert_eq!(cands[0].server_id, "s1");
}
}
@@ -0,0 +1,158 @@
//! Cluster-mode cross-server search — dedup by `cluster_key` + priority (spec §5).
use std::collections::HashSet;
use rusqlite::types::Value as SqlValue;
use crate::dto::{LibraryCrossServerSearchResponse, LibraryTrackDto};
use crate::repos;
use crate::search::{aliased_track_columns, fts_query, like_contains, PAGE_LIMIT_MAX};
use crate::store::LibraryStore;
use super::db::ATTACH_ALIAS;
use super::merge::solo_partition_key;
use super::priority::{in_list_sql, priority_case_sql};
const FUZZY_PER_SERVER_CAP: usize = 20;
/// FTS union over ordered servers; dedup by `cluster_key` (not canonical id).
pub fn run_cluster_search(
store: &LibraryStore,
query: &str,
limit: u32,
servers_ordered: &[String],
) -> Result<LibraryCrossServerSearchResponse, String> {
let limit = limit.clamp(1, PAGE_LIMIT_MAX);
if servers_ordered.is_empty() {
return Ok(LibraryCrossServerSearchResponse::default());
}
let Some(fts) = fts_query(query) else {
return Ok(LibraryCrossServerSearchResponse::default());
};
let (in_placeholders, mut in_params) = in_list_sql(servers_ordered);
let (priority_sql, mut priority_params) = priority_case_sql("t.server_id", servers_ordered);
let cols = aliased_track_columns("t");
let canonical_idx = repos::track_columns().split(',').count();
let sql = format!(
"SELECT {cols}, k.cluster_key, ({priority_sql}) AS priority_rank \
FROM track_fts f \
JOIN track t ON t.rowid = f.rowid \
LEFT JOIN {ATTACH_ALIAS}.track_cluster_key k \
ON k.server_id = t.server_id AND k.track_id = t.id \
WHERE track_fts MATCH ? AND t.deleted = 0 AND t.server_id IN ({in_placeholders}) \
ORDER BY bm25(track_fts) LIMIT ?"
);
let mut params: Vec<SqlValue> = Vec::new();
params.push(SqlValue::Text(fts));
params.append(&mut priority_params);
params.append(&mut in_params);
params.push(SqlValue::Integer((limit as i64).saturating_mul(4)));
let rows: Vec<(LibraryTrackDto, Option<String>, u32)> = store.with_read_conn(|conn| {
let mut stmt = conn.prepare(&sql)?;
let collected = stmt.query_map(rusqlite::params_from_iter(params.iter()), |r| {
let track = repos::row_to_track_row(r).map(|row| LibraryTrackDto::from_row(&row))?;
let cluster_key: Option<String> = r.get(canonical_idx)?;
let rank: i64 = r.get(canonical_idx + 1)?;
Ok((track, cluster_key, rank as u32))
})?;
collected.collect::<rusqlite::Result<Vec<_>>>()
})?;
let mut best_by_key: std::collections::HashMap<String, (LibraryTrackDto, u32)> =
std::collections::HashMap::new();
for (track, cluster_key, priority_rank) in rows {
let dedup_key = cluster_key
.clone()
.unwrap_or_else(|| solo_partition_key(&track.server_id, &track.id));
match best_by_key.get(&dedup_key) {
Some((_, best_rank)) if *best_rank <= priority_rank => {}
_ => {
best_by_key.insert(dedup_key, (track, priority_rank));
}
}
}
let mut hits: Vec<LibraryTrackDto> = best_by_key.into_values().map(|(t, _)| t).collect();
hits.truncate(limit as usize);
let hit_keys: HashSet<(String, String)> = hits
.iter()
.map(|t| (t.server_id.clone(), t.id.clone()))
.collect();
let fuzzy = fuzzy_cluster_matches(
store,
servers_ordered,
query.trim(),
&hit_keys,
limit as usize,
)?;
Ok(LibraryCrossServerSearchResponse {
hits,
fuzzy,
servers_searched: servers_ordered.to_vec(),
})
}
fn fuzzy_cluster_matches(
store: &LibraryStore,
targets: &[String],
query: &str,
hit_keys: &HashSet<(String, String)>,
overall_cap: usize,
) -> Result<Vec<LibraryTrackDto>, String> {
let like = like_contains(query);
let cols = aliased_track_columns("t");
let key_idx = repos::track_columns().split(',').count();
let sql = format!(
"SELECT {cols}, k.cluster_key \
FROM track t \
LEFT JOIN {ATTACH_ALIAS}.track_cluster_key k \
ON k.server_id = t.server_id AND k.track_id = t.id \
WHERE t.server_id = ? AND t.deleted = 0 AND t.title LIKE ? ESCAPE '\\' \
ORDER BY t.title COLLATE NOCASE ASC LIMIT ?"
);
let mut out: Vec<LibraryTrackDto> = Vec::new();
let mut seen_keys: HashSet<String> = HashSet::new();
for server in targets {
if out.len() >= overall_cap {
break;
}
let bound = [
SqlValue::Text(server.clone()),
SqlValue::Text(like.clone()),
SqlValue::Integer(FUZZY_PER_SERVER_CAP as i64),
];
let rows: Vec<(LibraryTrackDto, Option<String>)> = store.with_read_conn(|conn| {
let mut stmt = conn.prepare(&sql)?;
let collected = stmt.query_map(rusqlite::params_from_iter(bound.iter()), |r| {
let track = repos::row_to_track_row(r).map(|row| LibraryTrackDto::from_row(&row))?;
let cluster_key: Option<String> = r.get(key_idx)?;
Ok((track, cluster_key))
})?;
collected.collect::<rusqlite::Result<Vec<_>>>()
})?;
for (track, cluster_key) in rows {
if out.len() >= overall_cap {
break;
}
if hit_keys.contains(&(track.server_id.clone(), track.id.clone())) {
continue;
}
let dedup_key = cluster_key
.clone()
.unwrap_or_else(|| solo_partition_key(&track.server_id, &track.id));
if !seen_keys.insert(dedup_key) {
continue;
}
out.push(track);
}
}
Ok(out)
}
+3
View File
@@ -712,6 +712,9 @@ pub fn run() {
psysonic_library::commands::library_list_albums_by_genre,
psysonic_library::commands::library_get_artist_lossless_browse,
psysonic_library::commands::library_search_cross_server,
psysonic_library::commands::library_cluster_list_tracks,
psysonic_library::commands::library_cluster_resolve_candidates,
psysonic_library::commands::library_search_cluster,
psysonic_library::commands::library_get_track,
psysonic_library::commands::library_get_tracks_batch,
psysonic_library::commands::library_get_tracks_by_album,