diff --git a/Data/DatabaseManager.cs b/Data/DatabaseManager.cs index db3f43b..2f815a7 100644 --- a/Data/DatabaseManager.cs +++ b/Data/DatabaseManager.cs @@ -89,6 +89,15 @@ public void Initialise() using var conn = OpenConnection(); ApplyPragmas(conn); CreateSchema(conn); + // Derived absence history before this policy was inflated by batching. + // Reset it once, without changing catalog/user/media or retry state. + conn.RunInTransaction(c => + { + c.Execute("CREATE TABLE IF NOT EXISTS catalog_pruning_policy (version INTEGER PRIMARY KEY);"); + c.Execute(@"UPDATE catalog_items SET global_absent_syncs = 0 + WHERE NOT EXISTS (SELECT 1 FROM catalog_pruning_policy WHERE version = 1);"); + c.Execute("INSERT OR IGNORE INTO catalog_pruning_policy VALUES (1);"); + }); BuildColumnMaps(conn); } @@ -427,62 +436,94 @@ public async Task> GetCatalogItemsBySourceAsync(string source) ReadCatalogItem); } - /// - /// Returns prune CANDIDATES: items whose global_absent_syncs has reached the - /// threshold, excluding admin-blocked items and items still referenced by ANY - /// user's collection/list (collection_membership) or recorded in playback_log. - /// Does NOT delete — the caller additionally filters out anything Emby marks as - /// played (ever-watched protection) before soft-deleting via - /// . Returns (strm_path, aio_id) pairs. - /// - public async Task> GetAbsentPruneCandidatesAsync( + // Only automatic broad feeds own this retention policy. Any active alias + // with different ownership, a block, or a newer observation protects the title. + private const string AbsentPruneGuard = @" + source IN ('aiostreams','aiostreams_secondary','cinemeta_default') + AND removed_at IS NULL AND blocked_at IS NULL + AND ((local_source = 'strm' AND strm_path IS NOT NULL AND strm_path <> '') + OR ((strm_path IS NULL OR strm_path = '') AND COALESCE(local_source,'') IN ('','strm') + AND (local_path IS NULL OR local_path = ''))) + AND item_state NOT IN (3,5,11) AND media_type IN ('movie','series','anime') + AND global_absent_syncs >= @threshold + AND NOT EXISTS (SELECT 1 FROM catalog_items other + WHERE (other.aio_id = catalog_items.aio_id + OR (other.tmdb_id IS NOT NULL AND other.tmdb_id <> '' AND other.tmdb_id = catalog_items.tmdb_id + AND ((other.media_type IN ('series','anime')) = (catalog_items.media_type IN ('series','anime')))) + OR (other.strm_path IS NOT NULL AND other.strm_path <> '' AND other.strm_path = catalog_items.strm_path)) + AND (other.removed_at IS NULL OR other.local_source = 'library' OR other.blocked_at IS NOT NULL) + AND (other.source NOT IN ('aiostreams','aiostreams_secondary','cinemeta_default') + OR other.local_source = 'library' OR other.item_state IN (3,5,11) + OR other.blocked_at IS NOT NULL OR other.global_absent_syncs < @threshold + OR EXISTS (SELECT 1 FROM collection_membership cm WHERE cm.aio_id=other.aio_id) + OR EXISTS (SELECT 1 FROM playback_log pl WHERE pl.aio_id=other.aio_id))) + AND NOT EXISTS (SELECT 1 FROM collection_membership cm WHERE cm.aio_id = catalog_items.aio_id) + AND NOT EXISTS (SELECT 1 FROM playback_log pl WHERE pl.aio_id = catalog_items.aio_id) + AND NOT EXISTS (SELECT 1 FROM media_items m + WHERE ((m.primary_id_type='imdb' AND m.primary_id=catalog_items.aio_id) + OR (m.primary_id_type='tmdb' AND m.primary_id=catalog_items.tmdb_id)) + AND (m.saved=1 OR m.blocked=1 OR m.favorited=1 OR m.watch_progress_pct>0) + AND ((m.media_type IN ('series','anime')) = (catalog_items.media_type IN ('series','anime')))) + AND NOT EXISTS (SELECT 1 FROM media_item_ids mi JOIN media_items m ON m.id=mi.media_item_id + WHERE ((mi.id_type='imdb' AND mi.id_value=catalog_items.aio_id) + OR (mi.id_type='tmdb' AND mi.id_value=catalog_items.tmdb_id)) + AND (m.saved=1 OR m.blocked=1 OR m.favorited=1 OR m.watch_progress_pct>0) + AND ((m.media_type IN ('series','anime')) = (catalog_items.media_type IN ('series','anime'))))"; + + /// Returns eligible broad-feed rows, including unpublished metadata, without retiring anything. + public async Task> GetAbsentPruneCandidatesAsync( CancellationToken cancellationToken = default) { - var threshold = Services.RuntimePolicy.GlobalAbsentSyncThreshold; - - // The "Respect user playlists when pruning" setting (Marvin tab) gates the - // collection_membership guard: when on (default), items in any user's - // list/collection are never pruned. The playback_log guard (ever-watched) - // is always applied regardless. - const bool respectPlaylists = true; - var membershipGuard = respectPlaylists - ? @"AND NOT EXISTS ( - SELECT 1 FROM collection_membership cm - WHERE cm.aio_id = catalog_items.aio_id - )" - : string.Empty; - - var selectSql = $@" - SELECT strm_path, aio_id FROM catalog_items - WHERE global_absent_syncs >= @threshold - AND removed_at IS NULL - AND blocked_at IS NULL - {membershipGuard} - AND NOT EXISTS ( - SELECT 1 FROM playback_log pl - WHERE pl.aio_id = catalog_items.aio_id - );"; - - var candidates = new List<(string? StrmPath, string AioId)>(); + var candidates = new List<(string Id, string? StrmPath, string AioId)>(); await _dbWriteGate.WaitAsync(cancellationToken); try { using var conn = OpenConnection(); - using var stmt = conn.PrepareStatement(selectSql); - BindInt(stmt, "@threshold", threshold); + using var stmt = conn.PrepareStatement($"SELECT id, strm_path, aio_id FROM catalog_items WHERE {AbsentPruneGuard};"); + BindInt(stmt, "@threshold", Services.RuntimePolicy.GlobalAbsentSyncThreshold); foreach (var row in stmt.AsRows()) - { - candidates.Add(( - row.IsDBNull(0) ? null : row.GetString(0), - row.GetString(1))); - } + candidates.Add((row.GetString(0), row.IsDBNull(1) ? null : row.GetString(1), row.GetString(2))); } - finally + finally { _dbWriteGate.Release(); } + return candidates; + } + + /// + /// Retires only the candidate's broad-feed row, rechecking database retention + /// under the publication lock. Returns the actual row count; user removal + /// and reconciliation continue to use their separate explicit-removal methods. + /// + public async Task RetireAbsentCatalogItemsAsync( + IReadOnlyList<(string Id, string? StrmPath, string AioId)> candidates, + CancellationToken cancellationToken = default) + { + var retired = 0; + await Services.ImportReconciliationService.MutationGate.WaitAsync(cancellationToken); + try { - _dbWriteGate.Release(); + await _dbWriteGate.WaitAsync(cancellationToken); + try + { + using var conn = OpenConnection(); + conn.RunInTransaction(c => + { + foreach (var item in candidates) + { + cancellationToken.ThrowIfCancellationRequested(); + using var stmt = c.PrepareStatement($@"UPDATE catalog_items SET removed_at = datetime('now') + WHERE id = @id AND COALESCE(strm_path,'') = COALESCE(@path,'') AND {AbsentPruneGuard};"); + BindText(stmt, "@id", item.Id); BindNullableText(stmt, "@path", item.StrmPath); + BindInt(stmt, "@threshold", Services.RuntimePolicy.GlobalAbsentSyncThreshold); + while (stmt.MoveNext()) { } + using var changes = c.PrepareStatement("SELECT changes();"); + foreach (var row in changes.AsRows()) retired += row.GetInt(0); + } + }); + } + finally { _dbWriteGate.Release(); } } - - return candidates; + finally { Services.ImportReconciliationService.MutationGate.Release(); } + return retired; } /// @@ -558,7 +599,7 @@ await ExecuteWriteAsync(sql, cmd => } /// - /// Two-phase global absent-sync counter update across ALL sources. + /// Two-phase absent-sync update for automatic broad catalog sources. /// /// Phase 1: increment global_absent_syncs for all active, non-blocked items /// that are NOT in the union of present IDs and NOT protected by @@ -584,45 +625,25 @@ public async Task IncrementGlobalAbsentSyncsAsync( using var conn = OpenConnection(); conn.RunInTransaction(c => { - // Phase 1: increment counter for all unprotected absentees - foreach (var batch in allPresentAioIds.Chunk(500)) + // Populate the entire union before one increment. NOT IN one + // 500-ID chunk at a time counted one sync as many absences. + c.Execute("CREATE TEMP TABLE present_catalog_ids (aio_id TEXT PRIMARY KEY COLLATE NOCASE);"); + foreach (var id in allPresentAioIds) { - var presentPlaceholders = string.Join(",", batch.Select((_, i) => $"@pid{i}")); - var phase1Sql = $@" - UPDATE catalog_items - SET global_absent_syncs = global_absent_syncs + 1 - WHERE removed_at IS NULL - AND blocked_at IS NULL - AND aio_id NOT IN ({presentPlaceholders}) - AND NOT EXISTS ( - SELECT 1 FROM collection_membership cm - WHERE cm.aio_id = catalog_items.aio_id - )"; - - using (var stmt = c.PrepareStatement(phase1Sql)) - { - for (int i = 0; i < batch.Length; i++) - BindText(stmt, $"@pid{i}", batch[i]); - while (stmt.MoveNext()) { } - } - } - - // Phase 2: reset counter for items present in at least one source - foreach (var batch in allPresentAioIds.Chunk(500)) - { - var placeholders = string.Join(",", batch.Select((_, i) => $"@id{i}")); - var phase2Sql = $@" - UPDATE catalog_items - SET global_absent_syncs = 0, - last_verified_at = @now - WHERE aio_id IN ({placeholders})"; - - using var stmt2 = c.PrepareStatement(phase2Sql); - BindInt(stmt2, "@now", now); - for (int i = 0; i < batch.Length; i++) - BindText(stmt2, $"@id{i}", batch[i]); - while (stmt2.MoveNext()) { } + using var insert = c.PrepareStatement("INSERT OR IGNORE INTO present_catalog_ids VALUES (@id);"); + BindText(insert, "@id", id); while (insert.MoveNext()) { } } + c.Execute(@"UPDATE catalog_items SET global_absent_syncs = global_absent_syncs + 1 + WHERE source IN ('aiostreams','aiostreams_secondary','cinemeta_default') + AND removed_at IS NULL AND blocked_at IS NULL AND item_state NOT IN (3,5,11) + AND COALESCE(local_source,'') <> 'library' + AND NOT EXISTS (SELECT 1 FROM present_catalog_ids p WHERE p.aio_id = catalog_items.aio_id) + AND NOT EXISTS (SELECT 1 FROM collection_membership cm WHERE cm.aio_id = catalog_items.aio_id);"); + using var reset = c.PrepareStatement(@"UPDATE catalog_items + SET global_absent_syncs = 0, last_verified_at = @now + WHERE source IN ('aiostreams','aiostreams_secondary','cinemeta_default') + AND EXISTS (SELECT 1 FROM present_catalog_ids p WHERE p.aio_id = catalog_items.aio_id);"); + BindInt(reset, "@now", now); while (reset.MoveNext()) { } }); } finally @@ -1436,6 +1457,8 @@ media_type TEXT NOT NULL CHECK(media_type IN ('movie', 'series', 'a UNIQUE(aio_id, source) ); CREATE INDEX IF NOT EXISTS idx_catalog_aio ON catalog_items(aio_id); +CREATE INDEX IF NOT EXISTS idx_catalog_tmdb_retention ON catalog_items(tmdb_id, media_type); +CREATE INDEX IF NOT EXISTS idx_catalog_path_retention ON catalog_items(strm_path); CREATE INDEX IF NOT EXISTS idx_catalog_active ON catalog_items(removed_at) WHERE removed_at IS NULL; CREATE INDEX IF NOT EXISTS idx_catalog_media_type ON catalog_items(media_type, removed_at); diff --git a/Services/CatalogPruningPolicy.cs b/Services/CatalogPruningPolicy.cs new file mode 100644 index 0000000..6c619cf --- /dev/null +++ b/Services/CatalogPruningPolicy.cs @@ -0,0 +1,67 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using InfiniteDrive.Tasks; +using MediaBrowser.Controller.Entities; +using MediaBrowser.Controller.Entities.TV; + +namespace InfiniteDrive.Services; + +internal static class CatalogPruningPolicy +{ + internal static bool IsCompleteProviderSnapshot(CatalogFetchResult result) => + result.ProviderReachable && string.IsNullOrEmpty(result.ErrorMessage) + && result.CatalogOutcomes.Count > 0 && result.CatalogOutcomes.Values.All(x => x.Succeeded); + + internal static bool CanObserveAbsence(int completeProviders, int plannedProviders) => + plannedProviders > 0 && completeProviders == plannedProviders; + + internal static bool IsManagedPath(PluginConfiguration config, string? path) + { + try + { + return !string.IsNullOrWhiteSpace(path) && new[] { config.SyncPathMovies, config.SyncPathShows, config.SyncPathAnime } + .Any(root => !string.IsNullOrWhiteSpace(root) && ImportInventory.SafeManagedPath(root, path!)); + } + catch { return false; } + } + + // Every native version, every user's state and every series episode participates. + // Missing/failed lookup evidence must never be interpreted as permission to prune. + internal static bool HasProtectedNativeItems(IEnumerable? matches, IEnumerable users, + Func?> episodes, + Func userData, Func owned) + { + try + { + if (matches == null) return true; + var allUsers = users.ToArray(); + if (allUsers.Length == 0) return true; + foreach (var root in matches) + { + if (root == null || Protected(root)) return true; + if (root is Series) + { + var children = episodes(root); + if (children == null) return true; + foreach (var child in children) + if (child == null || Protected(child)) return true; + } + } + return false; + + bool Protected(BaseItem item) + { + if (owned(item)) return true; + foreach (var user in allUsers) + { + var state = userData(user, item); + if (state == null || state.Played || state.PlayCount > 0 || state.PlaybackPositionTicks > 0 + || state.LastPlayedDate.HasValue || state.IsFavorite || state.Rating.HasValue) return true; + } + return false; + } + } + catch { return true; } + } +} diff --git a/Tasks/CatalogProviders.cs b/Tasks/CatalogProviders.cs index c548b76..e4f5ffa 100644 --- a/Tasks/CatalogProviders.cs +++ b/Tasks/CatalogProviders.cs @@ -422,13 +422,20 @@ private static async Task> DiscoverCatalogsAsync( internal const int AioCatalogPageSize = 100; internal const int AioCatalogMaxPages = 200; - internal static async Task<(List Items, CatalogOutcome Outcome)> FetchOneCatalogAsync( + internal static Task<(List Items, CatalogOutcome Outcome)> FetchOneCatalogAsync( AioStreamsClient client, AioStreamsCatalogDef catalog, ILogger logger, int itemCap, CancellationToken cancellationToken, Func? onProgress = null) + => FetchCatalogPagesAsync(catalog, logger, itemCap, cancellationToken, + offset => offset == 0 ? client.GetCatalogAsync(catalog.Type!, catalog.Id!, cancellationToken) + : client.GetCatalogAsync(catalog.Type!, catalog.Id!, null, null, offset, cancellationToken), onProgress); + + internal static async Task<(List Items, CatalogOutcome Outcome)> FetchCatalogPagesAsync( + AioStreamsCatalogDef catalog, ILogger logger, int itemCap, CancellationToken cancellationToken, + Func> fetchPage, Func? onProgress = null) { var items = new List(); var seenIds = new HashSet(StringComparer.OrdinalIgnoreCase); @@ -442,12 +449,13 @@ private static async Task> DiscoverCatalogsAsync( cancellationToken.ThrowIfCancellationRequested(); // First page uses the simple URL; subsequent pages add skip=N - var response = offset == 0 - ? await client.GetCatalogAsync(catalog.Type!, catalog.Id!, cancellationToken) - : await client.GetCatalogAsync(catalog.Type!, catalog.Id!, null, null, offset, cancellationToken); + var response = await fetchPage(offset); - if (response?.Metas == null || response.Metas.Count == 0) - break; + if (response?.Metas == null) + return (items, new CatalogOutcome { Succeeded = false, ItemCount = items.Count, + Error = "Catalog page unavailable; partial data retained" }); + if (response.Metas.Count == 0) + return (items, new CatalogOutcome { Succeeded = true, ItemCount = items.Count }); var newIdsOnPage = 0; foreach (var meta in response.Metas) @@ -474,12 +482,14 @@ private static async Task> DiscoverCatalogsAsync( // every selected catalog at 20. Stop only when the addon returns // no rows or repeats a page without yielding a usable new ID. if (newIdsOnPage == 0) - break; + return (items, new CatalogOutcome { Succeeded = false, ItemCount = items.Count, + Error = "Catalog repeated a page before reaching its configured limit" }); offset += response.Metas.Count; } - return (items, new CatalogOutcome { Succeeded = true, ItemCount = items.Count }); + return (items, new CatalogOutcome { Succeeded = items.Count >= itemCap, ItemCount = items.Count, + Error = items.Count >= itemCap ? null : "Catalog page ceiling reached before its configured limit" }); } catch (OperationCanceledException) { diff --git a/Tasks/CatalogSyncTask.cs b/Tasks/CatalogSyncTask.cs index ec3c8b6..56b2191 100755 --- a/Tasks/CatalogSyncTask.cs +++ b/Tasks/CatalogSyncTask.cs @@ -101,6 +101,7 @@ internal async Task RunSyncAsync(CancellationToken cancellationToken, IProgress< Plugin.Pipeline.SetPhase("CatalogSync", "Fetch"); progress?.Report(5); _logger.LogDebug("[InfiniteDrive] Starting FetchFromAllProvidersAsync with {Count} providers", providers.Count); + var plannedProviders = providers.Count; var (allItems, fetchedSourceIds, attemptedProviders) = await FetchFromAllProvidersAsync(providers, config, db, cancellationToken); _logger.LogInformation( "[InfiniteDrive] Fetched {Count} raw catalog items from all sources (attempted {Attempted} providers)", @@ -118,6 +119,7 @@ internal async Task RunSyncAsync(CancellationToken cancellationToken, IProgress< var (cinemetaItems, cinemetaIds, _) = await FetchFromAllProvidersAsync( new List { new CinemetaDefaultProvider() }, config, db, cancellationToken); + plannedProviders++; allItems.AddRange(cinemetaItems); foreach (var kvp in cinemetaIds) fetchedSourceIds[kvp.Key] = kvp.Value; @@ -136,7 +138,7 @@ internal async Task RunSyncAsync(CancellationToken cancellationToken, IProgress< progress?.Report(40); // 3b. Prune items removed from their sources (with safety check) - await PruneRemovedItemsAsync(db, fetchedSourceIds, attemptedProviders, cancellationToken); + await PruneRemovedItemsAsync(db, fetchedSourceIds, plannedProviders, config, cancellationToken); // 4. Check that Emby libraries cover the sync paths; warn if not WarnIfLibrariesMissing(config); @@ -332,7 +334,7 @@ private static List BuildProviders(PluginConfiguration config) foreach (var fi in fetchResult.Items) results.Add(fi); - if (fetchResult.ProviderReachable && fetchResult.Items.Count > 0) + if (CatalogPruningPolicy.IsCompleteProviderSnapshot(fetchResult)) { var idSet = fetchedSourceIds.GetOrAdd(provider.SourceKey, _ => new HashSet(StringComparer.OrdinalIgnoreCase)); @@ -496,10 +498,11 @@ private async Task PruneRemovedItemsAsync( Data.DatabaseManager db, Dictionary> fetchedSourceIds, int totalProviders, + PluginConfiguration config, CancellationToken cancellationToken) { // Skip pruning if any source failed to fetch - if (fetchedSourceIds.Count < totalProviders) + if (!CatalogPruningPolicy.CanObserveAbsence(fetchedSourceIds.Count, totalProviders)) { _logger.LogInformation( "[InfiniteDrive] Skipping pruning — only {Fetched}/{Total} sources fetched successfully. Some resolvers may be down.", @@ -515,13 +518,15 @@ private async Task PruneRemovedItemsAsync( allPresentAioIds.Add(id); } + if (allPresentAioIds.Count == 0) return; // Empty snapshots cannot authorize retirement. + // Global increment + reset await db.IncrementGlobalAbsentSyncsAsync(allPresentAioIds, cancellationToken); // Get prune candidates (absent past threshold, not blocked, not in any // list/membership, not in playback_log), then filter out anything Emby // marks as played by ANY user — ever-watched content is never auto-pruned. - List<(string? StrmPath, string AioId)> candidates; + List<(string Id, string? StrmPath, string AioId)> candidates; try { candidates = await db.GetAbsentPruneCandidatesAsync(cancellationToken); @@ -535,54 +540,31 @@ private async Task PruneRemovedItemsAsync( if (candidates.Count == 0) return; - var toDeleteAioIds = new List(); - var removedPaths = new List(); - int keptWatched = 0; + var toRetire = new List<(string Id, string? StrmPath, string AioId)>(); + int keptProtected = 0; foreach (var c in candidates) { - if (HasBeenPlayedByAnyUser(c.AioId)) + cancellationToken.ThrowIfCancellationRequested(); + if ((!string.IsNullOrEmpty(c.StrmPath) && !CatalogPruningPolicy.IsManagedPath(config, c.StrmPath)) + || HasProtectedNativeState(c.AioId, config)) { - keptWatched++; - continue; // ever-watched → keep + keptProtected++; + continue; } - toDeleteAioIds.Add(c.AioId); - if (!string.IsNullOrEmpty(c.StrmPath)) removedPaths.Add(c.StrmPath!); + toRetire.Add(c); } - - if (keptWatched > 0) - _logger.LogInformation("[InfiniteDrive] CatalogSyncTask: kept {Count} ever-watched items from pruning", keptWatched); - - if (toDeleteAioIds.Count == 0) - return; - + if (keptProtected > 0) + _logger.LogInformation("[InfiniteDrive] CatalogSyncTask: retained {Count} protected or unverified prune candidates", keptProtected); + if (toRetire.Count == 0) return; try { - await db.SoftDeleteCatalogItemsAsync(toDeleteAioIds, cancellationToken); + var retiredCount = await db.RetireAbsentCatalogItemsAsync(toRetire, cancellationToken); + _logger.LogInformation("[InfiniteDrive] CatalogSyncTask: retired {Count} absent broad-feed rows; native orphan cleanup handles unreferenced STRMs", retiredCount); } + catch (OperationCanceledException) { throw; } catch (Exception ex) { - _logger.LogWarning(ex, "[InfiniteDrive] CatalogSyncTask: soft-delete failed"); - return; - } - - if (removedPaths.Count == 0) - return; - - _logger.LogInformation( - "[InfiniteDrive] CatalogSyncTask: pruning {Count} globally absent items", - removedPaths.Count); - - foreach (var strmPath in removedPaths) - { - try - { - Plugin.Instance!.StrmFileManager.DeleteWithVersions(strmPath); - } - catch (Exception ex) - { - _logger.LogDebug(ex, - "[InfiniteDrive] CatalogSyncTask: could not delete {Path}", strmPath); - } + _logger.LogWarning(ex, "[InfiniteDrive] CatalogSyncTask: absent retirement failed"); } } @@ -639,19 +621,11 @@ public static Dictionary BuildLibraryItemMapPublic( return map; } - /// - /// Authoritative "ever watched" check against Emby's own play state: returns true - /// if the item (located by provider id) has been played by ANY user. Used to keep - /// watched content from being auto-pruned. Defensive: any lookup failure returns - /// false (treat as not-played) so a transient error never blocks pruning forever — - /// the SQL guards (membership/playback_log) remain the first line of protection. - /// - private bool HasBeenPlayedByAnyUser(string aioId) + // Native ownership and user history participate in retention; lookup failures retain. + private bool HasProtectedNativeState(string aioId, PluginConfiguration config) { try { - if (string.IsNullOrEmpty(aioId)) return false; - var providerIds = new List>(); if (aioId.StartsWith("tt", StringComparison.OrdinalIgnoreCase)) providerIds.Add(new KeyValuePair("Imdb", aioId)); @@ -660,27 +634,23 @@ private bool HasBeenPlayedByAnyUser(string aioId) var num = new string(aioId.Where(char.IsDigit).ToArray()); if (!string.IsNullOrEmpty(num)) providerIds.Add(new KeyValuePair("Tmdb", num)); } - if (providerIds.Count == 0) return false; - + if (providerIds.Count == 0) return true; var items = _libraryManager.GetItemList(new InternalItemsQuery { AnyProviderIdEquals = providerIds, IncludeItemTypes = new[] { "Movie", "Series", "Episode" }, - Limit = 1, + Recursive = true, }); - if (items == null || items.Length == 0 || items[0] == null) return false; - - var item = items[0]!; - foreach (var user in _userManager.Users) - { - if (item.IsPlayed(user)) return true; - } - return false; + return CatalogPruningPolicy.HasProtectedNativeItems(items, _userManager.Users, + item => _libraryManager.GetItemList(new InternalItemsQuery + { AncestorIds = new[] { item.InternalId }, Recursive = true, IncludeItemTypes = new[] { "Episode" } }), + (user, item) => BaseItem.UserDataManager?.GetUserData(user, item), + item => !string.IsNullOrWhiteSpace(item.Path) && !CatalogPruningPolicy.IsManagedPath(config, item.Path)); } catch (Exception ex) { - _logger.LogDebug(ex, "[InfiniteDrive] played-by-any-user check failed for {AioId}", aioId); - return false; + _logger.LogWarning(ex, "[InfiniteDrive] Native retention lookup failed for {AioId}; retaining title", aioId); + return true; } } diff --git a/Tests/CatalogPruningTests.cs b/Tests/CatalogPruningTests.cs new file mode 100644 index 0000000..b96f016 --- /dev/null +++ b/Tests/CatalogPruningTests.cs @@ -0,0 +1,230 @@ +using System; +using System.Collections.Generic; +using System.IO; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using InfiniteDrive.Data; +using InfiniteDrive.Models; +using InfiniteDrive.Services; +using InfiniteDrive.Tasks; +using MediaBrowser.Controller.Entities; +using MediaBrowser.Controller.Entities.Movies; +using MediaBrowser.Controller.Entities.TV; +using Microsoft.Extensions.Logging.Abstractions; +using SQLitePCL.pretty; +using Xunit; + +namespace InfiniteDrive.Tests; + +public sealed class CatalogPruningTests +{ + [Fact] + public async Task FullUnionCountsOneAbsenceAcrossMoreThanTwoSqlBatches() + { + using var h = new Harness(); var absent = await h.Seed("ttabsent"); var present = await h.Seed("tt1000"); + h.Sql("UPDATE catalog_items SET global_absent_syncs=2;"); + var ids = Enumerable.Range(0, 1001).Select(x => "tt" + x).ToHashSet(); + await h.Db.IncrementGlobalAbsentSyncsAsync(ids); + Assert.Equal(3, (await h.Db.GetImportCatalogByIdAsync(absent.Id))!.GlobalAbsentSyncs); + Assert.Equal(0, (await h.Db.GetImportCatalogByIdAsync(present.Id))!.GlobalAbsentSyncs); + await h.Db.IncrementGlobalAbsentSyncsAsync(new()); + Assert.Equal(3, (await h.Db.GetImportCatalogByIdAsync(absent.Id))!.GlobalAbsentSyncs); + } + + [Fact] + public async Task InvalidHistoricalCounterIsResetOnlyOnceWithoutChangingUserOrRetryState() + { + using var h = new Harness(); var item = await h.Seed("ttlegacy"); + h.Sql("DELETE FROM catalog_pruning_policy; UPDATE catalog_items SET global_absent_syncs=100,retry_count=4,next_retry_at=1900000000,blocked_at='protected';"); + h.Db.Initialise(); var actual = (await h.Db.GetImportCatalogByIdAsync(item.Id))!; + Assert.Equal(0, actual.GlobalAbsentSyncs); Assert.Equal(4, actual.RetryCount); + Assert.Equal(1900000000, actual.NextRetryAt); Assert.Equal("protected", actual.BlockedAt); + Assert.Equal(item.StrmPath, actual.StrmPath); Assert.Null(actual.RemovedAt); + h.Sql("UPDATE catalog_items SET global_absent_syncs=2;"); h.Db.Initialise(); + Assert.Equal(2, (await h.Db.GetImportCatalogByIdAsync(item.Id))!.GlobalAbsentSyncs); + } + + [Theory] + [InlineData("external")] + [InlineData("owned")] + [InlineData("blocked")] + [InlineData("pinned")] + [InlineData("membership")] + [InlineData("played")] + [InlineData("saved")] + [InlineData("unknown-owner")] + public async Task ProtectedRowsNeverBecomeAbsentPruneCandidates(string protection) + { + using var h = new Harness(); var item = await h.Seed("ttprotected"); + h.Sql("UPDATE catalog_items SET global_absent_syncs=3;"); + h.Sql(protection switch { + "external" => "UPDATE catalog_items SET source='external_list';", + "owned" => "UPDATE catalog_items SET local_source='library',local_path='/owned/file.mkv',item_state=3;", + "blocked" => "UPDATE catalog_items SET blocked_at='protected';", + "pinned" => "UPDATE catalog_items SET item_state=5;", + "membership" => "INSERT INTO collection_membership(collection_name,aio_id,source,last_seen,created_at,updated_at) VALUES('Keep','ttprotected','aiostreams','now','now','now');", + "played" => "INSERT INTO playback_log(id,aio_id,resolution_mode,played_at) VALUES('played','ttprotected','cached','now');", + "saved" => "INSERT INTO media_items(id,primary_id_type,primary_id,media_type,title,status,saved,created_at,updated_at) VALUES('saved','imdb','ttprotected','movie','Keep','known',1,'now','now');", + "no-path" => "UPDATE catalog_items SET strm_path=NULL;", + _ => "UPDATE catalog_items SET local_source='unknown';" + }); + Assert.Empty(await h.Db.GetAbsentPruneCandidatesAsync()); + Assert.Equal(0, await h.Db.RetireAbsentCatalogItemsAsync(new[] { (item.Id, item.StrmPath, item.AioId) })); + Assert.Null((await h.Db.GetImportCatalogByIdAsync(item.Id))!.RemovedAt); + } + + [Theory] + [InlineData(null)] + [InlineData("")] + public async Task UnpublishedAbsentMetadataRetiresAfterThreeCompleteObservations(string? path) + { + using var h = new Harness(); var item=await h.Seed("ttunpublished"); + h.Sql("UPDATE catalog_items SET strm_path="+(path==null ? "NULL" : "''")+",local_path=NULL,local_source=NULL,item_state=6;"); + for (var i=0;i<3;i++) await h.Db.IncrementGlobalAbsentSyncsAsync(new() {"ttpresent"}); + var candidates=await h.Db.GetAbsentPruneCandidatesAsync(); Assert.Single(candidates); + Assert.Equal(path,candidates[0].StrmPath); Assert.Equal(item.Id,candidates[0].Id); + Assert.Equal(1,await h.Db.RetireAbsentCatalogItemsAsync(candidates)); + Assert.NotNull((await h.Db.GetImportCatalogByIdAsync(item.Id))!.RemovedAt); + } + + [Theory] + [InlineData("external_list", null, 6)] + [InlineData("aiostreams", "library", 3)] + [InlineData("aiostreams", null, 5)] + [InlineData("aiostreams", "unknown", 6)] + public async Task MetadataOnlyRetirementStillRespectsSourceOwnershipAndIntent(string source,string? owner,int state) + { + using var h = new Harness(); var item=await h.Seed("ttprotectedmetadata"); + item.Source=source; item.LocalSource=owner; item.StrmPath=null; item.ItemState=(ItemState)state; + h.Sql("DELETE FROM catalog_items;"); await h.Db.UpsertCatalogItemAsync(item); + h.Sql("UPDATE catalog_items SET global_absent_syncs=3;"); + Assert.Empty(await h.Db.GetAbsentPruneCandidatesAsync()); + Assert.Equal(0,await h.Db.RetireAbsentCatalogItemsAsync(new[] {(item.Id,item.StrmPath,item.AioId)})); + } + + [Fact] + public async Task MetadataOnlyRowsRetainWatchedAndActiveExternalTmdbAlias() + { + using var h=new Harness(); var item=await h.Seed("ttmetadata"); item.TmdbId="456"; + item.LocalSource=null; item.StrmPath=null; await h.Db.UpsertCatalogItemAsync(item); + // The normal sparse upsert preserves a materialized path; make the fixture truly unpublished. + h.Sql("UPDATE catalog_items SET strm_path=NULL,local_source=NULL,global_absent_syncs=3;"); + await h.Db.LogPlaybackAsync(new PlaybackEntry { AioId=item.AioId, ResolutionMode="cached" }); + Assert.Empty(await h.Db.GetAbsentPruneCandidatesAsync()); + h.Sql("DELETE FROM playback_log;"); + await h.Db.UpsertCatalogItemAsync(new CatalogItem { AioId="tmdb:456", TmdbId="456", Source="external_list", MediaType="movie",Title="Keep" }); + Assert.Empty(await h.Db.GetAbsentPruneCandidatesAsync()); + } + + [Theory] + [InlineData("same-imdb", "movie", true)] + [InlineData("different-imdb", "movie", true)] + [InlineData("different-imdb", "series", false)] + public async Task ActiveExternalAliasProtectsMatchingIdentityWithoutTmdbMovieTvCollision(string id, string type, bool protectedTitle) + { + using var h = new Harness(); var item = await h.Seed("same-imdb"); + item.TmdbId = "123"; await h.Db.UpsertCatalogItemAsync(item); + await h.Db.UpsertCatalogItemAsync(new CatalogItem { AioId = id, Source = "external_list", MediaType = type, Title = "List title", TmdbId = "123" }); + h.Sql("UPDATE catalog_items SET global_absent_syncs=3;"); + Assert.Equal(protectedTitle ? 0 : 1, (await h.Db.GetAbsentPruneCandidatesAsync()).Count); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task WatchedOrCollectedTmdbAliasKeepsBothBroadRows(bool collection) + { + using var h = new Harness(); var first=await h.Seed("ttfirst"); first.TmdbId="321"; await h.Db.UpsertCatalogItemAsync(first); + var alias=await h.Seed("tmdb:321"); alias.TmdbId="321"; await h.Db.UpsertCatalogItemAsync(alias); + h.Sql("UPDATE catalog_items SET global_absent_syncs=3;"); + h.Sql(collection ? "INSERT INTO collection_membership(collection_name,aio_id,source,last_seen,created_at,updated_at) VALUES('Keep','tmdb:321','aiostreams','now','now','now');" + : "INSERT INTO playback_log(id,aio_id,resolution_mode,played_at) VALUES('played','tmdb:321','cached','now');"); + Assert.Empty(await h.Db.GetAbsentPruneCandidatesAsync()); + } + + [Fact] + public async Task RetirementRechecksNewProtectionAndNeverDeletesAnotherSource() + { + using var h = new Harness(); var item = await h.Seed("ttstale"); + for (int i=0;i<2;i++) { await h.Db.IncrementGlobalAbsentSyncsAsync(new() { "ttother" }); Assert.Empty(await h.Db.GetAbsentPruneCandidatesAsync()); } + await h.Db.IncrementGlobalAbsentSyncsAsync(new() { "ttother" }); + var candidates = await h.Db.GetAbsentPruneCandidatesAsync(); Assert.Single(candidates); + var external = new CatalogItem { AioId=item.AioId, Source="external_list", MediaType="movie", Title="Keep" }; + await h.Db.UpsertCatalogItemAsync(external); + Assert.Equal(0, await h.Db.RetireAbsentCatalogItemsAsync(candidates)); + Assert.Null((await h.Db.GetImportCatalogByIdAsync(item.Id))!.RemovedAt); + h.Sql("UPDATE catalog_items SET removed_at='old' WHERE source='external_list';"); + Assert.Equal(1, await h.Db.RetireAbsentCatalogItemsAsync(candidates)); + Assert.NotNull((await h.Db.GetImportCatalogByIdAsync(item.Id))!.RemovedAt); + Assert.Equal("old", (await h.Db.GetImportCatalogByIdAsync(external.Id))!.RemovedAt); + } + + [Fact] + public void SkippedProvidersAndPartialCatalogFailuresCannotAuthorizeAbsence() + { + Assert.False(CatalogPruningPolicy.CanObserveAbsence(1,2)); Assert.False(CatalogPruningPolicy.CanObserveAbsence(0,0)); + Assert.True(CatalogPruningPolicy.CanObserveAbsence(2,2)); + var result = new CatalogFetchResult { Items = new() { new CatalogItem() } }; + result.CatalogOutcomes["success"] = new() { Succeeded=true }; + Assert.True(CatalogPruningPolicy.IsCompleteProviderSnapshot(result)); + result.CatalogOutcomes["failed"] = new() { Succeeded=false }; + Assert.False(CatalogPruningPolicy.IsCompleteProviderSnapshot(result)); + result.CatalogOutcomes.Remove("failed"); result.ProviderReachable=false; + Assert.False(CatalogPruningPolicy.IsCompleteProviderSnapshot(result)); + } + + [Theory] + [InlineData("null", false)] + [InlineData("missing-metas", false)] + [InlineData("repeat", false)] + [InlineData("empty", true)] + [InlineData("cap", true)] + public async Task PartialPageFailuresRetainFetchedDataWithoutAuthorizingPruning(string ending, bool complete) + { + int calls=0; var page = new AioStreamsCatalogResponse { Metas=new() { new AioStreamsMeta { Id="tt1234567", Type="movie", Name="Fixture" } } }; + var result = await AioStreamsCatalogProvider.FetchCatalogPagesAsync(new() { Id="popular", Type="movie" }, NullLogger.Instance, + ending=="cap" ? 1:2, CancellationToken.None, _ => Task.FromResult(++calls==1 ? page : ending switch { + "null" => null, "missing-metas" => new AioStreamsCatalogResponse { Metas=null! }, + "repeat" => page, _ => new AioStreamsCatalogResponse { Metas=new() } + })); + Assert.Single(result.Items); Assert.Equal(complete, result.Outcome.Succeeded); + } + + [Theory] + [InlineData("played")] + [InlineData("resume")] + [InlineData("favorite")] + [InlineData("history")] + public void LaterNativeVersionAndAnySeriesEpisodeProtectUserState(string state) + { + var first = new Movie(); var second = new Movie(); var series = new Series(); var episode = new Episode(); + var user=new User(); var data = state switch { "played"=>new UserItemData { PlayCount=1 }, "resume"=>new UserItemData { PlaybackPositionTicks=1 }, + "favorite"=>new UserItemData { IsFavorite=true }, _=>new UserItemData { LastPlayedDate=DateTimeOffset.UtcNow } }; + Assert.True(CatalogPruningPolicy.HasProtectedNativeItems(new BaseItem[] {first,second}, new[] {user}, _=>Array.Empty(), + (_, item)=>ReferenceEquals(item,second) ? data:new(), _=>false)); + Assert.True(CatalogPruningPolicy.HasProtectedNativeItems(new BaseItem[] {series}, new[] {user}, _=>new BaseItem[] {episode}, + (_, item)=>ReferenceEquals(item,episode) ? data:new(), _=>false)); + } + + [Fact] + public void UnknownNativeEvidenceAndOwnedMatchesRetainButUnwatchedManagedMatchesAreEligible() + { + var items=new BaseItem[] {new Movie()}; var users=new[] {new User()}; + Assert.True(CatalogPruningPolicy.HasProtectedNativeItems(null,users,_=>Array.Empty(),(_,_)=>new(),_=>false)); + Assert.True(CatalogPruningPolicy.HasProtectedNativeItems(items,users,_=>Array.Empty(),(_,_)=>throw new IOException(),_=>false)); + Assert.True(CatalogPruningPolicy.HasProtectedNativeItems(items,users,_=>Array.Empty(),(_,_)=>new(),_=>true)); + Assert.False(CatalogPruningPolicy.HasProtectedNativeItems(items,users,_=>Array.Empty(),(_,_)=>new(),_=>false)); + } + + private sealed class Harness : IDisposable + { + private readonly string _path=Path.Combine(Path.GetTempPath(),"pruning-"+Guid.NewGuid().ToString("N")); + public DatabaseManager Db {get;} + public Harness() { SqliteTestRuntime.EnsureInitialized(); Db=new(_path,NullLogger.Instance); Db.Initialise(); } + public async Task Seed(string id) { var item=new CatalogItem { AioId=id, Source="aiostreams", Title="Fixture", MediaType="movie", + LocalSource="strm", StrmPath="/managed/movies/"+id, ItemState=ItemState.Ready }; await Db.UpsertCatalogItemAsync(item); return item; } + public void Sql(string sql) { using var c=SQLite3.Open(Path.Combine(_path,"infinitedrive.db"),ConnectionFlags.ReadWrite,null,true); foreach (var statement in sql.Split(';', StringSplitOptions.RemoveEmptyEntries)) c.Execute(statement+";"); } + public void Dispose() { try { Directory.Delete(_path,true); } catch {} } + } +} diff --git a/docs/import-reconciliation.md b/docs/import-reconciliation.md index 1d733ad..b7cab3f 100644 --- a/docs/import-reconciliation.md +++ b/docs/import-reconciliation.md @@ -30,7 +30,29 @@ leave the library after they disappear from every catalog unless retained by the existing saved/list/collection or watched-content policy. A pruned title is not recreated until catalog re-entry or new authorized intent. Another active source keeps a shared title eligible. This feature does not change retention thresholds -or replace the pruning policy. Successful +or replace the pruning policy. + +Broad-feed pruning requires three fresh, complete observations of all configured +providers. An interval-skipped provider, failed catalog, missing page, repeated page +before its cap, or empty combined snapshot postpones absence counting and retirement. +Each complete sync increments an absent title once, regardless of catalog size. +On upgrade the invalid older absence counters are reset once, with a durable policy +marker; catalog history, retry/backoff, blocks and user state remain intact. +Automatic retirement is limited to broad-feed rows: validated managed STRMs and +unpublished metadata with no local-media path. Retiring metadata alone does not +authorize file deletion. External-list titles, +owned-media aliases, pins, collections, playback history, saves, favorites and partial +watch progress retain a title. IMDb and same-type TMDB aliases are checked together. +Native checks cover every version, every user and series episodes; an unavailable +check retains the title. Retirement rechecks database protections under the publication +lock. Existing orphan cleanup removes unreferenced STRMs afterward; retirement does +not recursively delete a title directory or prove its files have already disappeared. +Before upgrading, back up the database and configuration. The additive policy marker +and indexes can remain if binaries are rolled back. Older binaries still have the +unsafe pruning behavior: defer broad catalog sync while running them. Do not restore +an older database over newer playback, list, block or acquisition state. + +Successful native indexing requires matching provider identity, episode numbering, and a managed path. A file alone is awaiting indexing. Multiple versions count as one episode; an explicit native multi-episode range may cover multiple episode keys. @@ -193,6 +215,14 @@ status never creates tables, contacts providers, or triggers work. ## Verification and rollout +September 30, 2026 pruning safety verification: the pinned Emby 4.10.0.40 build +passed all 199 tests and published a release DLL. New cases exercise real SQLite +absence counting across more than two batches, the one-time policy reset, guarded +retirement including unpublished metadata, typed source aliases and late list +protection; controlled native adapters cover all versions, episodes, user history +and unavailable checks. This verifies the implementation, not production deletion, +playback or a completed capped-catalog transition. + The regression suite uses the actual Emby SQLite provider with temporary databases, plus controlled metadata/stream adapters. It covers unchanged inventory, partial failure, database reopen, automatic refill, block/prune exclusions, metadata failure, future/unknown