Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
191 changes: 107 additions & 84 deletions Data/DatabaseManager.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}

Expand Down Expand Up @@ -427,62 +436,94 @@ public async Task<List<CatalogItem>> GetCatalogItemsBySourceAsync(string source)
ReadCatalogItem);
}

/// <summary>
/// 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
/// <see cref="SoftDeleteCatalogItemsAsync"/>. Returns (strm_path, aio_id) pairs.
/// </summary>
public async Task<List<(string? StrmPath, string AioId)>> 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'))))";

/// <summary>Returns eligible broad-feed rows, including unpublished metadata, without retiring anything.</summary>
public async Task<List<(string Id, string? StrmPath, string AioId)>> 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;
}

/// <summary>
/// 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.
/// </summary>
public async Task<int> 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;
}

/// <summary>
Expand Down Expand Up @@ -558,7 +599,7 @@ await ExecuteWriteAsync(sql, cmd =>
}

/// <summary>
/// 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
Expand All @@ -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
Expand Down Expand Up @@ -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);

Expand Down
67 changes: 67 additions & 0 deletions Services/CatalogPruningPolicy.cs
Original file line number Diff line number Diff line change
@@ -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<BaseItem>? matches, IEnumerable<User> users,
Func<BaseItem, IEnumerable<BaseItem>?> episodes,
Func<User, BaseItem, UserItemData?> userData, Func<BaseItem, bool> 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; }
}
}
Loading
Loading