fix(library): resolve multi-server correctness blockers

This commit is contained in:
cucadmuh
2026-07-21 04:06:58 +03:00
parent ebbbc39014
commit 4c030bfb91
18 changed files with 1002 additions and 120 deletions
@@ -4099,7 +4099,7 @@ mod tests {
}
#[test]
fn multi_scope_text_fts_dedupes_tracks() {
fn multi_scope_text_fts_preserves_same_server_track_occurrences() {
let store = LibraryStore::open_in_memory();
seed_and_rebuild(
&store,
@@ -4134,8 +4134,9 @@ mod tests {
r.library_scopes = Some(vec![scope_pair("s1", "lib-a"), scope_pair("s1", "lib-b")]);
r.query = Some("aurora".into());
let resp = run_advanced_search(&store, &r).unwrap();
assert_eq!(resp.tracks.len(), 1);
assert_eq!(resp.tracks.len(), 2);
assert_eq!(resp.tracks[0].id, "t-a");
assert_eq!(resp.tracks[1].id, "t-b");
}
#[test]
@@ -395,11 +395,12 @@ fn inspect_album(store: &LibraryStore) -> Result<ScopeBrowseProjectionInspectDto
pub fn inspect(store: &LibraryStore) -> Result<ScopeBrowseProjectionInspectDto, String> {
let album = inspect_album(store)?;
let composer = crate::composer_projection::inspect(store)?;
let identity_needed = crate::identity::identity_maintenance_needed(store)?;
let pending = [album.clone(), composer.clone()]
.into_iter()
.filter(|item| item.needed)
.collect::<Vec<_>>();
if pending.is_empty() {
if pending.is_empty() && !identity_needed {
return Ok(ScopeBrowseProjectionInspectDto {
needed: false,
total_tracks: album.total_tracks.max(composer.total_tracks),
@@ -408,8 +409,16 @@ pub fn inspect(store: &LibraryStore) -> Result<ScopeBrowseProjectionInspectDto,
}
Ok(ScopeBrowseProjectionInspectDto {
needed: true,
total_tracks: pending.iter().map(|item| item.total_tracks).max().unwrap_or(0),
done_tracks: pending.iter().map(|item| item.done_tracks).min().unwrap_or(0),
total_tracks: pending
.iter()
.map(|item| item.total_tracks)
.max()
.unwrap_or_else(|| album.total_tracks.max(composer.total_tracks)),
done_tracks: if identity_needed {
0
} else {
pending.iter().map(|item| item.done_tracks).min().unwrap_or(0)
},
})
}
@@ -429,11 +438,34 @@ pub fn is_ready(store: &LibraryStore) -> Result<bool, String> {
}
pub fn run_backfill(store: &LibraryStore, app: &AppHandle) -> Result<(), String> {
run_backfill_impl(store, Some(app))?;
crate::composer_projection::run_backfill(store, Some(app))
run_backfill_impl(store, Some(app))
}
fn run_backfill_impl(store: &LibraryStore, app: Option<&AppHandle>) -> Result<(), String> {
// Projection batches intentionally write physical fallback identities. Persist
// a server rebuild request first so a crash can never leave a completed
// projection marker without a later canonical reconcile.
store.with_conn_mut("browse_projection.mark_identity_dirty", |conn| {
let server_ids = {
let mut statement = conn.prepare(
"SELECT DISTINCT server_id FROM track WHERE deleted = 0 ORDER BY server_id",
)?;
let rows = statement
.query_map([], |row| row.get::<_, String>(0))?
.collect::<rusqlite::Result<Vec<_>>>()?;
rows
};
let tx = conn.transaction()?;
crate::identity::mark_cluster_keys_dirty(&tx, server_ids.iter().map(String::as_str))?;
tx.commit()
})?;
run_album_backfill_impl(store, app)?;
crate::composer_projection::run_backfill(store, app)?;
crate::identity::ensure_pending_cluster_keys(store)?;
Ok(())
}
fn run_album_backfill_impl(store: &LibraryStore, app: Option<&AppHandle>) -> Result<(), String> {
let inspect_result = inspect_album(store)?;
if !inspect_result.needed {
return Ok(());
@@ -472,7 +504,17 @@ fn run_backfill_impl(store: &LibraryStore, app: Option<&AppHandle>) -> Result<()
for (_, server_id, library_id, album_id) in rows {
add_scope(&mut scopes, &server_id, Some(library_id), album_id);
}
let server_ids = scopes
.iter()
.map(|(server_id, _, _)| server_id.clone())
.collect::<HashSet<_>>();
refresh_album_scopes(&tx, scopes)?;
// Keep the reconcile request atomic with every physical-key batch.
// An unrelated read may drain an earlier request while migration runs.
crate::identity::mark_cluster_keys_dirty(
&tx,
server_ids.iter().map(String::as_str),
)?;
tx.execute(
"UPDATE library_data_migration SET cursor_rowid = ?2 WHERE id = ?1",
params![MIGRATION_ID, last_rowid],
@@ -776,6 +818,51 @@ mod tests {
assert_eq!(keys.len(), 1);
}
#[test]
fn completed_backfill_reconciles_physical_projection_keys_before_readiness() {
let store = LibraryStore::open_in_memory();
insert_artist(&store, "s1", "artist-1", "Artist");
insert_artist(&store, "s2", "artist-2", "Artist");
TrackRepository::new(&store)
.upsert_batch(&[
album_track(
"s1", "t1", "Artist", "artist-1", "album-1", "Shared", "Artist", "lib-a",
),
album_track(
"s2", "t2", "Artist", "artist-2", "album-2", "Shared", "Artist", "lib-b",
),
])
.unwrap();
crate::identity::rebuild_cluster_keys(&store, None).unwrap();
store
.with_conn_mut("test.reset_projection", |conn| {
conn.execute("DELETE FROM album_browse_projection", [])?;
conn.execute(
"DELETE FROM library_data_migration WHERE id = ?1",
params![MIGRATION_ID],
)?;
Ok(())
})
.unwrap();
run_backfill_impl(&store, None).unwrap();
assert!(is_ready(&store).unwrap());
let keys = store
.with_read_conn(|conn| {
let mut statement = conn.prepare(
"SELECT DISTINCT identity_key FROM album_browse_projection ORDER BY identity_key",
)?;
let rows = statement
.query_map([], |row| row.get::<_, String>(0))?
.collect::<rusqlite::Result<Vec<_>>>()?;
Ok(rows)
})
.unwrap();
assert_eq!(keys.len(), 1);
assert!(!keys[0].starts_with("physical:"));
}
#[test]
fn ordinary_browse_keeps_ambiguous_physical_albums_separate() {
let store = LibraryStore::open_in_memory();
@@ -8,7 +8,7 @@ use rusqlite::Connection;
pub const CLUSTER_SCHEMA: &str = "cluster";
pub const CLUSTER_DB_FILENAME: &str = "library-cluster.db";
const CLUSTER_SCHEMA_VERSION: i64 = 1;
const CLUSTER_SCHEMA_VERSION: i64 = 2;
const CLUSTER_SCHEMA_SQL: &str = "
CREATE TABLE IF NOT EXISTS cluster.track_cluster_key (
@@ -19,6 +19,7 @@ CREATE TABLE IF NOT EXISTS cluster.track_cluster_key (
album_key TEXT,
artist_key TEXT,
duration_sec INTEGER,
occurrence_rank INTEGER NOT NULL DEFAULT 0,
PRIMARY KEY (server_id, track_id)
);
CREATE INDEX IF NOT EXISTS cluster.idx_ck_scope_album
@@ -26,13 +27,13 @@ CREATE INDEX IF NOT EXISTS cluster.idx_ck_scope_album
CREATE INDEX IF NOT EXISTS cluster.idx_ck_scope_artist
ON track_cluster_key(server_id, library_id, artist_key);
CREATE INDEX IF NOT EXISTS cluster.idx_ck_scope_track
ON track_cluster_key(server_id, library_id, cluster_key);
ON track_cluster_key(server_id, library_id, cluster_key, duration_sec, occurrence_rank);
CREATE INDEX IF NOT EXISTS cluster.idx_ck_server_album
ON track_cluster_key(server_id, album_key);
CREATE INDEX IF NOT EXISTS cluster.idx_ck_server_artist
ON track_cluster_key(server_id, artist_key);
CREATE INDEX IF NOT EXISTS cluster.idx_ck_server_track
ON track_cluster_key(server_id, cluster_key);
ON track_cluster_key(server_id, cluster_key, duration_sec, occurrence_rank);
CREATE TABLE IF NOT EXISTS cluster.cluster_meta (
key TEXT PRIMARY KEY,
value TEXT
@@ -338,6 +339,44 @@ mod tests {
assert_eq!(version, CLUSTER_SCHEMA_VERSION);
}
#[test]
fn version_one_sidecar_is_recreated_with_occurrence_rank() {
let directory = TestDirectory::new("version-one");
let library_path = directory.library_path();
let cluster_path = cluster_db_path_for_library(&library_path);
{
let conn = Connection::open(&cluster_path).unwrap();
conn.execute_batch(
"CREATE TABLE track_cluster_key ( \
server_id TEXT NOT NULL, library_id TEXT NOT NULL, track_id TEXT NOT NULL, \
cluster_key TEXT, album_key TEXT, artist_key TEXT, duration_sec INTEGER, \
PRIMARY KEY (server_id, track_id) \
); \
CREATE TABLE cluster_meta (key TEXT PRIMARY KEY, value TEXT); \
INSERT INTO track_cluster_key(server_id, library_id, track_id) \
VALUES ('s1', 'lib', 'old'); \
PRAGMA user_version = 1;",
)
.unwrap();
}
let conn = Connection::open_in_memory().unwrap();
attach_cluster_write_file(&conn, &library_path).unwrap();
let key_count: i64 = conn
.query_row("SELECT COUNT(*) FROM cluster.track_cluster_key", [], |row| row.get(0))
.unwrap();
let rank_columns: i64 = conn
.query_row(
"SELECT COUNT(*) FROM pragma_table_info('track_cluster_key', 'cluster') \
WHERE name = 'occurrence_rank'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(key_count, 0);
assert_eq!(rank_columns, 1);
}
#[test]
fn versioned_sidecar_with_unknown_objects_is_recreated() {
let directory = TestDirectory::new("versioned-incompatible");
@@ -16,11 +16,11 @@ pub(crate) use norm::norm_part;
pub(crate) use invalidation::{record_album_scopes, record_artists, record_tracks};
pub use rebuild::{
cluster_rebuild_needed, ensure_cluster_keys_built, ensure_pending_cluster_keys,
rebuild_cluster_keys,
identity_maintenance_needed, rebuild_cluster_keys,
};
pub(crate) use rebuild::{
concrete_physical_album_key, delete_cluster_keys_for_tracks, mark_cluster_keys_dirty,
prune_cluster_keys_for_scope, refresh_library_ids_for_albums,
concrete_physical_album_key, mark_cluster_keys_dirty, prune_cluster_keys_for_scope,
refresh_library_ids_for_albums,
};
pub use keys::{build_track_cluster_keys, TrackClusterKeys};
@@ -40,20 +40,6 @@ pub(crate) fn mark_cluster_keys_dirty<'a>(
super::invalidation::record_servers(tx, server_ids)
}
pub(crate) fn delete_cluster_keys_for_tracks(
tx: &Transaction<'_>,
server_id: &str,
track_ids: &[String],
) -> rusqlite::Result<()> {
let mut statement = tx.prepare_cached(
"DELETE FROM cluster.track_cluster_key WHERE server_id = ?1 AND track_id = ?2",
)?;
for track_id in track_ids {
statement.execute(params![server_id, track_id])?;
}
Ok(())
}
pub(crate) fn prune_cluster_keys_for_scope(
tx: &Transaction<'_>,
server_id: &str,
@@ -134,6 +120,29 @@ pub fn cluster_rebuild_needed(conn: &Connection) -> rusqlite::Result<bool> {
Ok(stored.as_deref() != Some(NORM_VERSION))
}
pub fn identity_maintenance_needed(store: &LibraryStore) -> Result<bool, String> {
store.with_read_conn(|conn| {
let has_sources: bool = conn.query_row(
"SELECT EXISTS(SELECT 1 FROM track WHERE deleted = 0) \
OR EXISTS(SELECT 1 FROM identity_invalidation)",
[],
|row| row.get(0),
)?;
if !has_sources {
return Ok(false);
}
if cluster_rebuild_needed(conn)? {
return Ok(true);
}
conn.query_row(
"SELECT EXISTS(SELECT 1 FROM identity_invalidation) \
OR EXISTS(SELECT 1 FROM cluster.cluster_meta WHERE key LIKE ?1)",
params![format!("{DIRTY_META_PREFIX}%")],
|row| row.get(0),
)
})
}
fn set_cluster_meta(conn: &Connection) -> rusqlite::Result<()> {
let now = now_unix().to_string();
conn.execute(
@@ -172,6 +181,150 @@ pub(crate) fn concrete_physical_album_key(server_id: &str, album_id: &str) -> St
format!("physical:{}:{server_id}:{album_id}", server_id.len())
}
fn occurrence_rank_order_sql() -> &'static str {
"CASE WHEN t.disc_number IS NULL THEN 1 ELSE 0 END, t.disc_number, \
CASE WHEN t.track_number IS NULL THEN 1 ELSE 0 END, t.track_number, \
COALESCE(t.server_path, ''), ck.track_id"
}
fn recompute_all_occurrence_ranks(
tx: &Transaction<'_>,
server_id: Option<&str>,
) -> rusqlite::Result<()> {
let server_filter = if server_id.is_some() {
" AND ck.server_id = ?1"
} else {
""
};
let sql = format!(
"WITH ranked AS MATERIALIZED ( \
SELECT ck.server_id, ck.track_id, \
ROW_NUMBER() OVER ( \
PARTITION BY ck.server_id, ck.cluster_key, ck.duration_sec / 5 \
ORDER BY {} \
) - 1 AS occurrence_rank \
FROM cluster.track_cluster_key ck \
INNER JOIN track t ON t.server_id = ck.server_id AND t.id = ck.track_id \
WHERE t.deleted = 0 AND ck.cluster_key IS NOT NULL{server_filter} \
) \
UPDATE cluster.track_cluster_key AS ck \
SET occurrence_rank = ranked.occurrence_rank \
FROM ranked \
WHERE ck.server_id = ranked.server_id AND ck.track_id = ranked.track_id \
AND ck.occurrence_rank IS NOT ranked.occurrence_rank",
occurrence_rank_order_sql(),
);
match server_id {
Some(server_id) => {
tx.execute(&sql, params![server_id])?;
tx.execute(
"UPDATE cluster.track_cluster_key SET occurrence_rank = 0 \
WHERE server_id = ?1 AND cluster_key IS NULL AND occurrence_rank != 0",
params![server_id],
)?;
}
None => {
tx.execute(&sql, [])?;
tx.execute(
"UPDATE cluster.track_cluster_key SET occurrence_rank = 0 \
WHERE cluster_key IS NULL AND occurrence_rank != 0",
[],
)?;
}
}
Ok(())
}
fn reset_affected_rank_partitions(tx: &Transaction<'_>) -> rusqlite::Result<()> {
tx.execute_batch(
"CREATE TEMP TABLE IF NOT EXISTS identity_rank_partition ( \
cluster_key TEXT NOT NULL, \
duration_bucket INTEGER NOT NULL, \
PRIMARY KEY (cluster_key, duration_bucket) \
) WITHOUT ROWID; \
DELETE FROM temp.identity_rank_partition;",
)
}
fn capture_invalidated_rank_partitions(
tx: &Transaction<'_>,
server_id: &str,
) -> rusqlite::Result<()> {
tx.execute(
"WITH invalidated_artist AS MATERIALIZED ( \
SELECT entity_id FROM identity_invalidation \
WHERE server_id = ?1 AND kind = 'artist' \
), \
invalidated_album AS MATERIALIZED ( \
SELECT entity_id FROM identity_invalidation \
WHERE server_id = ?1 AND kind = 'album' \
UNION \
SELECT DISTINCT t.album_id FROM invalidated_artist ia \
CROSS JOIN track t \
WHERE t.server_id = ?1 AND t.artist_id = ia.entity_id \
AND t.album_id IS NOT NULL AND t.album_id != '' \
), \
candidate_track AS MATERIALIZED ( \
SELECT entity_id FROM identity_invalidation \
WHERE server_id = ?1 AND kind = 'track' \
UNION \
SELECT t.id FROM invalidated_album ia \
CROSS JOIN track t \
WHERE t.server_id = ?1 AND t.album_id = ia.entity_id \
UNION \
SELECT t.id FROM invalidated_artist ia \
CROSS JOIN track t \
WHERE t.server_id = ?1 AND t.artist_id = ia.entity_id \
) \
INSERT OR IGNORE INTO temp.identity_rank_partition(cluster_key, duration_bucket) \
SELECT ck.cluster_key, ck.duration_sec / 5 \
FROM candidate_track candidate \
INNER JOIN cluster.track_cluster_key ck \
ON ck.server_id = ?1 AND ck.track_id = candidate.entity_id \
WHERE ck.cluster_key IS NOT NULL",
params![server_id],
)?;
Ok(())
}
fn recompute_affected_occurrence_ranks(
tx: &Transaction<'_>,
server_id: &str,
) -> rusqlite::Result<()> {
let sql = format!(
"WITH ranked AS MATERIALIZED ( \
SELECT ck.server_id, ck.track_id, \
ROW_NUMBER() OVER ( \
PARTITION BY ck.server_id, ck.cluster_key, ck.duration_sec / 5 \
ORDER BY {} \
) - 1 AS occurrence_rank \
FROM temp.identity_rank_partition affected \
CROSS JOIN cluster.track_cluster_key ck \
ON ck.server_id = ?1 AND ck.cluster_key = affected.cluster_key \
AND ck.duration_sec / 5 = affected.duration_bucket \
INNER JOIN track t ON t.server_id = ck.server_id AND t.id = ck.track_id \
WHERE t.deleted = 0 \
) \
UPDATE cluster.track_cluster_key AS ck \
SET occurrence_rank = ranked.occurrence_rank \
FROM ranked \
WHERE ck.server_id = ranked.server_id AND ck.track_id = ranked.track_id \
AND ck.occurrence_rank IS NOT ranked.occurrence_rank",
occurrence_rank_order_sql(),
);
tx.execute(&sql, params![server_id])?;
tx.execute(
"UPDATE cluster.track_cluster_key SET occurrence_rank = 0 \
WHERE server_id = ?1 AND cluster_key IS NULL AND track_id IN ( \
SELECT entity_id FROM identity_invalidation \
WHERE server_id = ?1 AND kind = 'track' \
)",
params![server_id],
)?;
tx.execute("DELETE FROM temp.identity_rank_partition", [])?;
Ok(())
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum PendingRebuild<'a> {
All,
@@ -339,6 +492,7 @@ fn rebuild_cluster_keys_on_conn(
[],
)?;
}
recompute_all_occurrence_ranks(&tx, server_id)?;
crate::browse_projection::reconcile_identity_keys(&tx, server_id)?;
match server_id {
Some(server_id) => {
@@ -365,6 +519,8 @@ fn apply_identity_invalidations_on_conn(
server_id: &str,
) -> rusqlite::Result<u64> {
let tx = conn.transaction()?;
reset_affected_rank_partitions(&tx)?;
capture_invalidated_rank_partitions(&tx, server_id)?;
let select = "WITH invalidated_artist AS MATERIALIZED ( \
SELECT entity_id FROM identity_invalidation \
WHERE server_id = ?1 AND kind = 'artist' \
@@ -429,6 +585,8 @@ fn apply_identity_invalidations_on_conn(
drop(statement);
drop(upsert);
capture_invalidated_rank_partitions(&tx, server_id)?;
tx.execute(
"DELETE FROM cluster.track_cluster_key AS ck \
WHERE ck.server_id = ?1 \
@@ -442,6 +600,7 @@ fn apply_identity_invalidations_on_conn(
)",
params![server_id],
)?;
recompute_affected_occurrence_ranks(&tx, server_id)?;
set_cluster_meta(&tx)?;
tx.commit()?;
@@ -694,6 +853,53 @@ mod tests {
assert!(empty_artist.2.is_none());
}
#[test]
fn incremental_tombstone_reranks_remaining_track_occurrence() {
let store = LibraryStore::open_in_memory();
let mut first = track_row(
"s1", "t1", "Tyrion", Some("Narrator"), "Book", Some("Narrator"), 300, "lib",
);
first.track_number = Some(1);
let mut second = first.clone();
second.id = "t2".into();
second.track_number = Some(2);
TrackRepository::new(&store)
.upsert_batch(&[first, second])
.unwrap();
rebuild_cluster_keys(&store, Some("s1")).unwrap();
let before = store
.with_read_conn(|conn| {
let mut statement = conn.prepare(
"SELECT track_id, occurrence_rank FROM cluster.track_cluster_key \
WHERE server_id = 's1' ORDER BY track_id",
)?;
let rows = statement
.query_map([], |row| Ok((row.get::<_, String>(0)?, row.get::<_, i64>(1)?)))?
.collect::<rusqlite::Result<Vec<_>>>()?;
Ok(rows)
})
.unwrap();
assert_eq!(before, vec![("t1".into(), 0), ("t2".into(), 1)]);
TrackRepository::new(&store)
.apply_tombstone_results("s1", "", &[], &["t1".into()])
.unwrap();
ensure_cluster_keys_built(&store, "s1").unwrap();
let after = store
.with_read_conn(|conn| {
conn.query_row(
"SELECT track_id, occurrence_rank FROM cluster.track_cluster_key \
WHERE server_id = 's1'",
[],
|row| Ok((row.get::<_, String>(0)?, row.get::<_, i64>(1)?)),
)
})
.unwrap();
assert_eq!(after, ("t2".into(), 0));
}
#[test]
fn rebuild_uses_canonical_artist_name_for_every_track_with_the_same_artist_id() {
let store = LibraryStore::open_in_memory();
@@ -829,7 +829,7 @@ mod tests {
}
#[test]
fn multi_scope_live_search_dedupes_album_and_artist_with_priority() {
fn multi_scope_live_search_dedupes_album_and_artist_but_preserves_tracks() {
use crate::dto::LibraryScopePair;
use crate::identity::rebuild_cluster_keys;
@@ -902,8 +902,9 @@ mod tests {
assert_eq!(resp.artists[0].id, "ar-a");
assert_eq!(resp.albums.len(), 1);
assert_eq!(resp.albums[0].id, "alb-a");
assert_eq!(resp.tracks.len(), 1);
assert_eq!(resp.tracks.len(), 2);
assert_eq!(resp.tracks[0].id, "t-a");
assert_eq!(resp.tracks[1].id, "t-b");
}
/// Manual: `cargo test -p psysonic-library bench_disk_live_search --release -- --ignored --nocapture`
@@ -352,7 +352,6 @@ impl<'a> TrackRepository<'a> {
params![server_id, track_id],
)?;
}
crate::identity::delete_cluster_keys_for_tracks(&tx, server_id, deleted_ids)?;
crate::identity::record_tracks(
&tx,
deleted_ids.iter().map(|track_id| (server_id, track_id.as_str())),
@@ -14,6 +14,7 @@ use crate::dto::{
LibraryScopeBrowseResponse, LibraryScopePair, LibrarySortClause, LibraryTrackDto,
};
use crate::repos::{row_to_track_row, TrackRow};
use crate::scope_merge::TRACK_CLUSTER_PARTITION_KEY;
use crate::store::LibraryStore;
const CANDIDATE_PAGE_SIZE: usize = 64;
@@ -457,7 +458,7 @@ fn query_track_scope_candidates(
};
let sql = format!(
"SELECT {columns}, CASE WHEN ck.cluster_key IS NOT NULL \
THEN ck.cluster_key || ':' || CAST((ck.duration_sec / 5) AS TEXT) END \
THEN {TRACK_CLUSTER_PARTITION_KEY} END \
FROM track t \
LEFT JOIN cluster.track_cluster_key ck ON ck.server_id = t.server_id AND ck.track_id = t.id \
WHERE t.server_id = ? {library_filter} AND t.deleted = 0 {seek} \
@@ -515,15 +516,14 @@ fn track_identity_priorities(
.map(|_| "?")
.collect::<Vec<_>>()
.join(", ");
let identity_sql = "ck.cluster_key || ':' || CAST((ck.duration_sec / 5) AS TEXT)";
let sql = format!(
"{scope_cte} SELECT {identity_sql}, MIN(scope.pr) \
"{scope_cte} SELECT {TRACK_CLUSTER_PARTITION_KEY}, MIN(scope.pr) \
FROM scoped_track scope \
INNER JOIN track t ON t.rowid = scope.rowid \
INNER JOIN cluster.track_cluster_key ck \
ON ck.server_id = t.server_id AND ck.track_id = t.id \
WHERE t.deleted = 0 AND {identity_sql} IN ({placeholders}) \
GROUP BY {identity_sql}",
WHERE t.deleted = 0 AND {TRACK_CLUSTER_PARTITION_KEY} IN ({placeholders}) \
GROUP BY {TRACK_CLUSTER_PARTITION_KEY}",
);
binds.extend(identities.into_iter().map(SqlValue::Text));
store
@@ -859,7 +859,7 @@ mod tests {
query_track_scope_candidates(&store, &scopes[0], 0, None, 10).unwrap(),
query_track_scope_candidates(&store, &scopes[1], 1, None, 10).unwrap(),
];
assert_eq!(track_identity_priorities(&store, &scopes, &candidates).unwrap().get("same:20"), Some(&0));
assert_eq!(track_identity_priorities(&store, &scopes, &candidates).unwrap().get("same:20:0"), Some(&0));
let first = browse(&store, &track_request(scopes.clone(), 1, None)).unwrap();
assert_eq!(first.tracks.iter().map(|track| track.id.as_str()).collect::<Vec<_>>(), vec!["high-dup"]);
@@ -867,4 +867,30 @@ mod tests {
assert!(second.tracks.is_empty());
assert!(!second.has_more);
}
#[test]
fn same_server_occurrence_ranks_survive_across_cursor_pages() {
let store = LibraryStore::open_in_memory();
insert_track(&store, "s1", "lib-a", "chapter-1", "Tyrion", Some("tyrion"));
insert_track(&store, "s1", "lib-b", "chapter-2", "Tyrion", Some("tyrion"));
store
.with_conn_mut("test.scope_browse.rank", |conn| {
conn.execute(
"UPDATE cluster.track_cluster_key SET occurrence_rank = 1 \
WHERE server_id = 's1' AND track_id = 'chapter-2'",
[],
)?;
Ok(())
})
.unwrap();
let scopes = vec![
LibraryScopePair { server_id: "s1".into(), library_id: Some("lib-a".into()) },
LibraryScopePair { server_id: "s1".into(), library_id: Some("lib-b".into()) },
];
let first = browse(&store, &track_request(scopes.clone(), 1, None)).unwrap();
assert_eq!(first.tracks.iter().map(|track| track.id.as_str()).collect::<Vec<_>>(), vec!["chapter-1"]);
let second = browse(&store, &track_request(scopes, 1, first.next_cursor)).unwrap();
assert_eq!(second.tracks.iter().map(|track| track.id.as_str()).collect::<Vec<_>>(), vec!["chapter-2"]);
}
}
@@ -32,14 +32,22 @@ pub(crate) const ALBUM_DEDUP_KEY: &str = "CASE WHEN ck.album_key IS NOT NULL THE
const ARTIST_DEDUP_KEY: &str = "CASE WHEN ck.artist_key IS NOT NULL THEN ck.artist_key \
ELSE ('null:' || t.server_id || ':' || COALESCE(NULLIF(t.artist_id, ''), t.id)) END";
/// Track dedup: `cluster_key` + a fixed 5-second duration bucket (`duration_sec / 5`).
/// Track dedup combines `cluster_key`, a fixed 5-second duration bucket
/// (`duration_sec / 5`), and a deterministic per-server occurrence rank. The
/// rank preserves repeated same-server tracks while still pairing corresponding
/// copies across servers.
///
/// This is a bucket, not a symmetric ±5 s window: two rips whose durations straddle
/// a bucket edge (e.g. 314 s → bucket 62, 316 s → bucket 63) stay separate, while
/// two up to ~4 s apart inside a bucket merge. Kept as a single GROUP BY key for
/// speed; a true tolerance window would need a self-join. Encoder-padding drift at
/// boundaries is the known trade-off.
pub(crate) const TRACK_CLUSTER_PARTITION_KEY: &str = "ck.cluster_key || ':' \
|| CAST((ck.duration_sec / 5) AS TEXT) || ':' || CAST(ck.occurrence_rank AS TEXT)";
pub(crate) const TRACK_DEDUP_KEY: &str = "CASE WHEN ck.cluster_key IS NOT NULL \
THEN ck.cluster_key || ':' || CAST((ck.duration_sec / 5) AS TEXT) \
|| ':' || CAST(ck.occurrence_rank AS TEXT) \
ELSE ('null:' || t.server_id || ':' || t.id) END";
/// Sortable representative key so a single `MIN()` (SQLite bare-column rule) picks the
@@ -204,7 +212,7 @@ fn keyed_detail_track_source(
detail_key(value) AS (VALUES (?)), \
detail_tracks AS MATERIALIZED ( \
SELECT ck.server_id, ck.library_id, ck.track_id, ck.cluster_key, \
ck.album_key, ck.artist_key, ck.duration_sec, s.pr \
ck.album_key, ck.artist_key, ck.duration_sec, ck.occurrence_rank, s.pr \
FROM exact_scope s \
CROSS JOIN detail_key key \
CROSS JOIN cluster.track_cluster_key ck INDEXED BY {scope_index} \
@@ -213,7 +221,7 @@ fn keyed_detail_track_source(
AND ck.{key_column} = key.value \
UNION ALL \
SELECT ck.server_id, ck.library_id, ck.track_id, ck.cluster_key, \
ck.album_key, ck.artist_key, ck.duration_sec, s.pr \
ck.album_key, ck.artist_key, ck.duration_sec, ck.occurrence_rank, s.pr \
FROM whole_scope s \
CROSS JOIN detail_key key \
CROSS JOIN cluster.track_cluster_key ck INDEXED BY {server_index} \
@@ -1669,13 +1677,13 @@ fn lookup_track_partition(
conn: &rusqlite::Connection,
server_id: &str,
track_id: &str,
) -> rusqlite::Result<Option<(Option<String>, i64)>> {
) -> rusqlite::Result<Option<(Option<String>, i64, i64)>> {
conn.query_row(
"SELECT ck.cluster_key, ck.duration_sec / 5 FROM track t \
"SELECT ck.cluster_key, ck.duration_sec / 5, ck.occurrence_rank FROM track t \
INNER JOIN cluster.track_cluster_key ck ON ck.server_id = t.server_id AND ck.track_id = t.id \
WHERE t.server_id = ? AND t.id = ? AND t.deleted = 0 LIMIT 1",
rusqlite::params![server_id, track_id],
|r| Ok((r.get(0)?, r.get(1)?)),
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
)
.optional()
}
@@ -1701,6 +1709,7 @@ fn fetch_track_sources(
scopes: &[LibraryScopePair],
cluster_key: Option<&str>,
duration_bucket: i64,
occurrence_rank: i64,
anchor_server: &str,
anchor_id: &str,
) -> rusqlite::Result<Vec<LibraryEntitySourceDto>> {
@@ -1711,7 +1720,7 @@ fn fetch_track_sources(
"AND t.server_id = ? AND t.id = ?",
);
let bucket_filter = if cluster_key.is_some() {
"AND ck.duration_sec / 5 = ?"
"AND ck.duration_sec / 5 = ? AND ck.occurrence_rank = ?"
} else {
""
};
@@ -1725,6 +1734,7 @@ fn fetch_track_sources(
if let Some(key) = cluster_key {
binds.push(SqlValue::Text(key.to_string()));
binds.push(SqlValue::Integer(duration_bucket));
binds.push(SqlValue::Integer(occurrence_rank));
} else {
binds.push(SqlValue::Text(anchor_server.to_string()));
binds.push(SqlValue::Text(anchor_id.to_string()));
@@ -1815,7 +1825,7 @@ pub fn resolve_entity_sources(
store.with_scope_detail_read_conn(|conn| match request.entity_type {
LibrarySourceEntityType::Track => {
let Some((cluster_key, duration_bucket)) =
let Some((cluster_key, duration_bucket, occurrence_rank)) =
lookup_track_partition(conn, anchor_server, anchor_id)?
else {
return Ok(Vec::new());
@@ -1825,6 +1835,7 @@ pub fn resolve_entity_sources(
scopes,
cluster_key.as_deref(),
duration_bucket,
occurrence_rank,
anchor_server,
anchor_id,
)
@@ -2995,7 +3006,7 @@ mod tests {
}
#[test]
fn dedup_collapses_same_album_and_priority_winner_flips() {
fn album_merge_preserves_same_server_track_multiplicity_and_priority_winner_flips() {
let store = LibraryStore::open_in_memory();
let rows = [
track(
@@ -3040,8 +3051,8 @@ mod tests {
assert_eq!(albums_a[0].id, "alb-a");
assert_eq!(albums_a[0].year, Some(2001));
assert_eq!(albums_a[0].genre.as_deref(), Some("Rock"));
assert_eq!(albums_a[0].song_count, Some(1));
assert_eq!(albums_a[0].duration_sec, Some(200));
assert_eq!(albums_a[0].song_count, Some(2));
assert_eq!(albums_a[0].duration_sec, Some(400));
let req_b_first = LibraryScopeListRequest {
scopes: vec![scope_pair("s1", "lib-b"), scope_pair("s1", "lib-a")],
@@ -3053,8 +3064,8 @@ mod tests {
assert_eq!(albums_b.len(), 1);
assert_eq!(albums_b[0].id, "alb-b");
assert_eq!(albums_b[0].year, Some(1999));
assert_eq!(albums_b[0].song_count, Some(1));
assert_eq!(albums_b[0].duration_sec, Some(200));
assert_eq!(albums_b[0].song_count, Some(2));
assert_eq!(albums_b[0].duration_sec, Some(400));
}
#[test]
@@ -3148,6 +3159,73 @@ mod tests {
assert_eq!(hits.len(), 2);
}
#[test]
fn same_server_occurrences_survive_and_cross_server_sources_pair_by_rank() {
let store = LibraryStore::open_in_memory();
let mut rows = vec![
track(
"s1", "a1", "Tyrion", Some("Narrator"), "Book", "album-a",
Some("narrator"), 300, "lib-a", None, None, None,
),
track(
"s1", "a2", "Tyrion", Some("Narrator"), "Book", "album-a",
Some("narrator"), 300, "lib-a", None, None, None,
),
track(
"s2", "b1", "Tyrion", Some("Narrator"), "Book", "album-b",
Some("narrator"), 300, "lib-b", None, None, None,
),
track(
"s2", "b2", "Tyrion", Some("Narrator"), "Book", "album-b",
Some("narrator"), 300, "lib-b", None, None, None,
),
track(
"s3", "c1", "Tyrion", Some("Narrator"), "Book", "album-c",
Some("narrator"), 300, "lib-c", None, None, None,
),
];
for (index, row) in rows.iter_mut().enumerate() {
row.track_number = Some((index % 2 + 1) as i64);
row.server_path = Some(format!("chapter-{}.mp3", index % 2 + 1));
}
seed_and_rebuild(&store, &rows);
let scopes = vec![whole_scope("s1"), whole_scope("s2"), whole_scope("s3")];
let detail = album_detail(
&store,
&LibraryScopeAlbumDetailRequest {
scopes: scopes.clone(),
album_id: "album-a".into(),
server_id: "s1".into(),
},
)
.unwrap();
assert_eq!(
detail.tracks.iter().map(|track| track.id.as_str()).collect::<Vec<_>>(),
vec!["a1", "a2"]
);
for (anchor_id, expected_ids) in [
("a1", vec!["a1", "b1", "c1"]),
("a2", vec!["a2", "b2"]),
] {
let sources = resolve_entity_sources(
&store,
&LibraryResolveEntitySourcesRequest {
entity_type: LibrarySourceEntityType::Track,
anchor_server_id: "s1".into(),
anchor_id: anchor_id.into(),
scopes: scopes.clone(),
},
)
.unwrap();
assert_eq!(
sources.iter().map(|source| source.id.as_str()).collect::<Vec<_>>(),
expected_ids
);
}
}
#[test]
fn single_scope_returns_correct_album() {
let store = LibraryStore::open_in_memory();
@@ -3287,7 +3365,7 @@ mod tests {
assert_eq!(detail.album.year, Some(2001));
assert_eq!(detail.album.genre, None);
assert_eq!(detail.album.cover_art_id, None);
assert_eq!(detail.tracks.len(), 1);
assert_eq!(detail.tracks.len(), 2);
}
#[test]
@@ -3933,4 +4011,5 @@ mod tests {
bench("1 lib", &[scopes[0].clone()]);
bench("2 libs", &scopes);
}
}
@@ -208,8 +208,12 @@ pub async fn tag_library_membership(
let untagged = tracks
.count_untagged_tracks(server_id)
.map_err(SyncError::Storage)?;
let cursor = read_tag_cursor(store, server_id)?;
if require_untagged && untagged == 0 {
if let Some(cursor) = cursor.as_ref() {
write_tag_completion(store, server_id, &cursor.folders_hash, 0)?;
}
return Ok(TagReport {
folders_processed: 0,
albums_processed: 0,
@@ -239,7 +243,6 @@ pub async fn tag_library_membership(
folders.sort_by(|a, b| a.id.cmp(&b.id));
let hash = folders_hash(&folders);
let prior = read_tag_state(store, server_id)?;
let cursor = read_tag_cursor(store, server_id)?;
let active_cursor = cursor.as_ref().filter(|cursor| cursor.folders_hash == hash);
if !should_run_tagging_pass(untagged, prior.as_ref(), active_cursor.is_some(), &hash) {
return Ok(TagReport {
@@ -59,6 +59,28 @@ fn test_client(base: &str) -> SubsonicClient {
)
}
async fn mount_bounded_full_album_pages(server: &MockServer) {
for page_index in 0..MAX_ALBUM_LIST_REQUESTS_PER_PASS {
let offset = page_index * ALBUM_PAGE_SIZE;
let albums = (0..ALBUM_PAGE_SIZE)
.map(|index| json!({ "id": format!("album-{}", offset + index), "name": "A" }))
.collect::<Vec<_>>();
Mock::given(method("GET"))
.and(path("/rest/getAlbumList2.view"))
.and(query_param("musicFolderId", "1"))
.and(query_param("offset", offset.to_string()))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"subsonic-response": {
"status": "ok",
"albumList2": { "album": albums }
}
})))
.expect(1)
.mount(server)
.await;
}
}
#[test]
fn folders_hash_is_order_independent() {
let a = vec![
@@ -234,25 +256,7 @@ async fn tag_library_membership_resumes_from_persisted_page_cursor() {
.mount(&server)
.await;
for page_index in 0..MAX_ALBUM_LIST_REQUESTS_PER_PASS {
let offset = page_index * ALBUM_PAGE_SIZE;
let albums = (0..ALBUM_PAGE_SIZE)
.map(|index| json!({ "id": format!("album-{}", offset + index), "name": "A" }))
.collect::<Vec<_>>();
Mock::given(method("GET"))
.and(path("/rest/getAlbumList2.view"))
.and(query_param("musicFolderId", "1"))
.and(query_param("offset", offset.to_string()))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"subsonic-response": {
"status": "ok",
"albumList2": { "album": albums }
}
})))
.expect(1)
.mount(&server)
.await;
}
mount_bounded_full_album_pages(&server).await;
let resume_offset = MAX_ALBUM_LIST_REQUESTS_PER_PASS * ALBUM_PAGE_SIZE;
Mock::given(method("GET"))
.and(path("/rest/getAlbumList2.view"))
@@ -322,3 +326,54 @@ async fn tag_library_membership_resumes_from_persisted_page_cursor() {
.unwrap();
assert!(third.skipped);
}
#[tokio::test(flavor = "multi_thread")]
async fn tag_library_membership_finalizes_cursor_when_no_untagged_tracks_remain() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/rest/getMusicFolders.view"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"subsonic-response": {
"status": "ok",
"musicFolders": { "musicFolder": { "id": 1, "name": "Main" } }
}
})))
.expect(1)
.mount(&server)
.await;
mount_bounded_full_album_pages(&server).await;
let store = LibraryStore::open_in_memory();
TrackRepository::new(&store)
.upsert_batch(&[track_row("srv", "tagged-on-first-page", "album-0")])
.unwrap();
let client = test_client(&server.uri());
let progress = Arc::new(super::super::progress::NoopProgress);
let first = tag_library_membership(
&store,
&client,
"srv",
None,
progress.clone(),
false,
)
.await
.unwrap();
assert!(!first.completed);
assert_eq!(first.tracks_tagged, 1);
assert_eq!(first.untagged_remaining, 0);
assert!(read_tag_cursor(&store, "srv").unwrap().is_some());
let second = tag_library_membership(&store, &client, "srv", None, progress, true)
.await
.unwrap();
assert!(second.skipped);
assert!(second.completed);
assert_eq!(second.untagged_remaining, 0);
assert!(read_tag_cursor(&store, "srv").unwrap().is_none());
let completion = read_tag_state(&store, "srv").unwrap().unwrap();
assert_eq!(completion.folders_hash, "1:Main");
assert_eq!(completion.last_untagged_count, 0);
}
@@ -676,6 +676,7 @@ mod tests {
.reconcile_chunk(10)
.await
.unwrap();
crate::identity::ensure_cluster_keys_built(&store, "s1").unwrap();
let (projection, identity): (i64, i64) = store
.with_read_conn(|conn| {
@@ -0,0 +1,211 @@
import type { ReactNode } from 'react';
import { act, screen, waitFor } from '@testing-library/react';
import { beforeEach, describe, expect, it, vi } from 'vitest';
import type { HomeFeedSnapshot } from '@/features/home/store/homeFeedCache';
const homeMocks = vi.hoisted(() => ({
connection: { status: 'checking' as 'checking' | 'connected' | 'disconnected' },
loadHomeFeedWithStatus: vi.fn(),
loadHomeChronologicalFeed: vi.fn(),
unavailableServerIds: new Set<string>(),
}));
vi.mock('@/features/album', () => ({
AlbumRow: () => null,
LosslessAlbumsRail: () => null,
}));
vi.mock('@/features/home/components/Hero', () => ({
default: ({ albums }: { albums: Array<{ name: string }> }) => (
<div data-testid="home-hero">{albums.map(album => album.name).join(',')}</div>
),
}));
vi.mock('@/features/home/components/SongRail', () => ({ default: () => null }));
vi.mock('@/features/home/components/BecauseYouLikeRail', () => ({ default: () => null }));
vi.mock('@/features/home/components/MainstageDiagnosticFrame', () => ({
default: ({ children }: { children: ReactNode }) => children,
}));
vi.mock('@/features/playback/utils/mixRatingFilter', () => ({
filterAlbumsByMixRatingsAcrossServers: vi.fn(async albums => albums),
getMixMinRatingsConfigFromAuth: () => ({
enabled: false, minSong: 0, minAlbum: 0, minArtist: 0,
}),
}));
vi.mock('@/lib/perf/perfFlags', () => ({
usePerfProbeFlags: () => ({
disableMainstageRails: true,
disableHomeAlbumRows: false,
disableHomeSongRails: false,
disableMainstageRailArtwork: true,
disableHomeRailArtwork: false,
disableMainstageHero: false,
disableMainstageGridCards: true,
disableHomeArtworkFx: false,
disableHomeArtworkClip: false,
}),
}));
vi.mock('@/lib/perf/psyLabDebugTraces', () => ({
usePsyLabDebugTraces: () => ({ mainstage: false }),
}));
vi.mock('@/lib/perf/perfTelemetry', () => ({ bumpPerfCounter: vi.fn() }));
vi.mock('@/cover/useLibraryCoverPrefetch', () => ({ useLibraryCoverPrefetch: vi.fn() }));
vi.mock('@/cover/warmDiskPeek', () => ({
primeAlbumCoversForDisplay: vi.fn(async () => undefined),
warmHomeMainstageCovers: vi.fn(async () => undefined),
}));
vi.mock('@/features/home/store/becauseYouLikeCache', () => ({
readBecauseYouLikeCache: () => null,
}));
vi.mock('@/lib/hooks/useConnectionStatus', () => ({
useConnectionStatus: () => ({ status: homeMocks.connection.status }),
}));
vi.mock('@/features/offline', () => ({
useOfflineBrowseContext: () => ({ active: false }),
useOfflineBrowseReloadToken: () => 0,
useDevOfflineBrowseStore: (selector: (state: { forceOffline: boolean }) => unknown) => (
selector({ forceOffline: false })
),
}));
vi.mock('@/lib/library/libraryBrowseScope', async importOriginal => ({
...(await importOriginal<typeof import('@/lib/library/libraryBrowseScope')>()),
deriveLibraryBrowseScope: () => ({
anchorServerId: 'server-a',
pairs: [{ serverId: 'server-a', libraryId: 'library-a' }],
}),
}));
vi.mock('@/lib/network/serverReachability', async importOriginal => ({
...(await importOriginal<typeof import('@/lib/network/serverReachability')>()),
useUnavailableServerIds: () => homeMocks.unavailableServerIds,
}));
vi.mock('@/features/home/pages/homeFeedLoader', () => ({
deriveHomeFeedScope: () => ({ serverIds: ['server-a'], scopeKey: 'scope' }),
loadHomeFeedWithStatus: homeMocks.loadHomeFeedWithStatus,
loadHomeChronologicalFeed: homeMocks.loadHomeChronologicalFeed,
loadMoreHomeAlbums: vi.fn(),
patchHomeChronologicalFeed: (snapshot: HomeFeedSnapshot) => snapshot,
preserveHomeChronologicalFeeds: (snapshot: HomeFeedSnapshot) => snapshot,
}));
vi.mock('@/features/home/pages/homeCoverPrefetch', () => ({
groupHomeCoverPrefetchBuckets: () => [],
homeDiscoverCoverPrefetchBucket: () => ({}),
shouldOfferHomeLoadMore: () => false,
}));
vi.mock('@/store/offlineLocalLibrarySyncRevision', () => ({
useLibraryScopeSyncRevision: () => 0,
}));
vi.mock('@/features/home/pages/homeDiagnosticHelpers', () => ({
homeSnapshotForEnabledCoverWarm: (snapshot: HomeFeedSnapshot) => snapshot,
preserveDisabledHomeSections: (snapshot: HomeFeedSnapshot) => snapshot,
reportCachedHomeDiagnostics: vi.fn(),
}));
vi.mock('@/app/startupSplash', () => ({ scheduleStartupSplashDismiss: vi.fn() }));
import Home from '@/features/home/pages/Home';
import { clearHomeFeedCache, readHomeFeedCache } from '@/features/home/store/homeFeedCache';
import { useHomeStore, DEFAULT_HOME_SECTIONS } from '@/features/home/store/homeStore';
import { useAuthStore } from '@/store/authStore';
import { useMigrationStore } from '@/store/migrationStore';
import { makeServer } from '@/test/helpers/factories';
import { resetAuthStore } from '@/test/helpers/storeReset';
import { renderWithProviders } from '@/test/helpers/renderWithProviders';
function deferred<T>() {
let resolve!: (value: T) => void;
const promise = new Promise<T>(res => { resolve = res; });
return { promise, resolve };
}
function snapshot(name: string): HomeFeedSnapshot {
return {
scopeKey: 'scope',
scopeVersion: 1,
savedAt: 1,
offsets: {
starred: { 'server-a': 0 },
recent: { offset: 0, hasMore: false },
random: { 'server-a': 0 },
mostPlayed: { 'server-a': 0 },
recentlyPlayed: { offset: 0, hasMore: false },
},
starred: [],
recent: [],
random: [],
heroAlbums: [{
id: name, name, artist: 'Artist', artistId: 'artist', songCount: 1, duration: 1,
}],
mostPlayed: [],
recentlyPlayed: [],
randomArtists: [],
discoverSongs: [],
};
}
describe('Home startup feed loading', () => {
beforeEach(() => {
resetAuthStore();
clearHomeFeedCache();
homeMocks.connection.status = 'checking';
homeMocks.loadHomeFeedWithStatus.mockReset();
homeMocks.loadHomeChronologicalFeed.mockReset();
homeMocks.loadHomeChronologicalFeed.mockResolvedValue({
status: 'success', albums: [], hasMore: false, durationMs: 0,
});
useMigrationStore.setState({ phase: 'idle' });
useHomeStore.setState({ sections: DEFAULT_HOME_SECTIONS });
const server = makeServer({ id: 'server-a' });
useAuthStore.setState({
servers: [server],
activeServerId: server.id,
libraryBrowseServerIds: [server.id],
musicFoldersByServer: { [server.id]: [] },
libraryBrowseSelectionByServer: { [server.id]: [] },
libraryBrowseScopeVersion: 1,
});
});
it('waits for migrations, retries on connection, and ignores invalidated loads', async () => {
const beforeBlocked = deferred<{ snapshot: HomeFeedSnapshot; emptySnapshotReliable: boolean }>();
const beforeConnected = deferred<{ snapshot: HomeFeedSnapshot; emptySnapshotReliable: boolean }>();
const connected = deferred<{ snapshot: HomeFeedSnapshot; emptySnapshotReliable: boolean }>();
homeMocks.loadHomeFeedWithStatus
.mockReturnValueOnce(beforeBlocked.promise)
.mockReturnValueOnce(beforeConnected.promise)
.mockReturnValueOnce(connected.promise);
const view = renderWithProviders(<Home />);
expect(homeMocks.loadHomeFeedWithStatus).not.toHaveBeenCalled();
await act(async () => {
useMigrationStore.setState({ phase: 'completed' });
});
await waitFor(() => expect(homeMocks.loadHomeFeedWithStatus).toHaveBeenCalledTimes(1));
await act(async () => {
useMigrationStore.setState({ phase: 'inspecting' });
beforeBlocked.resolve({ snapshot: snapshot('blocked-stale'), emptySnapshotReliable: true });
await beforeBlocked.promise;
});
expect(readHomeFeedCache('scope', 1)).toBeNull();
await act(async () => {
useMigrationStore.setState({ phase: 'completed' });
});
await waitFor(() => expect(homeMocks.loadHomeFeedWithStatus).toHaveBeenCalledTimes(2));
homeMocks.connection.status = 'connected';
view.rerender(<Home />);
await waitFor(() => expect(homeMocks.loadHomeFeedWithStatus).toHaveBeenCalledTimes(3));
await act(async () => {
beforeConnected.resolve({ snapshot: snapshot('connection-stale'), emptySnapshotReliable: true });
await beforeConnected.promise;
});
expect(readHomeFeedCache('scope', 1)).toBeNull();
await act(async () => {
connected.resolve({ snapshot: snapshot('fresh'), emptySnapshotReliable: true });
await connected.promise;
});
await waitFor(() => expect(readHomeFeedCache('scope', 1)?.heroAlbums[0]?.name).toBe('fresh'));
expect(screen.getByTestId('home-hero')).toHaveTextContent('fresh');
});
});
+16 -9
View File
@@ -26,6 +26,7 @@ import {
readHomeFeedCache,
readHomeFeedCacheStale,
patchHomeFeedCache,
shouldCacheColdHomeFeed,
writeHomeFeedCache,
type HomeFeedSnapshot,
} from '@/features/home/store/homeFeedCache';
@@ -39,7 +40,7 @@ import { useUnavailableServerIds } from '@/lib/network/serverReachability';
import {
deriveHomeFeedScope,
loadHomeChronologicalFeed,
loadHomeFeed,
loadHomeFeedWithStatus,
loadMoreHomeAlbums,
patchHomeChronologicalFeed,
preserveHomeChronologicalFeeds,
@@ -67,6 +68,7 @@ import {
type MainstageEnabledSections,
} from '@/features/home/pages/homeDiagnosticHelpers';
import { scheduleStartupSplashDismiss } from '@/app/startupSplash';
import { useMigrationStore } from '@/store/migrationStore';
/** Match Random Albums overshoot when mix filter uses album/artist axes so hero + discover row can still fill. */
const HOME_RANDOM_FETCH = 100;
@@ -147,6 +149,7 @@ export default function Home() {
],
);
const connStatus = useConnectionStatus().status;
const migrationReady = useMigrationStore(s => s.phase === 'completed');
const devForceOffline = useDevOfflineBrowseStore(s => s.forceOffline);
const offlineBrowseActive = useOfflineBrowseContext().active;
const offlineBrowseReloadTs = useOfflineBrowseReloadToken();
@@ -225,7 +228,7 @@ export default function Home() {
);
useEffect(() => {
if (serverIds.length === 0 || !scopeKey || !anchorServerId) return;
if (!migrationReady || serverIds.length === 0 || !scopeKey || !anchorServerId) return;
let cancelled = false;
const loadVersion = ++feedLoadVersionRef.current;
const isCurrentLoad = () => !cancelled && feedLoadVersionRef.current === loadVersion;
@@ -236,7 +239,7 @@ export default function Home() {
const albumMix =
mixCfg.enabled && (mixCfg.minAlbum > 0 || mixCfg.minArtist > 0);
const randomSize = albumMix ? HOME_RANDOM_FETCH : HOME_DISCOVER_SLICE;
const snapshot = loadHomeFeed({
const feed = loadHomeFeedWithStatus({
serverIds,
scopeKey,
anchorServerId,
@@ -291,7 +294,7 @@ export default function Home() {
};
if (mainstageDiagnosticsEnabled && chronological.recent) startDiagnostic('recent');
if (mainstageDiagnosticsEnabled && chronological.recentlyPlayed) startDiagnostic('recentlyPlayed');
return { snapshot, chronological };
return { feed, chronological };
};
const applyChronologicalResult = (
section: 'recent' | 'recentlyPlayed',
@@ -362,10 +365,10 @@ export default function Home() {
void (async () => {
try {
const freshLoad = startFreshHomeFeed();
const loaded = await freshLoad.snapshot;
const loaded = await freshLoad.feed;
if (!isCurrentLoad()) return;
const fresh = preserveDisabledHomeSections(
preserveHomeChronologicalFeeds(loaded, displayedSnapshotRef.current),
preserveHomeChronologicalFeeds(loaded.snapshot, displayedSnapshotRef.current),
displayedSnapshotRef.current,
getEffectiveEnabledSections(),
);
@@ -398,15 +401,17 @@ export default function Home() {
(async () => {
try {
const freshLoad = startFreshHomeFeed();
const loaded = await freshLoad.snapshot;
const loaded = await freshLoad.feed;
if (!isCurrentLoad()) return;
const snap = preserveDisabledHomeSections(
preserveHomeChronologicalFeeds(loaded, displayedSnapshotRef.current),
preserveHomeChronologicalFeeds(loaded.snapshot, displayedSnapshotRef.current),
displayedSnapshotRef.current,
getEffectiveEnabledSections(),
);
if (offlineBrowseActive && isHomeFeedSnapshotEmpty(snap)) return;
writeHomeFeedCache(snap);
if (shouldCacheColdHomeFeed(snap, loaded.emptySnapshotReliable, false)) {
writeHomeFeedCache(snap);
}
applyFeedSnapshot(snap);
await patchChronologicalFeeds(freshLoad.chronological);
if (!cancelled) setLoading(false);
@@ -437,6 +442,8 @@ export default function Home() {
offlineBrowseActive,
offlineBrowseReloadTs,
librarySyncRevision,
migrationReady,
connStatus,
]);
/** When offline toggles without a library-filter bump, re-apply stale cache if the feed was cleared. */
@@ -16,6 +16,7 @@ import {
deriveHomeFeedScope,
loadHomeChronologicalFeed,
loadHomeFeed,
loadHomeFeedWithStatus,
loadMoreHomeAlbums,
patchHomeChronologicalFeed,
preserveHomeChronologicalFeeds,
@@ -107,6 +108,89 @@ describe('homeFeedLoader pure helpers', () => {
});
describe('homeFeedLoader failure isolation', () => {
it('distinguishes failed all-empty loads from successful empty libraries', async () => {
const base = {
serverIds: ['a'], scopeKey: 'scope', scopeVersion: 1, randomSize: 0,
anchorServerId: 'a', scopes: [], showArtists: false, showSongs: false, mixConfig,
enabledSections: {
starred: true,
mostPlayed: false,
hero: false,
discover: false,
discoverArtists: false,
discoverSongs: false,
},
};
const failed = await loadHomeFeedWithStatus({
...base,
deps: {
getAlbumListForServer: vi.fn(async () => { throw new Error('offline'); }) as never,
filterAlbumsByMixRatingsAcrossServers: vi.fn(async albums => albums),
},
});
const successfulEmpty = await loadHomeFeedWithStatus({
...base,
deps: {
getAlbumListForServer: vi.fn(async () => []) as never,
filterAlbumsByMixRatingsAcrossServers: vi.fn(async albums => albums),
},
});
expect(failed.snapshot.starred).toEqual([]);
expect(failed.emptySnapshotReliable).toBe(false);
expect(successfulEmpty.snapshot.starred).toEqual([]);
expect(successfulEmpty.emptySnapshotReliable).toBe(true);
});
it('marks an all-timeout empty load as unreliable', async () => {
vi.useFakeTimers();
const result = loadHomeFeedWithStatus({
serverIds: ['a'], scopeKey: 'scope', scopeVersion: 1, randomSize: 0,
anchorServerId: 'a', scopes: [], showArtists: false, showSongs: false, mixConfig,
enabledSections: {
starred: true,
mostPlayed: false,
hero: false,
discover: false,
discoverArtists: false,
discoverSongs: false,
},
deps: {
getAlbumListForServer: vi.fn(() => new Promise<never>(() => undefined)) as never,
filterAlbumsByMixRatingsAcrossServers: vi.fn(async albums => albums),
},
});
await vi.advanceTimersByTimeAsync(HOME_REQUEST_TIMEOUT_MS);
await expect(result).resolves.toMatchObject({ emptySnapshotReliable: false });
vi.useRealTimers();
});
it('does not trust an empty snapshot when only some requested servers respond', async () => {
const result = await loadHomeFeedWithStatus({
serverIds: ['a', 'b'], scopeKey: 'scope', scopeVersion: 1, randomSize: 0,
anchorServerId: 'a', scopes: [], showArtists: false, showSongs: false, mixConfig,
enabledSections: {
starred: true,
mostPlayed: false,
hero: false,
discover: false,
discoverArtists: false,
discoverSongs: false,
},
deps: {
getAlbumListForServer: vi.fn(async serverId => {
if (serverId === 'b') throw new Error('offline');
return [];
}) as never,
filterAlbumsByMixRatingsAcrossServers: vi.fn(async albums => albums),
},
});
expect(result.snapshot.starred).toEqual([]);
expect(result.emptySnapshotReliable).toBe(false);
});
it('reuses a timed-out chronological invoke until the native read settles', async () => {
vi.useFakeTimers();
const libraryScopeListMainstageAlbums = vi.fn(() => new Promise<never>(() => {}));
+83 -41
View File
@@ -42,6 +42,11 @@ export interface HomeSectionResult {
detail?: string;
}
export interface HomeFeedLoadResult {
snapshot: HomeFeedSnapshot;
emptySnapshotReliable: boolean;
}
type OwnedAlbum = SubsonicAlbum & { serverId: string };
type OwnedArtist = SubsonicArtist & { serverId: string };
type OwnedSong = SubsonicSong & { serverId: string };
@@ -352,7 +357,7 @@ export function preserveHomeChronologicalFeeds(
type TimedServerItems<T> = {
items: T[];
durationMs: number;
outcome: 'rows' | 'empty' | 'timeout' | 'error';
outcome: 'rows' | 'empty' | 'timeout' | 'error' | 'skipped';
source?: 'local' | 'network';
};
@@ -363,7 +368,7 @@ async function loadServerAlbums(
deps: HomeFeedLoaderDeps,
): Promise<TimedServerItems<OwnedAlbum>> {
const startedAt = nowMs();
if (size <= 0) return { items: [], durationMs: 0, outcome: 'empty' };
if (size <= 0) return { items: [], durationMs: 0, outcome: 'skipped' };
let timer: ReturnType<typeof setTimeout> | undefined;
const request = deps.getAlbumListForServer(serverId, type, size, 0, {}, HOME_REQUEST_TIMEOUT_MS)
.then(albums => ({ albums, outcome: albums.length > 0 ? 'rows' as const : 'empty' as const }))
@@ -389,7 +394,7 @@ async function loadServerArtists(
deps: HomeFeedLoaderDeps,
): Promise<TimedServerItems<OwnedArtist>> {
const startedAt = nowMs();
if (size <= 0) return { items: [], durationMs: 0, outcome: 'empty' };
if (size <= 0) return { items: [], durationMs: 0, outcome: 'skipped' };
const request: Promise<{
artists: SubsonicArtist[];
source: 'local' | 'network';
@@ -436,27 +441,44 @@ async function loadServerSongs(
size: number,
scopeFlightKey: string,
deps: HomeFeedLoaderDeps,
): Promise<OwnedSong[]> {
if (size <= 0) return [];
): Promise<TimedServerItems<OwnedSong>> {
if (size <= 0) return { items: [], durationMs: 0, outcome: 'skipped' };
const startedAt = nowMs();
const songs = await withinHomeDeadline(
isolated(async () => {
const local = await withinDeadline(
runLibraryLocalReadSingleFlight(
JSON.stringify(['home-songs', scopeFlightKey, serverId, size]),
() => deps.runLocalRandomSongs(serverId, size),
),
null,
Math.min(HOME_LOCAL_READ_TIMEOUT_MS, remainingHomeDeadline(startedAt)),
);
if (local != null) return local;
const remainingMs = remainingHomeDeadline(startedAt);
if (remainingMs <= 0) return [];
return deps.getRandomSongsForServer(serverId, size, undefined, remainingMs);
}, [] as SubsonicSong[]),
[] as SubsonicSong[],
const request: Promise<{
songs: SubsonicSong[];
source: 'local' | 'network';
outcome: TimedServerItems<OwnedSong>['outcome'];
}> = (async () => {
const local = await withinDeadline(
runLibraryLocalReadSingleFlight(
JSON.stringify(['home-songs', scopeFlightKey, serverId, size]),
() => deps.runLocalRandomSongs(serverId, size),
),
null,
Math.min(HOME_LOCAL_READ_TIMEOUT_MS, remainingHomeDeadline(startedAt)),
);
if (local != null) {
return { songs: local, source: 'local' as const, outcome: local.length > 0 ? 'rows' as const : 'empty' as const };
}
const remainingMs = remainingHomeDeadline(startedAt);
if (remainingMs <= 0) {
return { songs: [] as SubsonicSong[], source: 'network' as const, outcome: 'timeout' as const };
}
const songs = await deps.getRandomSongsForServer(serverId, size, undefined, remainingMs);
return { songs, source: 'network' as const, outcome: songs.length > 0 ? 'rows' as const : 'empty' as const };
})();
const result = await withinHomeDeadline(
isolated(() => request, {
songs: [] as SubsonicSong[], source: 'local' as const, outcome: 'error' as const,
}),
{ songs: [] as SubsonicSong[], source: 'local' as const, outcome: 'timeout' as const },
);
return songs.map(song => ({ ...song, serverId }));
return {
items: result.songs.map(song => ({ ...song, serverId })),
durationMs: elapsedMs(startedAt),
outcome: result.outcome,
source: result.source,
};
}
function formatServerTimings<T>(
@@ -469,7 +491,7 @@ function formatServerTimings<T>(
}).join(', ');
}
export async function loadHomeFeed(options: LoadHomeFeedOptions): Promise<HomeFeedSnapshot> {
export async function loadHomeFeedWithStatus(options: LoadHomeFeedOptions): Promise<HomeFeedLoadResult> {
const deps = { ...defaultDeps, ...options.deps };
const enabled: HomeFeedEnabledSections = {
starred: true,
@@ -499,7 +521,6 @@ export async function loadHomeFeed(options: LoadHomeFeedOptions): Promise<HomeFe
options.scopeVersion,
options.syncRevision ?? 0,
]);
const emptyGroups = () => options.serverIds.map(() => [] as OwnedAlbum[]);
const loadAlbumGroups = (type: Parameters<typeof getAlbumListForServer>[1], quotas: number[]) => (
Promise.all(options.serverIds.map((serverId, index) => (
loadServerAlbums(serverId, type, quotas[index] ?? 0, deps)
@@ -514,9 +535,9 @@ export async function loadHomeFeed(options: LoadHomeFeedOptions): Promise<HomeFe
status: 'success', durationMs: elapsedMs(starredStartedAt), itemCount: items.length,
detail: formatServerTimings(options.serverIds, groups),
});
return { groups: groups.map(group => group.items), items };
return { groups, items };
})
: Promise.resolve({ groups: emptyGroups(), items: [] as OwnedAlbum[] });
: Promise.resolve({ groups: [] as TimedServerItems<OwnedAlbum>[], items: [] as OwnedAlbum[] });
const mostPlayedStartedAt = nowMs();
const mostPlayedPromise = enabled.mostPlayed
@@ -526,9 +547,9 @@ export async function loadHomeFeed(options: LoadHomeFeedOptions): Promise<HomeFe
status: 'success', durationMs: elapsedMs(mostPlayedStartedAt), itemCount: items.length,
detail: formatServerTimings(options.serverIds, groups),
});
return { groups: groups.map(group => group.items), items };
return { groups, items };
})
: Promise.resolve({ groups: emptyGroups(), items: [] as OwnedAlbum[] });
: Promise.resolve({ groups: [] as TimedServerItems<OwnedAlbum>[], items: [] as OwnedAlbum[] });
const randomEnabled = enabled.hero || enabled.discover;
const randomStartedAt = nowMs();
@@ -559,9 +580,13 @@ export async function loadHomeFeed(options: LoadHomeFeedOptions): Promise<HomeFe
formatServerTimings(options.serverIds, groups),
].join('; '),
});
return { groups: groups.map(group => group.items), heroAlbums, random };
return { groups, heroAlbums, random };
})
: Promise.resolve({ groups: emptyGroups(), heroAlbums: [] as OwnedAlbum[], random: [] as OwnedAlbum[] });
: Promise.resolve({
groups: [] as TimedServerItems<OwnedAlbum>[],
heroAlbums: [] as OwnedAlbum[],
random: [] as OwnedAlbum[],
});
const artistsStartedAt = nowMs();
const artistsPromise = enabled.discoverArtists
@@ -576,22 +601,22 @@ export async function loadHomeFeed(options: LoadHomeFeedOptions): Promise<HomeFe
status: 'success', durationMs: elapsedMs(artistsStartedAt), itemCount: items.length,
detail: formatServerTimings(options.serverIds, groups),
});
return items;
return { groups, items };
})
: Promise.resolve([] as OwnedArtist[]);
: Promise.resolve({ groups: [] as TimedServerItems<OwnedArtist>[], items: [] as OwnedArtist[] });
const songsStartedAt = nowMs();
const songsPromise = enabled.discoverSongs
? Promise.all(options.serverIds.map((serverId, index) => (
loadServerSongs(serverId, songQuotas[index] ?? 0, scopeFlightKey, deps)
))).then(groups => {
const items = dedupeOwned(stableRoundRobin(groups, HOME_DISCOVER_SONGS_SIZE));
const items = dedupeOwned(stableRoundRobin(groups.map(group => group.items), HOME_DISCOVER_SONGS_SIZE));
report('discoverSongs', {
status: 'success', durationMs: elapsedMs(songsStartedAt), itemCount: items.length,
});
return items;
return { groups, items };
})
: Promise.resolve([] as OwnedSong[]);
: Promise.resolve({ groups: [] as TimedServerItems<OwnedSong>[], items: [] as OwnedSong[] });
const [starredResult, mostPlayedResult, randomResult, artists, songs] = await Promise.all([
starredPromise,
@@ -607,11 +632,11 @@ export async function loadHomeFeed(options: LoadHomeFeedOptions): Promise<HomeFe
options.serverIds.map((serverId, index) => [serverId, groups[index]?.length ?? 0]),
));
};
advanceInitial('starred', starredResult.groups);
advanceInitial('random', randomResult.groups);
advanceInitial('mostPlayed', mostPlayedResult.groups);
advanceInitial('starred', starredResult.groups.map(group => group.items));
advanceInitial('random', randomResult.groups.map(group => group.items));
advanceInitial('mostPlayed', mostPlayedResult.groups.map(group => group.items));
return {
const snapshot = {
scopeKey: options.scopeKey,
scopeVersion: options.scopeVersion,
savedAt: Date.now(),
@@ -622,9 +647,26 @@ export async function loadHomeFeed(options: LoadHomeFeedOptions): Promise<HomeFe
random: randomResult.random,
mostPlayed: mostPlayedResult.items,
recentlyPlayed: [],
randomArtists: artists,
discoverSongs: songs,
randomArtists: artists.items,
discoverSongs: songs.items,
};
const attemptedGroups = [
...(enabled.starred ? starredResult.groups : []),
...(enabled.mostPlayed ? mostPlayedResult.groups : []),
...(randomEnabled ? randomResult.groups : []),
...(enabled.discoverArtists ? artists.groups : []),
...(enabled.discoverSongs ? songs.groups : []),
];
const requestedGroups = attemptedGroups.filter(group => group.outcome !== 'skipped');
return {
snapshot,
emptySnapshotReliable: requestedGroups.length === 0
|| requestedGroups.every(group => group.outcome === 'rows' || group.outcome === 'empty'),
};
}
export async function loadHomeFeed(options: LoadHomeFeedOptions): Promise<HomeFeedSnapshot> {
return (await loadHomeFeedWithStatus(options)).snapshot;
}
export async function loadMoreHomeAlbums(options: LoadMoreHomeAlbumsOptions): Promise<HomeFeedSnapshot> {
@@ -4,8 +4,10 @@ import {
patchHomeFeedCache,
readHomeFeedCache,
readHomeFeedCacheStale,
shouldCacheColdHomeFeed,
writeHomeFeedCache,
} from '@/features/home/store/homeFeedCache';
import type { HomeFeedSnapshot } from '@/features/home/store/homeFeedCache';
function write(scopeKey: string, scopeVersion: number) {
writeHomeFeedCache({
@@ -23,6 +25,23 @@ function write(scopeKey: string, scopeVersion: number) {
});
}
function emptySnapshot(): HomeFeedSnapshot {
return {
scopeKey: 'scope',
scopeVersion: 1,
savedAt: 1,
offsets: {
starred: {},
recent: { offset: 0, hasMore: false },
random: {},
mostPlayed: {},
recentlyPlayed: { offset: 0, hasMore: false },
},
starred: [], recent: [], random: [], heroAlbums: [], mostPlayed: [],
recentlyPlayed: [], randomArtists: [], discoverSongs: [],
};
}
describe('homeFeedCache', () => {
beforeEach(() => {
clearHomeFeedCache();
@@ -72,4 +91,17 @@ describe('homeFeedCache', () => {
expect(patched?.recent.map(album => album.id)).toEqual(['new']);
expect(readHomeFeedCache('scope', 1)?.offsets.recent).toEqual({ offset: 1, hasMore: true });
});
it('rejects only unreliable cold empty snapshots', () => {
const empty = emptySnapshot();
expect(shouldCacheColdHomeFeed(empty, false, false)).toBe(false);
expect(shouldCacheColdHomeFeed(empty, true, false)).toBe(true);
expect(shouldCacheColdHomeFeed(empty, true, true)).toBe(false);
expect(shouldCacheColdHomeFeed({
...empty,
heroAlbums: [{
id: 'album', name: 'Album', artist: 'Artist', artistId: 'artist', songCount: 1, duration: 1,
}],
}, false, false)).toBe(true);
});
});
+9
View File
@@ -79,6 +79,15 @@ export function isHomeFeedSnapshotEmpty(snap: HomeFeedSnapshot): boolean {
&& snap.randomArtists.length === 0;
}
export function shouldCacheColdHomeFeed(
snap: HomeFeedSnapshot,
emptySnapshotReliable: boolean,
offlineBrowseActive: boolean,
): boolean {
if (!isHomeFeedSnapshotEmpty(snap)) return true;
return !offlineBrowseActive && emptySnapshotReliable;
}
export function writeHomeFeedCache(data: Omit<HomeFeedSnapshot, 'savedAt'>): void {
const snapshot = { ...data, savedAt: Date.now() };
const key = cacheKey(snapshot.scopeKey, snapshot.scopeVersion);