Fixes for tests
This commit is contained in:
@@ -55,14 +55,12 @@ impl HybridRepository {
|
||||
self.online.get_video_stream_url(item_id, media_source_id, start_time_seconds, audio_stream_index).await
|
||||
}
|
||||
|
||||
/// Race cache vs server, return first valid result
|
||||
/// Prefer cache if it has meaningful content, otherwise use server
|
||||
/// Cache-first query: try cache, fall back to server on miss.
|
||||
///
|
||||
/// Core algorithm of the cache-first parallel racing strategy.
|
||||
/// Runs both cache and server queries concurrently, then:
|
||||
/// 1. If cache has meaningful content → return cache (fast path)
|
||||
/// 2. If cache is empty/stale → return server (fresh data)
|
||||
/// 3. If server fails → return cache even if empty (offline fallback)
|
||||
/// 1. Check cache (100ms timeout applied by caller via cache_with_timeout)
|
||||
/// 2. If cache has meaningful content → return immediately (fast path)
|
||||
/// 3. If cache is empty/stale → query server (fresh data)
|
||||
/// 4. If server fails → return cache even if empty (offline fallback)
|
||||
///
|
||||
/// @req: UR-002 - Access media when online or offline
|
||||
/// @req: DR-013 - Repository pattern for online/offline data access
|
||||
@@ -76,24 +74,20 @@ impl HybridRepository {
|
||||
F1: std::future::Future<Output = Result<T, RepoError>> + Send,
|
||||
F2: std::future::Future<Output = Result<T, RepoError>> + Send,
|
||||
{
|
||||
// Wait for both to complete (cache has 100ms timeout)
|
||||
let (cache_result, server_result) = tokio::join!(cache_future, server_future);
|
||||
// Try cache first (100ms timeout already applied by callers)
|
||||
let cache_result = cache_future.await;
|
||||
|
||||
// Prefer cache if it has meaningful content
|
||||
if let Ok(data) = &cache_result {
|
||||
if data.has_content() {
|
||||
debug!("[HybridRepo] Using cache result (has content)");
|
||||
debug!("[HybridRepo] Cache hit, returning immediately");
|
||||
return Ok(data.clone());
|
||||
}
|
||||
}
|
||||
|
||||
// Fall back to server result
|
||||
match server_result {
|
||||
Ok(data) => {
|
||||
debug!("[HybridRepo] Using server result");
|
||||
// TODO: Spawn background cache update
|
||||
Ok(data)
|
||||
}
|
||||
// Cache miss — fall back to server
|
||||
debug!("[HybridRepo] Cache miss, querying server");
|
||||
match server_future.await {
|
||||
Ok(data) => Ok(data),
|
||||
Err(e) => {
|
||||
// Server failed, try to return cache even if empty
|
||||
cache_result.or(Err(e))
|
||||
@@ -135,46 +129,57 @@ impl MediaRepository for HybridRepository {
|
||||
let parent_id_for_save = parent_id.clone();
|
||||
let opts_clone = options.clone();
|
||||
|
||||
// Check cache first to see if we have data
|
||||
let cache_future = self.cache_with_timeout(async move {
|
||||
offline.get_items(&parent_id, opts_clone).await
|
||||
// Start server request in background (non-blocking)
|
||||
let server_handle = tokio::spawn(async move {
|
||||
online.get_items(&parent_id_clone, options).await
|
||||
});
|
||||
|
||||
let server_future = async move {
|
||||
online.get_items(&parent_id_clone, options).await
|
||||
};
|
||||
// Check cache first (fast, 100ms timeout)
|
||||
let cache_result = self.cache_with_timeout(async move {
|
||||
offline.get_items(&parent_id, opts_clone).await
|
||||
}).await;
|
||||
|
||||
// Wait for both, prefer cache if available
|
||||
let (cache_result, server_result) = tokio::join!(cache_future, server_future);
|
||||
|
||||
// Check if cache had meaningful content
|
||||
let cache_had_content = cache_result.as_ref()
|
||||
.map(|data| data.has_content())
|
||||
.unwrap_or(false);
|
||||
|
||||
// Prefer cache if it has content
|
||||
let result = if cache_had_content {
|
||||
debug!("[HybridRepo] Using cached data for parent {}", &parent_id_for_save[..8.min(parent_id_for_save.len())]);
|
||||
cache_result?
|
||||
} else {
|
||||
// Use server result and save to cache for next time
|
||||
let server_data = server_result?;
|
||||
|
||||
if !server_data.items.is_empty() {
|
||||
let items_clone = server_data.items.clone();
|
||||
// Cache hit: return immediately, update cache in background
|
||||
if let Ok(data) = &cache_result {
|
||||
if data.has_content() {
|
||||
debug!("[HybridRepo] Cache hit for get_items, returning immediately for parent {}", &parent_id_for_save[..8.min(parent_id_for_save.len())]);
|
||||
// Background: save server result to cache when it arrives
|
||||
tokio::spawn(async move {
|
||||
if let Err(e) = offline_for_save.save_to_cache(&parent_id_for_save, &items_clone).await {
|
||||
warn!("[HybridRepo] Failed to save {} items to cache: {:?}", items_clone.len(), e);
|
||||
} else {
|
||||
debug!("[HybridRepo] Saved {} items to cache for parent {}", items_clone.len(), &parent_id_for_save[..8.min(parent_id_for_save.len())]);
|
||||
match server_handle.await {
|
||||
Ok(Ok(server_data)) if !server_data.items.is_empty() => {
|
||||
if let Err(e) = offline_for_save.save_to_cache(&parent_id_for_save, &server_data.items).await {
|
||||
warn!("[HybridRepo] Background cache update failed: {:?}", e);
|
||||
} else {
|
||||
debug!("[HybridRepo] Background updated {} cached items for parent {}", server_data.items.len(), &parent_id_for_save[..8.min(parent_id_for_save.len())]);
|
||||
}
|
||||
}
|
||||
_ => {} // Server failed or returned empty — keep existing cache
|
||||
}
|
||||
});
|
||||
return Ok(data.clone());
|
||||
}
|
||||
}
|
||||
|
||||
server_data
|
||||
};
|
||||
|
||||
Ok(result)
|
||||
// Cache miss — wait for server result
|
||||
match server_handle.await {
|
||||
Ok(Ok(server_data)) => {
|
||||
if !server_data.items.is_empty() {
|
||||
let items_clone = server_data.items.clone();
|
||||
tokio::spawn(async move {
|
||||
if let Err(e) = offline_for_save.save_to_cache(&parent_id_for_save, &items_clone).await {
|
||||
warn!("[HybridRepo] Failed to save {} items to cache: {:?}", items_clone.len(), e);
|
||||
} else {
|
||||
debug!("[HybridRepo] Saved {} items to cache for parent {}", items_clone.len(), &parent_id_for_save[..8.min(parent_id_for_save.len())]);
|
||||
}
|
||||
});
|
||||
}
|
||||
Ok(server_data)
|
||||
}
|
||||
Ok(Err(e)) => cache_result.or(Err(e)),
|
||||
Err(join_err) => cache_result.or(Err(RepoError::Network {
|
||||
message: format!("Server task failed: {}", join_err),
|
||||
})),
|
||||
}
|
||||
}
|
||||
|
||||
async fn get_item(&self, item_id: &str) -> Result<MediaItem, RepoError> {
|
||||
@@ -447,35 +452,34 @@ impl MediaRepository for HybridRepository {
|
||||
let playlist_id_clone = playlist_id.clone();
|
||||
let playlist_id_for_save = playlist_id.clone();
|
||||
|
||||
let cache_future = self.cache_with_timeout(async move {
|
||||
offline.get_playlist_items(&playlist_id).await
|
||||
// Start server request in background (non-blocking)
|
||||
let server_handle = tokio::spawn(async move {
|
||||
online.get_playlist_items(&playlist_id_clone).await
|
||||
});
|
||||
|
||||
let server_future = async move {
|
||||
online.get_playlist_items(&playlist_id_clone).await
|
||||
};
|
||||
// Check cache first (fast, 100ms timeout)
|
||||
let cache_result = self.cache_with_timeout(async move {
|
||||
offline.get_playlist_items(&playlist_id).await
|
||||
}).await;
|
||||
|
||||
let (cache_result, server_result) = tokio::join!(cache_future, server_future);
|
||||
|
||||
let cache_had_content = cache_result.as_ref()
|
||||
.map(|data| data.has_content())
|
||||
.unwrap_or(false);
|
||||
|
||||
if cache_had_content {
|
||||
// If server also succeeded, update cache in background
|
||||
if let Ok(server_entries) = server_result {
|
||||
// Cache hit: return immediately, update cache in background
|
||||
if let Ok(data) = &cache_result {
|
||||
if data.has_content() {
|
||||
debug!("[HybridRepo] Cache hit for playlist items, returning immediately");
|
||||
tokio::spawn(async move {
|
||||
if let Err(e) = offline_for_save.save_playlist_items_to_cache(&playlist_id_for_save, &server_entries).await {
|
||||
warn!("[HybridRepo] Failed to update playlist cache: {:?}", e);
|
||||
if let Ok(Ok(server_entries)) = server_handle.await {
|
||||
if let Err(e) = offline_for_save.save_playlist_items_to_cache(&playlist_id_for_save, &server_entries).await {
|
||||
warn!("[HybridRepo] Failed to update playlist cache: {:?}", e);
|
||||
}
|
||||
}
|
||||
});
|
||||
return cache_result;
|
||||
}
|
||||
return cache_result;
|
||||
}
|
||||
|
||||
// Cache miss - use server result
|
||||
match server_result {
|
||||
Ok(entries) => {
|
||||
// Cache miss — wait for server result
|
||||
match server_handle.await {
|
||||
Ok(Ok(entries)) => {
|
||||
let entries_clone = entries.clone();
|
||||
tokio::spawn(async move {
|
||||
if let Err(e) = offline_for_save.save_playlist_items_to_cache(&playlist_id_for_save, &entries_clone).await {
|
||||
@@ -484,7 +488,10 @@ impl MediaRepository for HybridRepository {
|
||||
});
|
||||
Ok(entries)
|
||||
}
|
||||
Err(e) => cache_result.or(Err(e)),
|
||||
Ok(Err(e)) => cache_result.or(Err(e)),
|
||||
Err(join_err) => cache_result.or(Err(RepoError::Network {
|
||||
message: format!("Server task failed: {}", join_err),
|
||||
})),
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user