diff --git a/Data/DatabaseManager.ImportCoverage.cs b/Data/DatabaseManager.ImportCoverage.cs index 4518579..bbf599f 100644 --- a/Data/DatabaseManager.ImportCoverage.cs +++ b/Data/DatabaseManager.ImportCoverage.cs @@ -138,6 +138,15 @@ AND json_extract(e.value,'$.State') IN ('missing','retrying','awaiting_indexing' AND (json_extract(e.value,'$.LeaseUntil') IS NULL OR json_extract(e.value,'$.LeaseUntil')<=@now)) ORDER BY s.checked_at,c.id LIMIT 10;", c => BindText(c, "@now", now.ToString("o")), ReadCatalogItem); + public Task> GetDiversityImportCatalogAsync(DateTimeOffset now) => QueryListAsync( + @"SELECT c.* FROM import_coverage s JOIN catalog_items c ON c.id=json_extract(s.payload,'$.CatalogIds[0]') + WHERE json_extract(s.payload,'$.Exclusion')='' + AND EXISTS (SELECT 1 FROM json_each(s.payload,'$.Items') e + WHERE json_extract(e.value,'$.Diversity.Status') IN ('waiting','in_flight') + AND julianday(json_extract(e.value,'$.Diversity.NextAttempt'))<=julianday(@now) + AND (json_extract(e.value,'$.LeaseUntil') IS NULL OR julianday(json_extract(e.value,'$.LeaseUntil'))<=julianday(@now))) + ORDER BY s.checked_at,c.id LIMIT 10;", c => BindText(c, "@now", now.ToString("o")), ReadCatalogItem); + // Revisit unfinished existing versions before advancing the catalog cursor. Failed // episodes retain their own backoff; successfully refreshed siblings leave this queue. public Task> GetCatchUpImportCatalogAsync(DateTimeOffset started, DateTimeOffset now) => QueryListAsync( diff --git a/InfiniteDrive.csproj b/InfiniteDrive.csproj index ef7ce1a..0231cc6 100755 --- a/InfiniteDrive.csproj +++ b/InfiniteDrive.csproj @@ -2,9 +2,9 @@ net8.0 - 0.42.13.0 - 0.42.13.0 - 0.42.13.0 + 0.42.14.0 + 0.42.14.0 + 0.42.14.0 enable InfiniteDrive InfiniteDrive diff --git a/Models/ImportCoverage.cs b/Models/ImportCoverage.cs index 6104444..d2dc6a6 100644 --- a/Models/ImportCoverage.cs +++ b/Models/ImportCoverage.cs @@ -59,6 +59,8 @@ public sealed class ImportEpisode public string Failure { get; set; } = ""; public string? Lease { get; set; } public DateTimeOffset? LeaseUntil { get; set; } + // Additive JSON migration: old observations deserialize without queue state. + public ImportDiversityEntry? Diversity { get; set; } } /// Fresh per-file metadata bound to its current URL hash, without storing the URL. diff --git a/Models/ImportDiversityEntry.cs b/Models/ImportDiversityEntry.cs new file mode 100644 index 0000000..61d3b35 --- /dev/null +++ b/Models/ImportDiversityEntry.cs @@ -0,0 +1,78 @@ +using System; +using System.Collections.Generic; +using System.Linq; + +namespace InfiniteDrive.Models; + +/// One durable queue entry per coverage identity/episode, never a playback claim. +public sealed class ImportDiversityEntry +{ + public string Status { get; set; } = "waiting"; + public string Profile { get; set; } = ""; + public long Generation { get; set; } + public string Reason { get; set; } = "single_source"; + public int Streak { get; set; } + public int Attempts { get; set; } + public DateTimeOffset CreatedAt { get; set; } + public DateTimeOffset? CheckedAt { get; set; } + public DateTimeOffset? NextAttempt { get; set; } +} + +public static class ImportDiversityPolicy +{ + private static readonly string[] Labels = { "TorBox", "Usenet", "Real-Debrid", "AllDebrid", "Premiumize", + "Debrid-Link", "Offcloud", "EasyDebrid", "Debrider", "Easynews", "NZBDav", "AltMount", "StremThru" }; + public static string? Provider(string? label) => label?.ToLowerInvariant() is "nntp" or "stremionntp" + ? "Usenet" : Labels.FirstOrDefault(x => string.Equals(x, label, StringComparison.OrdinalIgnoreCase)); + + // Every retained path must have one unambiguous known label. Historical paths don't count. + public static List? Providers(ImportEpisode ep) + { + if (ep.Paths.Count == 0) return null; + var result = new HashSet(StringComparer.Ordinal); + foreach (var path in ep.Paths.Distinct(StringComparer.Ordinal)) + { + var labels = ep.Versions.Where(x => x.Path == path).Select(x => Provider(x.ServiceLabel)).Distinct().ToList(); + if (labels.Count != 1 || labels[0] == null) return null; + result.Add(labels[0]!); + } + return result.OrderBy(x => x, StringComparer.Ordinal).ToList(); + } + + public static bool Eligible(ImportCoverage coverage, ImportEpisode ep) => + coverage.Exclusion.Length == 0 && coverage.SnapshotStatus == "success" && ep.Expected && ep.Eligible && + ep.State == "indexed" && ep.Failure.Length == 0 && Providers(ep)?.Count == 1; + + public static void Reconcile(ImportCoverage coverage, ImportEpisode ep, string profile, DateTimeOffset now) + { + var providers = Providers(ep); + if (!Eligible(coverage, ep)) + { + if (ep.Diversity != null) + { ep.Diversity.Status = providers?.Count > 1 ? "resolved" : "retired"; + ep.Diversity.Reason = providers?.Count > 1 ? "multiple_sources" : "ineligible"; } + return; + } + if (ep.Diversity == null || ep.Diversity.Profile != profile || ep.Diversity.Generation != coverage.Generation) + ep.Diversity = new() { Profile = profile, Generation = coverage.Generation, CreatedAt = now, + NextAttempt = now.AddHours(6) }; + else if (ep.Diversity.Status is "retired" or "resolved") + { ep.Diversity.Status = "waiting"; ep.Diversity.Reason = "single_source"; + ep.Diversity.NextAttempt = now.AddHours(6); } + if (ep.Diversity.Status == "in_flight" && (ep.LeaseUntil == null || ep.LeaseUntil <= now)) + { ep.Diversity.Status = "waiting"; ep.Diversity.Reason = "interrupted"; + ep.Diversity.NextAttempt = now.AddMinutes(5); } + if (ep.Paths.Distinct().Count() >= 8) + { ep.Diversity.Status = "capacity"; ep.Diversity.Reason = "version_limit"; } + else if (ep.Diversity.Status == "capacity") + { ep.Diversity.Status = "waiting"; ep.Diversity.Reason = "single_source"; } + } + + public static void Defer(ImportDiversityEntry entry, string reason, DateTimeOffset now, + double jitter, DateTimeOffset cooldown) + { + entry.Streak++; entry.Status = "waiting"; entry.Reason = reason; entry.CheckedAt = now; + entry.NextAttempt = ImportCoveragePolicy.RetryAt(entry.Streak, + reason is "same_source" or "source_unavailable" or "no_valid_addition", now, jitter, cooldown); + } +} diff --git a/Models/ImportRunTelemetry.cs b/Models/ImportRunTelemetry.cs index b99d2a9..717f1f3 100644 --- a/Models/ImportRunTelemetry.cs +++ b/Models/ImportRunTelemetry.cs @@ -14,7 +14,7 @@ public sealed record ImportRunSnapshot(string Id, DateTimeOffset StartedAt, Date int RefreshAttempts, int ActiveLookups, int Matched, int EmptyResults, int TransportFailures, int LookupDeadlines, int Http429, int ProviderConfigurationFailures, int CancelledLookups, int Published, int Refreshed, int PublicationFailures, IReadOnlyDictionary Timings, - int HttpRequests = 0, int HttpRetries = 0); + int HttpRequests = 0, int HttpRetries = 0, int DiversityAttempts = 0, int DiversityPublished = 0); /// One bounded, thread-safe run observation. Contains no target URLs or exception messages. public sealed class ImportRunTelemetry @@ -32,6 +32,7 @@ public sealed class ImportRunTelemetry private int _checked, _missing, _refresh, _active, _matched, _empty, _transport, _deadlines, _http429, _providerConfiguration, _cancelled, _published, _refreshed, _publicationFailures; private int _httpRequests, _httpRetries; + private int _diversityAttempts, _diversityPublished; public ImportRunTelemetry(string id) { _id = id; } public static ImportRunTelemetry Start(string id) { @@ -49,7 +50,7 @@ public static void RecordCurrentHttp(bool retry) lock (current._gate) { current._httpRequests++; if (retry) current._httpRetries++; } } public static void RecordCurrentTiming(string stage, double seconds) => Volatile.Read(ref _current)?.RecordTiming(stage, seconds); - public void Attempt(bool refresh) { lock (_gate) { if (refresh) _refresh++; else _missing++; } } + public void Attempt(bool refresh, bool diversity = false) { lock (_gate) { if (diversity) _diversityAttempts++; else if (refresh) _refresh++; else _missing++; } } public void LookupStarted() { lock (_gate) _active++; } public void LookupFinished(string failure) { @@ -68,7 +69,7 @@ public void LookupFinished(string failure) } } } - public void Published(bool refresh) { lock (_gate) { _published++; if (refresh) _refreshed++; } } + public void Published(bool refresh, bool diversity = false) { lock (_gate) { _published++; if (refresh) _refreshed++; if (diversity) _diversityPublished++; } } public void PublicationFailed() { lock (_gate) _publicationFailures++; } public void RecordTiming(string stage, double seconds) { @@ -100,7 +101,7 @@ public ImportRunSnapshot Snapshot() return new(_id, _started, DateTimeOffset.UtcNow, _finished, _status, _phase, _title, _episode, Math.Round(_elapsed.Elapsed.TotalSeconds, 3), Math.Round(_phaseTime.Elapsed.TotalSeconds, 3), _checked, _missing, _refresh, _active, _matched, _empty, _transport, _deadlines, _http429, - _providerConfiguration, _cancelled, _published, _refreshed, _publicationFailures, times, _httpRequests, _httpRetries); + _providerConfiguration, _cancelled, _published, _refreshed, _publicationFailures, times, _httpRequests, _httpRetries, _diversityAttempts, _diversityPublished); } } private sealed class Stage diff --git a/PluginConfiguration.cs b/PluginConfiguration.cs index 7a319fd..3188d47 100755 --- a/PluginConfiguration.cs +++ b/PluginConfiguration.cs @@ -46,6 +46,7 @@ namespace InfiniteDrive public class PluginConfiguration : BasePluginConfiguration { [DataMember] public Models.ImportMode ImportRecoveryMode { get; set; } = Models.ImportMode.Observe; + [DataMember] public bool ImportProviderDiversityEnabled { get; set; } = true; [DataMember] public string ImportCatchUpStartedAt { get; set; } = ""; [DataMember] public string ImportCatchUpUntil { get; set; } = ""; [DataMember] public string ImportHouseholdTimezone { get; set; } = "America/Chicago"; diff --git a/README.md b/README.md index 84eef32..1fba9c3 100755 --- a/README.md +++ b/README.md @@ -9,9 +9,9 @@ imports titles and repairs missing files. A stream will appear. Probably. -**Current released version: 0.42.13** +**Current version: 0.42.14** -Release 0.42.13 adds administrator whole-title blocks and recoverable managed-stream cleanup. It targets **Emby 4.10.0.40**. +Release 0.42.14 adds eventual provider diversity while retaining the administrator whole-title blocks and recoverable managed-stream cleanup from 0.42.13. It targets **Emby 4.10.0.40**. ## How we run it @@ -234,6 +234,18 @@ when ready it uses the existing run budget. **Cancel pending reprocessing** leav in-flight work alone. Blocks, removals, owned media and disputed episode identities always win. It does not add a missing provider slot to already successful choices. +### Eventual source diversity + +With **Repair & refresh** enabled, Marvin gradually revisits indexed movies and episodes that have one known saved delivery source. Several TorBox quality variants still count as one source. A series is counted per episode; having TorBox on one episode and Usenet on another does not resolve either episode's gap. + +The queue defaults on. **Pause source diversity** stops its lookups/additions while essential repair continues; **Enable eventual source diversity** resumes normal scheduled processing. Entries, backoff and attempt history persist. Neither button triggers a pass or resets retries. There is no force-drain action. Unknown labels, blocked/removed/owned titles, unconfirmed identities and essential failures are excluded. + +New entries wait six hours. Missing items and ordinary refreshes have priority; diversity borrows at most two checks per normal pass or eight per catch-up pass from the existing shared allowance, serially. With stale health it can make one shared probe only if essential work has not already reserved a check. Same-source/empty diversity results escalate through 6h, 1d, 2d, 3d and 7d, plus up to 20% jitter and longer cooldowns. The queue persists beyond catch-up expiry; gradual processing does not guarantee a second source exists. + +A successful diversity check **adds** one validated variant and preserves every working file. New variants require a known different label and positive exact upstream size no greater than 40GB. An item with all eight version slots occupied waits for capacity; Marvin does not delete a working choice to force diversity. A saved addition does not prove playback or independent provider infrastructure. Ordinary refresh retains its existing replacement behavior and can change coverage later. + +Inspect dated queue state, reasons, streak and due time in native import status or the Failure Library. A one-source flag is not proof of a queue entry. **[Read the exact queue contract](docs/provider-diversity-queue.md)** for schema/leases, state transitions, health predicate, every ladder/counter, restart and partial-write recovery, operator controls, examples and rollback. + ### Marvin dashboard Open **Plugins → InfiniteDrive → Marvin** and use **Refresh dashboard**. diff --git a/Services/Api/ImportHealthService.cs b/Services/Api/ImportHealthService.cs index 93171c1..0ad192a 100644 --- a/Services/Api/ImportHealthService.cs +++ b/Services/Api/ImportHealthService.cs @@ -45,12 +45,14 @@ public async Task Get(ImportHealthRequest request) PipelinePhase = Plugin.Pipeline.Current, CurrentRun = ImportRunTelemetry.Current, NextCreditAt = await db.GetNextImportCreditAsync(DateTimeOffset.UtcNow, speed.AttemptsPerDay), RecentRuns = await db.GetImportRunHistoryAsync(), LastRun = db.GetMetadata("import_last_run"), CollectionHealth = db.GetMetadata("import_collection_health"), NextOffset = request.Offset + page.Count, SourceHealth = db.GetImportSourceHealth(), Reprocessing = await db.GetImportReprocessSummaryAsync(), + DiversityEnabled = Plugin.Instance.Configuration.ImportProviderDiversityEnabled, Items = page.Select(x => new { x.Identity, x.Title, Complete = x.Complete && x.SnapshotAt >= DateTimeOffset.UtcNow.AddHours(-6), x.SnapshotAt, x.CheckedAt, x.SnapshotStatus, x.ProviderStatus, x.Exclusion, Stale = x.SnapshotAt == null || x.SnapshotAt < DateTimeOffset.UtcNow.AddHours(-6), Expected = x.Items.Count(i => i.Eligible), Indexed = x.Items.Count(i => i.Eligible && i.State == "indexed"), Episodes = x.Items.Select(i => new { i.Key, i.Season, i.Episode, i.State, i.Eligibility, - i.Failure, i.Attempts, i.ConsecutiveFailures, i.LastFailureAt, i.NextAttempt, i.ObservedAt }) }) }; + i.Failure, i.Attempts, i.ConsecutiveFailures, i.LastFailureAt, i.NextAttempt, i.ObservedAt, + Providers = ImportDiversityPolicy.Providers(i), i.Diversity }) }) }; } public async Task Post(ImportActionRequest request) @@ -71,6 +73,12 @@ internal static async Task ApplyAsync(ImportActionRequest request) await db.EnsureImportCoverageAsync(); switch (request.Action) { + case "enable_diversity": + case "disable_diversity": + plugin.Configuration.ImportProviderDiversityEnabled = request.Action == "enable_diversity"; + plugin.SaveConfiguration(); + // No kickoff, early retry, queue deletion or repair-policy change. + return new { Status = "setting_changed", DiversityEnabled = plugin.Configuration.ImportProviderDiversityEnabled }; case "mode": if (!Enum.TryParse(request.Mode, false, out var mode) || !Enum.IsDefined(mode)) return new { Status = "invalid_mode" }; diff --git a/Services/ImportInventory.cs b/Services/ImportInventory.cs index 8a3de36..604c5a2 100644 --- a/Services/ImportInventory.cs +++ b/Services/ImportInventory.cs @@ -24,6 +24,8 @@ public sealed record ImportObservation(List Paths, List NativeId public interface IImportInventory { bool ProviderPaused => false; + string ProfileFingerprint => ""; + bool CanAddSource(CatalogItem item, ImportEpisode episode) => HasSavedFiles(episode); bool HasSavedFiles(ImportEpisode episode) => episode.Paths.Any(File.Exists); Task FetchAsync(CatalogItem item, CancellationToken ct); Task ObserveAsync(CatalogItem item, ImportEpisode episode, CancellationToken ct); @@ -32,6 +34,10 @@ public interface IImportInventory Task> ResolveCatchUpAsync(CatalogItem item, ImportEpisode episode, CancellationToken ct) => ResolveAsync(item, episode, ct); Task> PublishAsync(CatalogItem item, ImportEpisode episode, List versions, CancellationToken ct); + Task> ResolveDiversityAsync(CatalogItem item, ImportEpisode episode, CancellationToken ct) + => ResolveAsync(item, episode, ct); + Task> PublishAdditionAsync(CatalogItem item, ImportEpisode episode, List versions, CancellationToken ct) + => throw new NotSupportedException("additive_publication_unavailable"); void Notify(CatalogItem item); } @@ -78,6 +84,8 @@ public bool ProviderPaused } public bool IsOwned(CatalogItem item) => new OwnedMediaPreferenceService(_library, _logger).FindOwnedMedia(item) != null; + public string ProfileFingerprint + { get { using var client = AioStreamsClientFactory.Create(_logger); return client.ConfigurationFingerprint; } } public async Task FetchAsync(CatalogItem item, CancellationToken ct) { @@ -286,7 +294,10 @@ public async Task> ResolveAsync(CatalogItem item, ImportEp public Task> ResolveCatchUpAsync(CatalogItem item, ImportEpisode episode, CancellationToken ct) => ResolveCoreAsync(item, episode, true, ct); - private async Task> ResolveCoreAsync(CatalogItem item, ImportEpisode episode, bool maintenance, CancellationToken ct) + public Task> ResolveDiversityAsync(CatalogItem item, ImportEpisode episode, CancellationToken ct) + => ResolveCoreAsync(item, episode, true, ct, true); + + private async Task> ResolveCoreAsync(CatalogItem item, ImportEpisode episode, bool maintenance, CancellationToken ct, bool diversity = false) { using var client = AioStreamsClientFactory.Create(_logger); // Every Repair lookup uses one HTTP submission and the coordinator's @@ -303,7 +314,15 @@ private async Task> ResolveCoreAsync(CatalogItem item, Imp } if (client.LastHttpStatus == 429) throw new ImportHttpRateLimitException(); if (response == null) throw new IOException("stream_transport_unavailable"); - return VersionSelectorService.SelectBestVersions(StreamParser.ParseAll(response.Streams), + var streams = StreamParser.ParseAll(response.Streams); + if (diversity) + { + var existing = ImportDiversityPolicy.Providers(episode); + var additions = streams.Where(x => ImportDiversityPolicy.Provider(x.ServiceLabel) is { } label && + existing != null && !existing.Contains(label) && x.SizeBytes is > 0 and <= 40_000_000_000).ToList(); + if (additions.Count > 0) streams = additions; + } + return VersionSelectorService.SelectBestVersions(streams, _config.DesiredVersions, RuntimePolicy.EmbyVersionLimit, _config); } @@ -333,6 +352,45 @@ public static List BuildVersionEvidence(List path return evidence; } + public async Task> PublishAdditionAsync(CatalogItem item, ImportEpisode episode, + List versions, CancellationToken ct) + { + var (folder, name) = Destination(item, episode); + // Recheck all retained files against their own evidence before an additive write. + foreach (var path in episode.Paths) + { + if (!SafeManagedPath(folder, path) || !File.Exists(path) || new FileInfo(path).LinkTarget != null) + throw new IOException("retained_file_changed"); + var hash = Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(File.ReadAllText(path).Trim()))).ToLowerInvariant(); + if (!episode.Versions.Any(x => x.Path == path && x.UrlSha256 == hash)) + throw new IOException("retained_file_changed"); + } + var before = FindFiles(folder, name); + if (!before.ToHashSet(StringComparer.Ordinal).SetEquals(episode.Paths)) throw new IOException("retained_set_changed"); + var label = ImportDiversityPolicy.Providers(episode)?.SingleOrDefault(); + var addition = versions.FirstOrDefault(x => ImportDiversityPolicy.Provider(x.Stream.ServiceLabel) is { } provider && + provider != label && x.Stream.SizeBytes is > 0 and <= 40_000_000_000); + if (addition == null || before.Count >= RuntimePolicy.EmbyVersionLimit) return new(); + var added = await _writer.WriteAdditionalStrmAsync(folder, name, addition, ct); + // Keep old evidence verbatim; never bind new metadata to an old URL. + episode.Versions.AddRange(BuildVersionEvidence(new() { added }, new() { addition }, DateTimeOffset.UtcNow) + .Where(x => !episode.Versions.Any(old => old.Path == x.Path))); + return before.Append(added).Distinct(StringComparer.Ordinal).ToList(); + } + + public bool CanAddSource(CatalogItem item, ImportEpisode episode) + { + try + { + var (folder, name) = Destination(item, episode); + if (!FindFiles(folder, name).ToHashSet(StringComparer.Ordinal).SetEquals(episode.Paths)) return false; + return episode.Paths.All(path => SafeManagedPath(folder, path) && File.Exists(path) && + new FileInfo(path).LinkTarget == null && episode.Versions.Any(x => x.Path == path && + x.UrlSha256 == Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(File.ReadAllText(path).Trim()))).ToLowerInvariant())); + } + catch { return false; } + } + public void Notify(CatalogItem item) { // ReportFileSystemChanged in the writer discovers new children. The coordinator diff --git a/Services/ImportReconciliationService.cs b/Services/ImportReconciliationService.cs index 31711c7..598e422 100644 --- a/Services/ImportReconciliationService.cs +++ b/Services/ImportReconciliationService.cs @@ -23,12 +23,14 @@ public sealed class ImportReconciliationService private readonly TimeZoneInfo _timezone; private readonly Func _workBudget; private ImportRunTelemetry? _activeTelemetry; + private readonly Func _diversityEnabled; public ImportReconciliationService(DatabaseManager db, IImportInventory inventory, Func mode, TimeZoneInfo timezone, Func? clock = null, - Func? cooldown = null, Func? workBudget = null) + Func? cooldown = null, Func? workBudget = null, Func? diversityEnabled = null) { _db = db; _inventory = inventory; _mode = mode; _timezone = timezone; _workBudget = workBudget ?? (() => ImportWorkBudget.Normal); + _diversityEnabled = diversityEnabled ?? (() => true); _clock = clock ?? (() => DateTimeOffset.UtcNow); _cooldown = cooldown ?? (() => DateTimeOffset.MinValue); } public static ImportReconciliationService Create() @@ -38,7 +40,8 @@ public static ImportReconciliationService Create() p.Configuration, p.ProviderManager, p.StrmFileManager!), () => p.Configuration.ImportRecoveryMode, TimeZoneInfo.FindSystemTimeZoneById(p.Configuration.ImportHouseholdTimezone), cooldown: () => p.CooldownGate?.GlobalCooldownUntil ?? DateTimeOffset.MinValue, - workBudget: () => ImportWorkBudget.For(p.Configuration, DateTimeOffset.UtcNow)); + workBudget: () => ImportWorkBudget.For(p.Configuration, DateTimeOffset.UtcNow), + diversityEnabled: () => p.Configuration.ImportProviderDiversityEnabled); } public async Task RunAsync(CancellationToken ct, IReadOnlyList? selected = null) @@ -61,6 +64,9 @@ public async Task RunAsync(CancellationToken ct, IReadOnlyList? sel _activeTelemetry = telemetry; var pending = new List(); var deferred = new List(); + var diversity = new List(); + var diversityAttempts = 0; + var diversityPublished = 0; // Rotate the first claim on scarce credits across native slices. The other // lane can borrow unused credits after inventory discovery; no new quota. var priority = allowance.IsCatchUp && _db.GetMetadata("import_catch_up_last_priority") == "refresh" @@ -73,7 +79,7 @@ public async Task RunAsync(CancellationToken ct, IReadOnlyList? sel var healthLoaded = false; var rollingUsed = 0; - async Task ResolveAsync(CatalogItem item, ImportEpisode episode) + async Task ResolveAsync(CatalogItem item, ImportEpisode episode, bool diverse = false) { telemetry.LookupStarted(); var timer = Stopwatch.StartNew(); @@ -83,7 +89,7 @@ async Task ResolveAsync(CatalogItem item, ImportEpisode episod deadline.CancelAfter(TimeSpan.FromSeconds(allowance.IsCatchUp ? 120 : 60)); try { - versions = allowance.IsCatchUp + versions = diverse ? await _inventory.ResolveDiversityAsync(item, episode, deadline.Token) : allowance.IsCatchUp ? await _inventory.ResolveCatchUpAsync(item, episode, deadline.Token) : await _inventory.ResolveAsync(item, episode, deadline.Token); failure = versions.Count == 0 ? "source_unavailable" : ""; @@ -116,6 +122,52 @@ async Task CompleteAsync(ResolutionWork work) !episode.Eligible || !episode.Expected || !await IsAuthorizedAsync(live, work.Item, token) || _inventory.IsOwned(work.Item)) return; var now = _clock(); + if (work.Diversity) + { + var entry = episode.Diversity; + if (!_diversityEnabled() || entry == null || entry.Profile != _inventory.ProfileFingerprint || + entry.Generation != live.Generation || !ImportDiversityPolicy.Eligible(live, episode) || + DiversitySignature(episode) != work.Signature) + { + if (entry != null) { entry.Status = "waiting"; entry.Reason = "guard_changed"; entry.NextAttempt = now.AddMinutes(5); } + episode.Lease = null; episode.LeaseUntil = null; + await _db.SaveImportCoverageAsync(live, token); return; + } + sourceHealth.Record(result.Failure, now); + await _db.SaveImportSourceHealthAsync(sourceHealth, token); + var reason = result.Failure; + var providers = ImportDiversityPolicy.Providers(episode)!; + var additions = result.Versions?.Where(x => ImportDiversityPolicy.Provider(x.Stream.ServiceLabel) is { } label && + !providers.Contains(label) && x.Stream.SizeBytes is > 0 and <= 40_000_000_000).Take(1).ToList() ?? new(); + if (reason.Length == 0 && additions.Count == 0) + reason = result.Versions?.Any(x => providers.Contains(ImportDiversityPolicy.Provider(x.Stream.ServiceLabel) ?? "")) == true + ? "same_source" : "no_valid_addition"; + if (reason.Length == 0) + { + var timer = Stopwatch.StartNew(); + try + { + var paths = await _inventory.PublishAdditionAsync(work.Item, episode, additions, token); + if (paths.Count == 0) throw new IOException("addition_failed"); + episode.Paths = paths; + if (ImportDiversityPolicy.Providers(episode)?.Count is not >= 2) throw new IOException("addition_unverified"); + entry.Status = "resolved"; entry.Reason = "multiple_sources"; entry.CheckedAt = now; + entry.NextAttempt = null; entry.Streak = 0; + await _db.SaveImportCoverageAsync(live, token); + await RegisterPathsAsync(live, work.Item, episode, additions, token, true); + diversityPublished++; published++; telemetry.Published(false, true); + } + catch (OperationCanceledException) when (token.IsCancellationRequested) { throw; } + catch { reason = "publication_failed"; telemetry.PublicationFailed(); } + finally { telemetry.RecordTiming("publication", timer.Elapsed.TotalSeconds); } + } + if (reason.Length > 0) ImportDiversityPolicy.Defer(entry, reason, now, Random.Shared.NextDouble() * .2, _cooldown()); + episode.Lease = null; episode.LeaseUntil = null; + await _db.SaveImportCoverageAsync(live, token); + var ix = work.Coverage.Items.FindIndex(x => x.Key == episode.Key); + if (ix >= 0) work.Coverage.Items[ix] = episode; + return; + } var failure = result.Failure; sourceHealth.Record(failure, now); await _db.SaveImportSourceHealthAsync(sourceHealth, token); @@ -151,6 +203,9 @@ async Task CompleteAsync(ResolutionWork work) ImportCoveragePolicy.Failed(episode, failure, now, Random.Shared.NextDouble() * .2, _cooldown()); } episode.Lease = null; episode.LeaseUntil = null; + ImportDiversityPolicy.Reconcile(live, episode, _inventory.ProfileFingerprint, now); + if (failure.Length == 0 && episode.Diversity?.Status == "waiting") + episode.Diversity.NextAttempt = now.AddHours(6); await _db.SaveImportCoverageAsync(live, token); if (work.Reprocess) await _db.SetImportReprocessStatusAsync(live.Identity, episode.Key, failure.Length == 0 ? "published" : "failed", now, token); @@ -207,25 +262,45 @@ async Task DispatchAsync(ResolutionCandidate candidate) episode.NextAttempt > now && !reprocess || episode.LeaseUntil > now || !await IsAuthorizedAsync(live, candidate.Item, token) || _inventory.IsOwned(candidate.Item)) return; if (reprocess && !sourceHealth.RecoveryReady(now) && reprocessAttempts >= 1) return; - if (candidate.Upgrade && !allowance.NeedsRefresh(episode, now) && !reprocess) return; + if (candidate.Diversity) + { + var entry = episode.Diversity; + if (!_diversityEnabled() || !_inventory.HasSavedFiles(episode) || !ImportDiversityPolicy.Eligible(live, episode) || entry == null || + entry.Profile != _inventory.ProfileFingerprint || entry.Generation != live.Generation || + entry.Status != "waiting" || entry.NextAttempt > now || + episode.Paths.Distinct().Count() >= RuntimePolicy.EmbyVersionLimit || + diversityAttempts >= (allowance.IsCatchUp ? 8 : 2) || + !sourceHealth.RecoveryReady(now) && attempts >= 1) return; + // A diverse lookup is additive work, never an early repair/reprocess retry. + if (reprocess) return; + if (!_inventory.CanAddSource(candidate.Item, episode)) + { + entry.Reason = "retained_file_changed"; entry.NextAttempt = now.AddHours(6); + await _db.SaveImportCoverageAsync(live, token); return; + } + entry.Status = "in_flight"; entry.Attempts++; entry.CheckedAt = now; + diversityAttempts++; + } + if (candidate.Upgrade && !candidate.Diversity && !allowance.NeedsRefresh(episode, now) && !reprocess) return; attempts++; if (candidate.Upgrade) upgrades++; var lease = Guid.NewGuid().ToString("N"); episode.Lease = lease; episode.LeaseUntil = now.AddMinutes(10); - if (!candidate.Upgrade) episode.InitialFailure = true; - episode.Attempts++; + if (!candidate.Upgrade && !candidate.Diversity) episode.InitialFailure = true; + if (!candidate.Diversity) episode.Attempts++; await _db.RecordImportAttemptAsync(lease, now, token); if (reprocess) { reprocessAttempts++; await _db.SetImportReprocessStatusAsync(live.Identity, episode.Key, "in_flight", now, token); } - telemetry.Attempt(candidate.Upgrade); + telemetry.Attempt(candidate.Upgrade, candidate.Diversity); await _db.SaveImportCoverageAsync(live, token); var index = candidate.Coverage.Items.FindIndex(x => x.Key == episode.Key); if (index >= 0) candidate.Coverage.Items[index] = episode; start = () => new(candidate.Coverage, candidate.Item, episode, live.Generation, lease, - candidate.Upgrade, episode.State, reprocess, ResolveAsync(candidate.Item, episode)); + candidate.Upgrade, episode.State, reprocess, candidate.Diversity, DiversitySignature(episode), + ResolveAsync(candidate.Item, episode, candidate.Diversity)); } finally { MutationGate.Release(); } pending.Add(start()); @@ -235,6 +310,8 @@ async Task DispatchAsync(ResolutionCandidate candidate) try { await _db.EnsureImportCoverageAsync(token); + await _db.PersistMetadataAsync("import_diversity_schema", "1", token); + await _db.PersistMetadataAsync("import_diversity_enabled", _diversityEnabled() ? "true" : "false", token); sourceHealth = _db.GetImportSourceHealth(); healthLoaded = true; rollingUsed = await _db.GetRecentImportAttemptsAsync(_clock()); @@ -265,7 +342,7 @@ async Task DispatchAsync(ResolutionCandidate candidate) var catchUp = allowance.CatchUpStartedAt.HasValue ? await _db.GetCatchUpImportCatalogAsync(allowance.CatchUpStartedAt.Value, _clock()) : new List(); - page = (await _db.GetReprocessImportCatalogAsync()).Concat(await _db.GetDueImportCatalogAsync(_clock())).Concat(catchUp).Concat(page) + page = (await _db.GetReprocessImportCatalogAsync()).Concat(await _db.GetDueImportCatalogAsync(_clock())).Concat(catchUp).Concat(page).Concat(await _db.GetDiversityImportCatalogAsync(_clock())) .DistinctBy(x => x.Id).ToList(); } var mayHaveRefresh = page.Any(x => !string.IsNullOrEmpty(x.StrmPath)); @@ -308,7 +385,10 @@ async Task DispatchAsync(ResolutionCandidate candidate) if (coverage.Exclusion.Length > 0) { foreach (var e in coverage.Items) + { await _db.SetImportReprocessStatusAsync(coverage.Identity, e.Key, "cancelled", now, token); + ImportDiversityPolicy.Reconcile(coverage, e, _inventory.ProfileFingerprint, now); + } coverage.CheckedAt = now; await SaveObservedAsync(coverage, token); continue; @@ -351,6 +431,8 @@ async Task DispatchAsync(ResolutionCandidate candidate) if (coverage.Items.Count == 0 && !ImportInventory.IsSeries(item)) coverage.Items.Add(new ImportEpisode { Key = "movie" }); + foreach (var entryEpisode in coverage.Items) + ImportDiversityPolicy.Reconcile(coverage, entryEpisode, _inventory.ProfileFingerprint, now); var due = coverage.Items.Where(x => x.Expected && x.Eligible && (!x.LeaseUntil.HasValue || x.LeaseUntil <= now) && (!x.NextAttempt.HasValue || x.NextAttempt <= now) && (x.State is "missing" or "retrying" || allowance.NeedsRefresh(x, now))).ToList(); @@ -391,6 +473,7 @@ async Task DispatchAsync(ResolutionCandidate candidate) episode.State = ImportCoveragePolicy.Classify(episode, file, observation.NativeIds.Count > 0, observation.Conflict, now); episode.ObservedAt = now; + ImportDiversityPolicy.Reconcile(coverage, episode, _inventory.ProfileFingerprint, now); if (episode.State == "indexed") { episode.LastSuccess = now; @@ -431,6 +514,9 @@ async Task DispatchAsync(ResolutionCandidate candidate) deferred.Add(candidate); else await DispatchAsync(candidate); } + foreach (var entryEpisode in coverage.Items) + if (ImportDiversityPolicy.Eligible(coverage, entryEpisode) && entryEpisode.Diversity is { Status: "waiting" } entry && entry.NextAttempt <= now) + diversity.Add(new(coverage, item, entryEpisode, coverage.Generation, false, true)); await SaveObservedAsync(coverage, token); scanned++; if (coverage.SnapshotStatus == "success" && coverage.Items.Any(x => x.ObservedAt.HasValue)) @@ -457,6 +543,13 @@ async Task DispatchAsync(ResolutionCandidate candidate) } telemetry.Work("waiting_for_sources"); while (pending.Count > 0) await DrainOneAsync(); + // Lower priority than every discovered missing/refresh candidate, serial publication. + foreach (var candidate in diversity.OrderBy(x => x.Episode.Diversity!.NextAttempt).ThenBy(x => x.Coverage.Identity, StringComparer.Ordinal).ThenBy(x => x.Episode.Key, StringComparer.Ordinal)) + { + await DispatchAsync(candidate); + while (pending.Count > 0) await DrainOneAsync(); + if (diversityAttempts >= (allowance.IsCatchUp ? 8 : 2)) break; + } } catch (OperationCanceledException) { status = ct.IsCancellationRequested ? "cancelled" : "budget_deferred"; } catch { status = "failed"; throw; } @@ -479,6 +572,18 @@ async Task DispatchAsync(ResolutionCandidate candidate) var episode = live?.Items.FirstOrDefault(x => x.Key == work.Episode.Key); if (live == null || episode == null || episode.Lease != work.Lease) continue; var result = await work.Resolution; + if (work.Diversity) + { + if (episode.Diversity != null) + { + if (result.Failure.Length > 0 && result.Failure != "cancelled") + { ImportDiversityPolicy.Defer(episode.Diversity, result.Failure, _clock(), Random.Shared.NextDouble() * .2, _cooldown()); sourceHealth.Record(result.Failure, _clock()); } + else { episode.Diversity.Status = "waiting"; episode.Diversity.Reason = "slice_cancelled"; + episode.Diversity.NextAttempt = _clock().AddMinutes(5); } + } + episode.Lease = null; episode.LeaseUntil = null; + await _db.SaveImportCoverageAsync(live, CancellationToken.None); continue; + } if (result.Failure.Length > 0 && result.Failure != "cancelled") { ImportCoveragePolicy.Failed(episode, result.Failure, _clock(), Random.Shared.NextDouble() * .2, _cooldown()); sourceHealth.Record(result.Failure, _clock()); } @@ -494,7 +599,7 @@ async Task DispatchAsync(ResolutionCandidate candidate) telemetry.Finish(status); var report = System.Text.Json.JsonSerializer.Serialize(new { Id = runId, Status = status, FinishedAt = _clock(), Scanned = scanned, Attempts = attempts, Metadata = metadata, Published = published, - Upgrades = upgrades, Refreshed = refreshed, ReprocessAttempts = reprocessAttempts, SourceHealth = sourceHealth, SourceFailures = sourceFailures, + Upgrades = upgrades, Refreshed = refreshed, ReprocessAttempts = reprocessAttempts, DiversityAttempts = diversityAttempts, DiversityPublished = diversityPublished, SourceHealth = sourceHealth, SourceFailures = sourceFailures, TransportFailures = transportFailures, Priority = allowance.IsCatchUp ? priority : "normal", ElapsedSeconds = elapsed.Elapsed.TotalSeconds, Speed = allowance.IsCatchUp ? "catch_up" : "normal", allowance.AttemptsPerDay, allowance.Parallelism, allowance.CatchUpStartedAt, allowance.CatchUpUntil, Analytics = telemetry.Snapshot() }); @@ -505,10 +610,14 @@ async Task DispatchAsync(ResolutionCandidate candidate) } } - private sealed record ResolutionCandidate(ImportCoverage Coverage, CatalogItem Item, ImportEpisode Episode, long Generation, bool Upgrade); + private static string DiversitySignature(ImportEpisode episode) => string.Join("\n", + episode.Paths.OrderBy(x => x, StringComparer.Ordinal).Select(path => path + "|" + string.Join(",", + episode.Versions.Where(x => x.Path == path).Select(x => x.UrlSha256 + ":" + x.ServiceLabel).OrderBy(x => x, StringComparer.Ordinal)))); + + private sealed record ResolutionCandidate(ImportCoverage Coverage, CatalogItem Item, ImportEpisode Episode, long Generation, bool Upgrade, bool Diversity = false); private sealed record ResolutionResult(List? Versions, string Failure); private sealed record ResolutionWork(ImportCoverage Coverage, CatalogItem Item, ImportEpisode Episode, - long Generation, string Lease, bool Upgrade, string PreviousState, bool Reprocess, Task Resolution); + long Generation, string Lease, bool Upgrade, string PreviousState, bool Reprocess, bool Diversity, string Signature, Task Resolution); internal static async Task IsBlockedAsync(DatabaseManager db, CatalogItem item, CancellationToken ct) { @@ -549,7 +658,7 @@ private async Task SaveObservedAsync(ImportCoverage state, CancellationToken ct) finally { MutationGate.Release(); _activeTelemetry?.RecordTiming("checkpoint", timer.Elapsed.TotalSeconds); } } - private async Task RegisterPathsAsync(ImportCoverage state, CatalogItem item, ImportEpisode episode, List versions, CancellationToken ct) + private async Task RegisterPathsAsync(ImportCoverage state, CatalogItem item, ImportEpisode episode, List versions, CancellationToken ct, bool additive = false) { var directory = Path.GetDirectoryName(episode.Paths[0])!; var titleRoot = episode.Season.HasValue ? Path.GetDirectoryName(directory)! : directory; @@ -560,7 +669,10 @@ private async Task RegisterPathsAsync(ImportCoverage state, CatalogItem item, Im current.StrmPath = titleRoot; current.LocalSource = "strm"; current.LocalPath = titleRoot; current.ItemState = ItemState.Written; - current.SelectedVersionsJson = StrmFileManager.SerializeVersions(versions); + current.SelectedVersionsJson = additive ? System.Text.Json.JsonSerializer.Serialize( + StrmFileManager.DeserializeVersions(current.SelectedVersionsJson).Concat( + StrmFileManager.DeserializeVersions(StrmFileManager.SerializeVersions(versions))).DistinctBy(x => x.Url)) + : StrmFileManager.SerializeVersions(versions); current.LastVersionRefreshAt = _clock().ToString("o"); current.UpdatedAt = _clock().ToString("o"); if (ImportInventory.IsSeries(item)) diff --git a/Services/StrmFileManager.cs b/Services/StrmFileManager.cs index 4512fad..9ae7260 100644 --- a/Services/StrmFileManager.cs +++ b/Services/StrmFileManager.cs @@ -159,6 +159,48 @@ public async Task WriteOrReplaceStrmFilesAsync( finally { gate.Release(); } } + /// Add one deterministic variant without overwriting or deleting any existing file. + public async Task WriteAdditionalStrmAsync(string folder, string name, + SelectedVersion version, CancellationToken ct) + { + if (!Uri.TryCreate(version.Stream.Url, UriKind.Absolute, out var uri) || + uri.Scheme is not ("http" or "https") || version.Stream.SizeBytes is not (> 0 and <= 40_000_000_000) || + ImportDiversityPolicy.Provider(version.Stream.ServiceLabel) is not { } provider) + throw new InvalidOperationException("invalid_additional_source"); + var gate = _rewriteLocks.GetOrAdd(Path.GetFullPath(folder), _ => new SemaphoreSlim(1, 1)); + await gate.WaitAsync(ct); + string? tmp = null; + try + { + ct.ThrowIfCancellationRequested(); + if (!Directory.Exists(folder) || new DirectoryInfo(folder).LinkTarget != null) + throw new IOException("destination_changed"); + var hash = Convert.ToHexString(System.Security.Cryptography.SHA256.HashData(Encoding.UTF8.GetBytes(version.Stream.Url))).ToLowerInvariant(); + var path = Path.Combine(folder, $"{SanitiseFileName(name)} - {SanitiseFileName(version.VersionLabel)} - addition {hash[..16]}.strm"); + if (File.Exists(path)) + { + if (new FileInfo(path).LinkTarget != null || File.ReadAllText(path) != version.Stream.Url) + throw new IOException("addition_collision"); + return path; + } + if (ImportInventory.FindFiles(folder, name).Count >= RuntimePolicy.EmbyVersionLimit) + throw new IOException("version_limit"); + tmp = path + "." + Guid.NewGuid().ToString("N") + ".tmp"; + await File.WriteAllTextAsync(tmp, version.Stream.Url, new UTF8Encoding(false), ct); + ct.ThrowIfCancellationRequested(); + File.Move(tmp, path, overwrite: false); + tmp = null; + if (File.ReadAllText(path) != version.Stream.Url) throw new IOException("addition_verification_failed"); + NotifyLibraryMonitor(path); + return path; + } + finally + { + if (tmp != null && File.Exists(tmp)) File.Delete(tmp); + gate.Release(); + } + } + /// /// Deletes the primary .strm file at plus any /// multi-version variants sharing the same base name (e.g. "Title - 1080p.strm"), diff --git a/Tests/ImportDiversityTests.cs b/Tests/ImportDiversityTests.cs new file mode 100644 index 0000000..b94595b --- /dev/null +++ b/Tests/ImportDiversityTests.cs @@ -0,0 +1,119 @@ +using System; +using System.IO; +using System.Linq; +using System.Reflection; +using System.Threading.Tasks; +using InfiniteDrive.Models; +using InfiniteDrive.Services; +using Xunit; + +namespace InfiniteDrive.Tests; + +public sealed class ImportDiversityTests +{ + private static readonly DateTimeOffset Now = DateTimeOffset.Parse("2026-10-03T18:00:00Z"); + private static ImportEpisode Single() => new() { Key = "movie", Expected = true, Eligible = true, + State = "indexed", Paths = new() { "/fixture/one.strm", "/fixture/two.strm" }, Versions = new() + { new("/fixture/one.strm", "a", "TorBox", "4K", null, null, Now), + new("/fixture/two.strm", "b", "torbox", "1080p", null, 100, Now), + new("/retired/old.strm", "c", "Usenet", "720p", null, 100, Now) } }; + private static ImportCoverage Coverage(ImportEpisode ep) => new() { Identity = "movie:imdb:tt99999", SnapshotStatus = "success", Items = new() { ep } }; + + [Fact] public void QualityVariantsAndHistoricalSourcesDoNotInflateCoverage() + { + var ep = Single(); Assert.Equal(new[] { "TorBox" }, ImportDiversityPolicy.Providers(ep)); + ep.Versions[1] = ep.Versions[1] with { ServiceLabel = "Usenet" }; + Assert.Equal(2, ImportDiversityPolicy.Providers(ep)!.Count); + } + [Theory] [InlineData("")] [InlineData("unknown")] [InlineData("https://private.invalid/key")] + public void UnknownLabelsNeverProveOneProvider(string label) + { var ep = Single(); ep.Versions[1] = ep.Versions[1] with { ServiceLabel = label }; Assert.Null(ImportDiversityPolicy.Providers(ep)); } + [Fact] public void ConflictingLabelsAndMissingEvidenceStayUnknown() + { + var ep = Single(); ep.Versions.Add(ep.Versions[0] with { ServiceLabel = "Usenet" }); + Assert.Null(ImportDiversityPolicy.Providers(ep)); ep.Versions.Clear(); Assert.Null(ImportDiversityPolicy.Providers(ep)); + } + [Fact] public void QueueIsIdempotentAndProfileGenerationChangesRestartItsOwnWaitingWindow() + { + var ep = Single(); var c = Coverage(ep); + ImportDiversityPolicy.Reconcile(c, ep, "profile1", Now); var first = ep.Diversity; + ImportDiversityPolicy.Reconcile(c, ep, "profile1", Now.AddHours(1)); Assert.Same(first, ep.Diversity); + Assert.Equal(Now.AddHours(6), ep.Diversity!.NextAttempt); + c.Generation++; ImportDiversityPolicy.Reconcile(c, ep, "profile1", Now.AddHours(1)); + Assert.NotSame(first, ep.Diversity); Assert.Equal(c.Generation, ep.Diversity!.Generation); + var second = ep.Diversity; ImportDiversityPolicy.Reconcile(c, ep, "profile2", Now.AddHours(2)); + Assert.NotSame(second, ep.Diversity); Assert.Equal("profile2", ep.Diversity!.Profile); + } + [Theory] [InlineData("excluded")] [InlineData("index_mismatch")] [InlineData("retrying")] [InlineData("awaiting_indexing")] + public void NonIndexedStatesRetireInsteadOfConsumingDiversityCredits(string state) + { + var ep = Single(); var c = Coverage(ep); ImportDiversityPolicy.Reconcile(c, ep, "p", Now); + ep.State = state; ImportDiversityPolicy.Reconcile(c, ep, "p", Now); + Assert.Equal("retired", ep.Diversity!.Status); + } + [Fact] public void ActualFailuresAndBlocksWinOverDiversity() + { + var ep = Single(); var c = Coverage(ep); ImportDiversityPolicy.Reconcile(c, ep, "p", Now); + ep.Failure = "source_unavailable"; ep.NextAttempt = Now.AddDays(3); + ImportDiversityPolicy.Reconcile(c, ep, "p", Now); Assert.Equal("retired", ep.Diversity!.Status); + Assert.Equal(Now.AddDays(3), ep.NextAttempt); ep.Failure = ""; c.Exclusion = "not_authorized"; + Assert.False(ImportDiversityPolicy.Eligible(c, ep)); + } + [Fact] public void BackoffEscalatesWithoutTouchingRepairStateAndHonorsCooldown() + { + var ep = Single(); var c = Coverage(ep); ImportDiversityPolicy.Reconcile(c, ep, "p", Now); + var hours = new[] { 6, 24, 48, 72, 168, 168 }; + foreach (var h in hours) { ImportDiversityPolicy.Defer(ep.Diversity!, "same_source", Now, 0, DateTimeOffset.MinValue); Assert.Equal(Now.AddHours(h), ep.Diversity!.NextAttempt); } + ImportDiversityPolicy.Defer(ep.Diversity!, "http_429", Now, .2, Now.AddDays(20)); + Assert.Equal(Now.AddDays(20), ep.Diversity!.NextAttempt); Assert.Equal("", ep.Failure); Assert.Null(ep.NextAttempt); + } + [Fact] public void InterruptedLeaseRecoversButLiveLeaseDoesNot() + { + var ep = Single(); var c = Coverage(ep); ImportDiversityPolicy.Reconcile(c, ep, "p", Now); + ep.Diversity!.Status = "in_flight"; ep.LeaseUntil = Now.AddMinutes(1); + ImportDiversityPolicy.Reconcile(c, ep, "p", Now); Assert.Equal("in_flight", ep.Diversity.Status); + ImportDiversityPolicy.Reconcile(c, ep, "p", Now.AddMinutes(2)); + Assert.Equal("waiting", ep.Diversity.Status); Assert.Equal(Now.AddMinutes(7), ep.Diversity.NextAttempt); + } + [Fact] public void CapacityDoesNotSpendARequestOrEraseAnUnknownSizeVersion() + { + var ep=Single();var c=Coverage(ep); + for(var i=2;i<8;i++) {var path=$"/fixture/{i}.strm";ep.Paths.Add(path);ep.Versions.Add(new(path,"hash","TorBox","1080p",null,null,Now));} + ImportDiversityPolicy.Reconcile(c,ep,"p",Now); + Assert.Equal("capacity",ep.Diversity!.Status);Assert.Equal(0,ep.Diversity.Attempts);Assert.Equal(8,ep.Paths.Count); + ep.Paths.RemoveAt(7);ImportDiversityPolicy.Reconcile(c,ep,"p",Now.AddHours(7)); + Assert.Equal("waiting",ep.Diversity.Status);Assert.Equal(Now.AddHours(6),ep.Diversity.NextAttempt); + } + [Fact] public void MissingLeaseAfterInterruptionDoesNotStrandEntry() + { + var ep=Single();var c=Coverage(ep);ImportDiversityPolicy.Reconcile(c,ep,"p",Now); + ep.Diversity!.Status="in_flight";ImportDiversityPolicy.Reconcile(c,ep,"p",Now); + Assert.Equal("waiting",ep.Diversity.Status);Assert.Equal(Now.AddMinutes(5),ep.Diversity.NextAttempt); + } + [Fact] public void TelemetrySeparatesDiversityFromMissingAndRefreshAttempts() + { + var t=new ImportRunTelemetry("fixture");t.Attempt(false,true);t.Published(false,true); + var s=t.Snapshot();Assert.Equal(1,s.DiversityAttempts);Assert.Equal(0,s.MissingAttempts); + Assert.Equal(0,s.RefreshAttempts);Assert.Equal(1,s.DiversityPublished);Assert.Equal(1,s.Published);Assert.Equal(0,s.Refreshed); + } + [Fact] public async Task AdditiveWriterPreservesUnknownSizeWorkingFilesAndNeighborsAndIsIdempotent() + { + var folder = Path.Combine(Path.GetTempPath(), "diversity-" + Guid.NewGuid().ToString("N")); Directory.CreateDirectory(folder); + try + { + var log = DispatchProxy.Create(); + var writer = new StrmFileManager(log); + var old = Path.Combine(folder, "E1 - TorBox.strm"); File.WriteAllText(old, "https://example.invalid/old"); + var neighbor = Path.Combine(folder, "E10.strm"); File.WriteAllText(neighbor, "https://example.invalid/neighbor"); + var v = new SelectedVersion { VersionLabel = "1080p - Usenet", Stream = new() { Url = "https://example.invalid/new", ServiceLabel = "Usenet", SizeBytes = 100 } }; + var added = await writer.WriteAdditionalStrmAsync(folder, "E1", v, default); + Assert.Equal(added, await writer.WriteAdditionalStrmAsync(folder, "E1", v, default)); + Assert.Equal(2, ImportInventory.FindFiles(folder, "E1").Count); + Assert.Equal("https://example.invalid/old", File.ReadAllText(old)); Assert.Equal("https://example.invalid/neighbor", File.ReadAllText(neighbor)); + Assert.Empty(Directory.GetFiles(folder, "*.tmp")); + v.Stream.SizeBytes = null; + await Assert.ThrowsAsync(() => writer.WriteAdditionalStrmAsync(folder, "E1", v, default)); + } + finally { Directory.Delete(folder, true); } + } +} diff --git a/Tests/ImportReconciliationTests.cs b/Tests/ImportReconciliationTests.cs index d9f32c8..1475cad 100644 --- a/Tests/ImportReconciliationTests.cs +++ b/Tests/ImportReconciliationTests.cs @@ -726,6 +726,74 @@ public class NullLogProxy : System.Reflection.DispatchProxy } } + private static async Task SeedDiversity(Harness h, int count = 1) + { + h.Inventory.Count = count; h.Engage(); + foreach (var n in Enumerable.Range(1, count)) h.Inventory.Files.Add(n); + await h.Seed(); await h.Run(ImportMode.Observe); + var c = await h.State(); + foreach (var e in c.Items) + { + e.LastVersionRefresh = h.Now; e.Versions = new() { new(e.Paths[0], "saved", "TorBox", "1080p", null, null, h.Now) }; + ImportDiversityPolicy.Reconcile(c, e, h.Inventory.ProfileFingerprint, h.Now); + e.Diversity!.NextAttempt = h.Now; + } + await h.Db.SaveImportCoverageAsync(c); + } + [Fact] public async Task DiversitySameSourcePersistsWithoutLoopOrRepairFailure() + { + using var h = new Harness(); await SeedDiversity(h); h.Inventory.Service = "TorBox"; + await h.Run(); var e = (await h.State()).Items[0]; + Assert.Equal(1, h.Inventory.Resolutions); Assert.Equal(0, h.Inventory.Published); + Assert.Equal("same_source", e.Diversity!.Reason); Assert.Equal("", e.Failure); Assert.Null(e.NextAttempt); + h.Db = new DatabaseManager(h.Directory.Path, NullLogger.Instance); h.Db.Initialise(); + await h.Run(); Assert.Equal(1, h.Inventory.Resolutions); Assert.Equal(1, await h.Db.GetRecentImportAttemptsAsync(h.Now)); + } + [Fact] public async Task DiversityAddsSecondSourceWithoutReplacementAndCompletesOnlyThatEpisode() + { + using var h = new Harness(); await SeedDiversity(h, 2); h.Inventory.Service = "Usenet"; + await h.Run(); var items = (await h.State()).Items; + Assert.Equal(1, h.Inventory.Added); Assert.Equal(0, h.Inventory.Published); + Assert.Equal(2, items[0].Paths.Count); Assert.Equal("resolved", items[0].Diversity!.Status); + Assert.Equal("waiting", items[1].Diversity!.Status); Assert.Equal(0, items[0].Attempts); + } + [Fact] public async Task DiversityWaitsForCooldownAndDisabledSettingAndExhaustedCredit() + { + using var h = new Harness(); await SeedDiversity(h); h.Inventory.Service = "Usenet"; + h.Cooldown = h.Now.AddHours(1); await h.Run(); Assert.Equal(0, h.Inventory.Resolutions); + h.Cooldown = DateTimeOffset.MinValue; h.Config.ImportProviderDiversityEnabled = false; + await h.Run(); Assert.Equal(0, h.Inventory.Resolutions); + h.Config.ImportProviderDiversityEnabled = true; + await new ImportReconciliationService(h.Db,h.Inventory,()=>ImportMode.Repair,TimeZoneInfo.Utc,()=>h.Now, + workBudget:()=>new ImportWorkBudget(120,20,5,5,0,h.Now,h.Now.AddDays(7))).RunAsync(default,new[]{h.Item}); + Assert.Equal(0,h.Inventory.Resolutions); Assert.Equal("waiting",(await h.State()).Items[0].Diversity!.Status); + } + [Fact] public async Task DiversityProfileChangeAndBlockInFlightPreventPublication() + { + using var h = new Harness(); await SeedDiversity(h); h.Inventory.Service = "Usenet"; + h.Inventory.OnResolve = () => { h.Inventory.ProfileFingerprint = "changed"; return Task.CompletedTask; }; + await h.Run(); Assert.Equal(0,h.Inventory.Added); + h.Inventory.OnResolve = null; var c = await h.State(); + ImportDiversityPolicy.Reconcile(c,c.Items[0],h.Inventory.ProfileFingerprint,h.Now); c.Items[0].Diversity!.NextAttempt=h.Now; await h.Db.SaveImportCoverageAsync(c); + h.Inventory.OnResolve = async()=>{await h.Db.UpsertBlockedItemAsync(h.Item.AioId,null,null,"fixture","series","test");}; + await h.Run(); Assert.Equal(0,h.Inventory.Added); + } + [Fact] public async Task DiversityRateLimitKeepsWorkingSourceAndSeparateBackoff() + { + using var h = new Harness(); await SeedDiversity(h); h.Inventory.RateLimit=true; + await h.Run(); var e=(await h.State()).Items[0]; + Assert.Equal("http_429",e.Diversity!.Reason); Assert.True(e.Diversity.NextAttempt>h.Now); + Assert.Equal("",e.Failure); Assert.Single(e.Paths); Assert.True(h.Db.GetImportSourceHealth().Paused(h.Now)); + } + [Fact] public async Task DiversityEssentialAttemptWinsAndStaleHealthAllowsOnlyOneProbe() + { + using var h = new Harness(); await SeedDiversity(h,2); h.Inventory.Service="Usenet"; + var c=await h.State(); c.Items[1].Paths.Clear();c.Items[1].State="missing";c.Items[1].LastVersionRefresh=null; + h.Inventory.Files.Remove(2);await h.Db.SaveImportCoverageAsync(c); + await h.Run();Assert.Equal(1,h.Inventory.Published);Assert.Equal(0,h.Inventory.Added); + Assert.Equal(new[]{2},h.Inventory.ResolvedEpisodes); + } + private sealed class Harness : IDisposable { public TempDir Directory = new(); public DatabaseManager Db; public FakeInventory Inventory = new(); @@ -735,13 +803,15 @@ private sealed class Harness : IDisposable public PluginConfiguration Config = new() { ImportRecoveryMode = ImportMode.Repair }; public void Engage() { Config.ImportCatchUpStartedAt = Now.ToString("o"); Config.ImportCatchUpUntil = Now.AddDays(7).ToString("o"); } public Task Seed() => Db.UpsertCatalogItemAsync(Item); - public Task Run(ImportMode mode = ImportMode.Repair) => new ImportReconciliationService(Db, Inventory, () => mode, TimeZoneInfo.Utc, () => Now, () => Cooldown, () => ImportWorkBudget.For(Config, Now)).RunAsync(default, new[] { Item }); + public Task Run(ImportMode mode = ImportMode.Repair) => new ImportReconciliationService(Db, Inventory, () => mode, TimeZoneInfo.Utc, () => Now, () => Cooldown, () => ImportWorkBudget.For(Config, Now), () => Config.ImportProviderDiversityEnabled).RunAsync(default, new[] { Item }); public async Task State() => (await Db.GetImportCoverageAsync("series:imdb:tt999999991"))!; public void Dispose() => Directory.Dispose(); } private sealed class FakeInventory : IImportInventory { - public int Resolutions, Published, Notifications, FailEpisode, DisputedEpisode, Count = 2; + public int Resolutions, Published, Notifications, FailEpisode, DisputedEpisode, Count = 2, Added; + public string Service = ""; public string ProfileFingerprint { get; set; } = "fixture"; + public Dictionary Additions = new(); public bool HasSavedFiles(ImportEpisode episode) => episode.Episode.HasValue && Files.Contains(episode.Episode.Value); public List ResolvedEpisodes = new(); public bool Owned, BadSnapshot, Paused, RateLimit, Indexed = true; public bool ProviderPaused => Paused; public HashSet Files = new(); public Func? OnResolve; public Func? OnResolveWithToken; public Func? OnObserve; @@ -754,13 +824,21 @@ public Task FetchAsync(CatalogItem item, CancellationToken ct) public async Task ObserveAsync(CatalogItem item, ImportEpisode ep, CancellationToken ct) { if (OnObserve != null) await OnObserve(ep); - return new ImportObservation(Files.Contains(ep.Episode!.Value) ? new() { $"/fake/series/Season 01/e{ep.Episode}.strm" } : new(), Files.Contains(ep.Episode.Value) && Indexed ? new() { ep.Key } : new(), false); + var paths = Files.Contains(ep.Episode!.Value) ? new List { $"/fake/series/Season 01/e{ep.Episode}.strm" } : new(); + if (Additions.TryGetValue(ep.Episode.Value,out var addition)) paths.Add(addition); + return new ImportObservation(paths, Files.Contains(ep.Episode.Value) && Indexed ? new() { ep.Key } : new(), false); } public bool IsOwned(CatalogItem item) => Owned; public async Task> ResolveAsync(CatalogItem item, ImportEpisode ep, CancellationToken ct) - { Resolutions++; ResolvedEpisodes.Add(ep.Episode!.Value); if (RateLimit) throw new ImportHttpRateLimitException(); if (OnResolve != null) await OnResolve(); if (OnResolveWithToken != null) await OnResolveWithToken(ct); return ep.Episode == FailEpisode ? new() : new() { new() { Stream = new() { Url = "https://example.invalid/test" } } }; } + { Resolutions++; ResolvedEpisodes.Add(ep.Episode!.Value); if (RateLimit) throw new ImportHttpRateLimitException(); if (OnResolve != null) await OnResolve(); if (OnResolveWithToken != null) await OnResolveWithToken(ct); return ep.Episode == FailEpisode ? new() : new() { new() { Stream = new() { Url = "https://example.invalid/test", ServiceLabel=Service, SizeBytes=100 } } }; } public Task> PublishAsync(CatalogItem item, ImportEpisode ep, List versions, CancellationToken ct) { Published++; Files.Add(ep.Episode!.Value); return Task.FromResult(new List { $"/fake/series/Season 01/e{ep.Episode}.strm" }); } + public Task> PublishAdditionAsync(CatalogItem item,ImportEpisode ep,List versions,CancellationToken ct) + { + Added++; var path=$"/fake/series/Season 01/e{ep.Episode} - addition.strm";Additions[ep.Episode!.Value]=path; + ep.Versions.Add(new(path,"additional",versions[0].Stream.ServiceLabel,"1080p",null,100,Now)); + return Task.FromResult(ep.Paths.Append(path).ToList()); + } public void Notify(CatalogItem item) => Notifications++; } private sealed class TempDir : IDisposable diff --git a/UI/Settings/SyncAndMarvinTabView.cs b/UI/Settings/SyncAndMarvinTabView.cs index c44ca96..8561be5 100644 --- a/UI/Settings/SyncAndMarvinTabView.cs +++ b/UI/Settings/SyncAndMarvinTabView.cs @@ -49,6 +49,8 @@ public override async Task RunCommand(string itemId, string comma case "ImportCheck": result = await ImportHealthService.ApplyAsync(new() { Action = "check" }); break; case "ImportReprocess": result = await ImportHealthService.ApplyAsync(new() { Action = "reprocess_deferred" }); break; case "ImportCancelReprocess": result = await ImportHealthService.ApplyAsync(new() { Action = "cancel_reprocess" }); break; + case "ImportEnableDiversity": result = await ImportHealthService.ApplyAsync(new() { Action = "enable_diversity" }); break; + case "ImportDisableDiversity": result = await ImportHealthService.ApplyAsync(new() { Action = "disable_diversity" }); break; case "ImportIncludeSpecials": case "ImportExcludeSpecials": result = await ImportHealthService.ApplyAsync(new() { Action = cmd == "ImportIncludeSpecials" ? "include_specials" : "exclude_specials", Identity = itemId }); break; @@ -133,6 +135,9 @@ private async Task LoadImportsAsync() ? $"Paused until {When(health.PausedUntil, now)} · {Explain(health.LastOutcome)} · recovery step {health.Escalation}" : $"{health.ResponsiveReplies}/3 responsive replies · last match {When(health.LastMatchAt, now)} · does not test every indexer"); var recovery = await db.GetImportReprocessSummaryAsync(); + AddDashboard("Eventual source diversity", Plugin.Instance.Configuration.ImportProviderDiversityEnabled + ? "Enabled · after essential work · at most 2 normal / 8 catch-up checks per pass · one probe with stale health · see episode queue state" + : "Paused · durable entries and due times retained · essential repair continues"); AddDashboard("Recovery pass", recovery.RequestId == null ? "No recovery pass requested" : $"{recovery.Pending:N0} pending · {recovery.InFlight:N0} in flight · {recovery.Published:N0} published · {recovery.Failed:N0} failed · {recovery.Cancelled:N0} cancelled" + (!health.RecoveryReady(now) && recovery.Pending > 0 ? " · one probe per pass until three recent replies include a source match" : "")); @@ -198,12 +203,15 @@ private async Task LoadImportsAsync() Icon = IconNames.info, IconMode = ItemListIconMode.SmallRegular, Button1 = new ButtonItem(title.IncludeSpecials ? "Skip specials" : "Include released specials") { Data1 = title.Identity, CommandId = title.IncludeSpecials ? "ImportExcludeSpecials" : "ImportIncludeSpecials" } }); - foreach (var episode in title.Items.Where(x => x.State != "indexed" || x.Failure.Length > 0).Take(50)) + foreach (var episode in title.Items.Where(x => x.State != "indexed" || x.Failure.Length > 0 || x.Diversity?.Status is "waiting" or "in_flight" or "capacity").Take(50)) { var reason = episode.State == "indexed" && episode.Failure.Length > 0 ? "In Emby; source refresh failed" : episode.State == "excluded" ? Explain(episode.Eligibility) : Explain(episode.State); var more = episode.Failure.Length > 0 ? Explain(episode.Failure) : ""; if (episode.ConsecutiveFailures > 0) more += $" · failure step {episode.ConsecutiveFailures}"; if (episode.NextAttempt.HasValue) more += (more.Length > 0 ? " · " : "") + (episode.NextAttempt > now ? "Retry " + When(episode.NextAttempt, now) : "Retry due"); + if (episode.Diversity is { } diversity) + more += $" · source diversity {diversity.Status} ({diversity.Reason}) · {diversity.Attempts} reserved checks" + + (diversity.NextAttempt.HasValue ? $" · check after {When(diversity.NextAttempt, now)}" : ""); UI.ImportItems.Add(new GenericListItem { PrimaryText = (episode.Season.HasValue ? $"S{episode.Season:D2}E{episode.Episode:D2}" : "Movie") + " · " + reason, @@ -228,7 +236,7 @@ private void DrawRun(ImportRunSnapshot run, string label) AddDashboard(label, $"{run.Status} · {run.ElapsedSeconds:N0}s elapsed · started {run.StartedAt:HH:mm:ss} UTC"); if (run.Status == "running") AddDashboard("Doing", $"{run.Phase.Replace('_',' ')} · {run.PhaseSeconds:N0}s in this step" + (run.Title.Length > 0 ? $" · {run.Title} {run.EpisodeKey}" : "")); - AddDashboard("Lookups", $"{run.ActiveLookups:N0} in flight (includes paced waiting) · {run.MissingAttempts:N0} missing-file attempts · {run.RefreshAttempts:N0} refresh attempts"); + AddDashboard("Lookups", $"{run.ActiveLookups:N0} in flight (includes paced waiting) · {run.MissingAttempts:N0} missing-file attempts · {run.RefreshAttempts:N0} refresh attempts · {run.DiversityAttempts:N0} diversity attempts"); AddDashboard("AIO HTTP traffic", $"{run.HttpRequests:N0} submissions · {run.HttpRetries:N0} transport retries · excludes AIO's downstream indexer calls"); AddDashboard("Source results", $"{run.Matched:N0} matched · {run.EmptyResults:N0} empty · {run.CancelledLookups:N0} cancelled"); AddDashboard("Source failures", $"{run.TransportFailures:N0} transport · {run.LookupDeadlines:N0} lookup deadlines · {run.Http429:N0} HTTP 429 · {run.ProviderConfigurationFailures:N0} provider settings"); diff --git a/UI/Settings/SyncAndMarvinUI.cs b/UI/Settings/SyncAndMarvinUI.cs index 2788fe4..37395e1 100644 --- a/UI/Settings/SyncAndMarvinUI.cs +++ b/UI/Settings/SyncAndMarvinUI.cs @@ -50,6 +50,9 @@ public class SyncAndMarvinUI : EditableOptionsBase public LabelItem ReprocessHelp { get; set; } = new LabelItem("After an outage, queue one retry of the deferred items. Marvin waits through provider cooldowns, tests a small sample, then drains the queue as AIO responds. Existing files and retry history stay intact."); public ButtonItem ReprocessDeferred { get; set; } = new ButtonItem("Reprocess deferred items") { Data1 = "ImportReprocess" }; public ButtonItem CancelReprocess { get; set; } = new ButtonItem("Cancel pending reprocessing") { Data1 = "ImportCancelReprocess" }; + public LabelItem DiversityHelp { get; set; } = new LabelItem("Single-source items wait in a durable, low-priority queue. Marvin adds another validated source when credits and health allow; working files stay intact. Eight-version items wait for room. This does not guarantee playback or a second available source."); + public ButtonItem EnableDiversity { get; set; } = new ButtonItem("Enable eventual source diversity") { Data1 = "ImportEnableDiversity" }; + public ButtonItem DisableDiversity { get; set; } = new ButtonItem("Pause source diversity") { Data1 = "ImportDisableDiversity" }; public GenericItemList ImportItems { get; set; } = new GenericItemList(); public ButtonItem NextImports { get; set; } = new ButtonItem("Next 25 titles") { Data1 = "ImportNext" }; } diff --git a/docs/README.md b/docs/README.md index bfadee8..58600d7 100644 --- a/docs/README.md +++ b/docs/README.md @@ -7,6 +7,7 @@ - [Discover](USER_DISCOVER_UI.md): browser discovery and saved titles. - [External lists](EXTERNAL_LISTS.md): system and user catalog sources. - [Title blocks](title-blocks.md): administrator whole-title blocking and recoverable stream cleanup. +- [Eventual source diversity](provider-diversity-queue.md): durable low-priority queue, exact backoff, additive publication and operator controls. - [Import recovery](import-reconciliation.md): Observe/Repair, limits and retention. - [Troubleshooting](troubleshooting.md): diagnosis and recovery. - [Security](SECURITY.md): private configuration, STRMs and administrative access. @@ -27,4 +28,4 @@ old test plans and dated reports. These are not current configuration instructions. The implementation and current guides take precedence over historical claims. -Updated October 1, 2026 for the 0.42.13 candidate source tree. +Updated October 3, 2026 for the 0.42.14 source tree. diff --git a/docs/import-reconciliation.md b/docs/import-reconciliation.md index c0b725f..f478998 100644 --- a/docs/import-reconciliation.md +++ b/docs/import-reconciliation.md @@ -350,3 +350,7 @@ Prowlarr search submission globally. InfiniteDrive still validates and preselect STRMs ahead of playback. No reserved Usenet/TorBox slot or provider policy change is included. Back up the current plugin DB/configuration before deploying; retain additive recovery rows when rolling back the DLL, and never restore older DB state. + +## Eventual provider diversity (0.42.14) + +Successful single-source choices have a separate durable low-priority queue. It adds a validated source without replacing working files, shares native credits/health/backoff, and persists beyond catch-up. This is distinct from essential failure reprocessing. See the [exact provider-diversity contract](provider-diversity-queue.md), including capacity/unknown evidence, accounting, every state transition and rollback. diff --git a/docs/provider-diversity-queue.md b/docs/provider-diversity-queue.md new file mode 100644 index 0000000..d62466b --- /dev/null +++ b/docs/provider-diversity-queue.md @@ -0,0 +1,107 @@ +# Eventual source diversity — Marvin's durable queue + +Implemented for 0.42.14, targeting Emby 4.10.0.40. Deployment and dated QA evidence are recorded separately in the release notes and stack handoff. This queue improves already indexed items with one known delivery source. It does not guarantee a second source exists, represent a download, or establish playback. + +## Enable, pause and ownership + +`ImportProviderDiversityEnabled` defaults to **true**, including when an older configuration has no such XML member. The native Marvin page has **Enable eventual source diversity** and **Pause source diversity**. Both require Emby's normal settings administration. The administrator-only POST `/InfiniteDrive/Imports/Action` accepts `enable_diversity` or `disable_diversity`; GET never starts work. Neither action triggers Marvin, changes Repair/catch-up, clears backoff, deletes entries or expands source policy. + +Diversity dispatch/publication also requires Repair mode. Off runs no coordinator. Observe can discover/reconcile durable queue entries but makes no diversity requests/publications; it is not a global read-only switch for the legacy importer. Pausing diversity retains all queue state and lets essential repair continue. Disabling during a lookup prevents its additive publication at the final guard, clears its lease and leaves a waiting entry with a five-minute guard-change deferral. It cannot undo an addition already published. Re-enable waits for normal scheduled passes and existing due times/guards. There is no force-reprocess-diversity or mass-reset action; the queue drains automatically. **Reprocess deferred items** is the distinct essential-failure recovery queue and does not bypass diversity due times. + +## Unit, labels and eligibility + +One entry belongs to `(import_coverage.Identity, ImportEpisode.Key)`: `movie` or confirmed `aired:S:E`. Several resolutions/releases from TorBox still count as one source. A series with ten TorBox-only and ten Usenet-only episodes has twenty single-source episodes, regardless of the series' aggregate labels. + +Only current `Paths` and matching `Versions.Path` supply provider evidence. Every retained path must have exactly one unambiguous recognized label; missing/conflicting/unrecognized evidence is unknown. Retired version paths do not count. Recognized labels are TorBox, Usenet, Real-Debrid, AllDebrid, Premiumize, Debrid-Link, Offcloud, EasyDebrid, Debrider, Easynews, NZBDav, AltMount and StremThru; comparison ignores case. `nntp` and `stremionntp` normalize to Usenet. The Usenet label is a delivery lane, not a particular indexer; two labels do not prove independent infrastructure. + +Eligibility requires a successful inventory snapshot, no coverage exclusion, Expected and Eligible true, indexed state, no essential Failure, and exactly one known source. Authoritative alias blocks/removals, incompatible identities and physical owned media are rechecked before dispatch and publication. No-source items remain essential missing-file repair. Future/special/disputed/ineligible keys stay excluded. A failed essential refresh retains its own failure/backoff and takes the item out of active diversity. Deliberately archived blocked-title references are not diversity gaps. + +Before an outbound diversity lookup, production also verifies the retained file set, managed-root/symlink protections and each retained URL hash. A changed set/hash postpones this entry six hours as `retained_file_changed` without charging a request. Before writing, the same protections are checked again. Unknown **sizes** on working files do not authorize deletion and do not by themselves prevent an additive improvement; unknown **provider identity** prevents single-source classification. + +## Durable schema, migration and generations + +`ImportEpisode.Diversity` is an additive JSON object inside the existing `import_coverage.payload`; no new SQLite table or destructive migration is needed. Older rows deserialize with null Diversity. The episode's unique location makes census idempotent: repeated observations update one object, never append jobs. SQLite coverage saves are atomic document upserts under the existing `MutationGate`; only one native run holds `RunGate`. + +| Field | Meaning | +| --- | --- | +| Status | `waiting`, `in_flight`, `capacity`, `resolved` or `retired` | +| Profile | SHA-256 endpoint fingerprint from the configured AIO client; no endpoint/token is copied into the entry | +| Generation | Coverage's authorization/specials generation at creation | +| Reason | Last eligibility, deferral, verification or completion reason | +| Streak | Consecutive diversity deferrals, separate from essential failure streak | +| Attempts | Reserved diversity checks, not indexer HTTP calls | +| CreatedAt | First entry creation for this profile/generation | +| CheckedAt | Last reservation or completed/deferred check | +| NextAttempt | Earliest allowed check; not a promised execution time | + +Creation, endpoint/profile fingerprint change or coverage generation change creates one new entry with a **six-hour** initial wait and zero diversity streak/attempts. It does not reset the essential failure state, existing attempt ledger or provider cooldown. The fingerprint detects configured endpoint changes, not opaque AIO server-side policy edits at the same endpoint; those do not automatically shorten due times. + +Dispatch claims the existing episode lease with a random ID and **ten-minute** expiry, persists entry `in_flight`, increments its Attempts and reserves the shared attempt credit before starting transport. Publication rechecks live generation, lease, enabled mode, authorization/ownership, current profile and a signature of retained paths/hashes/labels. A changed guard prevents addition. Every completed path clears its lease; budget cancellation leaves queue work waiting five minutes without penalizing the repair item. After a crash, a missing/expired lease on `in_flight` becomes waiting with a five-minute interruption delay when reconciled. Active leases cannot be stolen. + +Markers `import_diversity_schema=1` and `import_diversity_enabled` are written by native passes for read-only reporting. Their timestamps/snapshots are not instantaneous UI control acknowledgements. Reading administrator coverage/status or the public collector never creates entries. + +## Priority, pages, deadlines and request accounting + +The queue uses **the existing native schedule**, not a new timer. All discovered missing, ordinary refresh and essential reprocess candidates are handled before diversity. Existing missing/refresh priority rotation remains unchanged. Diversity gets at most **two reservations per normal pass** or **eight per catch-up pass**, processes them serially and can only use capacity remaining in that pass's shared budget. Normal: 120-second pass, 20 total attempts, 5 ordinary upgrades, 200 rolling daily attempts. Catch-up: 480 seconds, 4,096 total attempts/upgrades, 40,000 rolling daily safeguard. Those are ceilings, not measured provider allowances. Diversity does not enlarge them or alter AIO/indexer concurrency. + +The inventory cursor discovers entries in existing catalog scans; every pass also loads up to ten due diversity titles ordered by coverage check time and catalog ID. Episode candidates are ordered by NextAttempt, then identity/key. Existing observations check up to 200 keys per title. Discovery is incremental, so the dashboard need not have a queue entry for every single-source reference immediately after upgrade. An ongoing urgent backlog can starve diversity; there is no independent minimum quota or completion deadline. Unresolved entries remain beyond catch-up expiry under normal limits. + +Diversity uses one maintenance AIO submission with the coordinator's lookup deadline (60 seconds normal, 120 catch-up), no transport retry. It selects a validated unfamiliar source from the response **before** quality-bucket selection can hide that source. AIO's configured providers/filters remain authoritative. It does not make separate per-provider/indexer searches. AIO fan-out can still make multiple downstream requests, so one logical reservation is not one indexer call. + +The shared `import_attempts` reservation is persisted **before** transport and remains even if cancellation happens before HTTP submission. Counters therefore distinguish reserved attempts, actual `HttpRequests`/`HttpRetries`, matching/empty outcomes, saved variants and group publications. Diversity attempts are reported separately from MissingAttempts and RefreshAttempts; DiversityPublished counts additive movie/episode-group events and is included in total Published. It does not increment Refreshed because old choices were not replaced. Same-source matches can increase Matched without resolving diversity. Old reports lack these new counters and read as zero. + +## Exact health and backoff rules + +Every reservation requires provider maintenance not paused, global cooldown elapsed, the native source circuit not paused, per-item essential backoff/lease respected, available daily/pass credit and current authorization. For full diversity capacity, `ImportSourceHealth.RecoveryReady(now)` requires no pause, ResponsiveReplies >=3, MatchedReplies >0, LastResponseAt and LastMatchAt within 15 minutes, and ConsecutiveErrors ==0. These are AIO-path observations, not tests of every indexer. The existing responsive counter resets on a response gap longer than 15 minutes; it is not a separate timestamped rolling list of three responses. + +If that readiness is stale/unproven, **at most one total reservation in the entire pass** may be a diversity probe. If essential work already reserved a check, diversity cannot add a stale-health probe. A paused/error circuit is not bypassed. Successful same-source replies are responsive source matches; empty usable selections are responsive empties. Transport/deadline/429/configuration retain the existing circuit semantics. No new all-provider health polling is introduced. + +A diversity deferral increments its own Streak, regardless of prior reason. `same_source`, `source_unavailable` and `no_valid_addition` use this ladder: **6h, 24h, 48h, 72h, then 7d**. Transport, lookup deadline, HTTP429, provider configuration and publication failure use **15m, 1h, 4h, 1d, 3d, then 7d** at the current streak index. Delays multiply by random jitter from **0 through 20%** and then take the maximum with global cooldown, including longer provider Retry-After handling. The capped base is seven days; jitter can make the resulting wait 8.4 days. A different error type does not restart the diversity streak. Confirmed additional-source publication clears it. A new endpoint/generation creates a new entry as described above. + +Slice cancellation, interrupted lease and changed publication guard wait five minutes without incrementing Streak. Changed retained file evidence waits six hours without a request or streak increment. Capacity is checked without making a request. Actual transport/empty/429/publication failures are stored **on the diversity entry**, without marking a working indexed item's essential repair failed, clearing its original Failure/NextAttempt, or moving LastVersionRefresh. + +## State transitions + +| Event / guard | Durable outcome | Reserved credit | Next check | +| --- | --- | --- | --- | +| New eligible single-source item | waiting / single_source | none | creation +6h | +| Same identity/profile/generation census | retains entry/streak/due time | none | unchanged | +| Eight current variants | capacity / version_limit | none | no diversity dispatch until room exists | +| Capacity becomes <8 | waiting / single_source | none | retained due time | +| Due + enabled Repair + guards + spare budget | in_flight; lease; Attempts+1 | one shared reservation | lookup deadline applies | +| Same-source or no usable additional selection | waiting; distinct reason; Streak+1 | retains reservation | no-source ladder + jitter/cooldown | +| Transport/deadline/429/config/publication failure | waiting; distinct reason; Streak+1 | retains reservation | error ladder + jitter/cooldown | +| Valid second source saved/evidence verified | resolved / multiple_sources; streak0 | retains reservation | null | +| Source/authorization/generation changes in flight | no addition; guard_changed if still live | retains reservation | +5m, then reconciliation | +| Block/ownership/ineligible snapshot | retired / ineligible on observation; dispatch/publication forbidden immediately | none beyond prior reservation | no dispatch | +| Multiple saved sources found by census | resolved / multiple_sources | none | no dispatch | +| Retired/resolved item later eligible with one source | waiting / single_source; no force reset of old streak | none | now +6h | +| Process/slice interruption | waiting / interrupted or slice_cancelled | prior reservation retained | +5m after reconciliation/cancellation | +| Disable, Observe, health pause, exhausted budget, future due | entry retained; no new dispatch/addition | none | guard must clear; original due remains | + +## Additive publication, capacity and recovery + +Only a new recognized label from the configured AIO response is accepted. The candidate must pass existing CAM/REMUX policy, have an HTTP(S) URL and **positive exact upstream SizeBytes <=40,000,000,000**. Rounded size text is insufficient. Unknown labels/sizes and ambiguous identities are never forced. + +The queue adds **one** variant per resolved item. The real writer shares the folder lock with ordinary replacement and uses a deterministic URL-hash-suffixed Emby variant name. It verifies an existing identical target for idempotence or moves a temporary file with overwrite **false**. It never deletes/rewrites an existing STRM, including unknown-size originals or adjacent episodes. Existing hash/size evidence is retained; the new file gets its own URL-bound evidence. Catalog stored versions are merged instead of replaced. No metadata/source/media ownership is reassigned. + +The existing eight-version limit remains. A single-source item already using eight slots is surfaced as capacity, not trimmed. Ordinary authorized refresh may later free a slot; diversity never does so itself. Normal refresh still has its original replacement semantics and can subsequently change source coverage, causing a new six-hour wait. Diversity does not guarantee permanent two-provider coverage against later source changes. + +Filesystem publication and SQLite are not one atomic transaction. After a crash between writing an addition and checkpointing evidence, preserve the file. A census with an unbound/unlabelled added path classifies coverage unknown, retires unsupported diversity work and requires evidence review/ordinary safe refresh; it does not invent metadata, delete the orphan or publish duplicates. A completed lookup/pass is not completion: only verified additional saved evidence resolves the entry. A new alternate need not yet be indexed even when the original item remains indexed. + +## Inspection, examples and rollback + +Native Marvin displays the enabled/paused setting, diversity attempts and active/capacity episode entries (up to 50 displayed per title). Administrator GET `/InfiniteDrive/Imports` returns paged episodes, recognized provider labels and Diversity fields. The public Failure Library reads a sanitized saved projection: one-source filter/counts per item, actual queue states, safe reason, reserved checks, streak and next-check time. Missing entries are labelled **no native entry**, never fabricated queued. Failed/older-than-five-minute snapshots are historical. No provider URLs, file paths, credentials, profile fingerprints or authentication storage are exported to the public collector. + +Examples: + +- Three TorBox variants: one queue entry, initial wait six hours. A TorBox-only response defers at least six hours, leaves all three files and never starts an immediate loop. +- TorBox then validated Usenet: one additive STRM, old hashes/paths unchanged, entry resolved. This does not prove playback or a particular Usenet indexer worked. +- HTTP429: retain TorBox and store the actual reason; apply the current error rung plus jitter and any longer global cooldown. Native circuit also closes; no queue-wide reset. +- Outage: no calls while health/cooldown pauses. Afterwards one bounded shared probe while readiness is unproven; full 2/8 cap only with readiness and spare credits. +- Restart: waiting state/due/streak persist. An expired in-flight lease waits five more minutes and revalidates before another reservation. +- Block during lookup: final alias/ownership/generation guards prevent publication. A later observation retires the entry; unblocking is never automatic. +- Endpoint change: a new six-hour entry for the new fingerprint; an old result cannot publish. Server-side edits at the same endpoint do not reset waiting times. +- Catch-up expires with 1,000 entries: they remain and contend only for normal remaining capacity. No window extension, quota increase or completion promise. + +Before upgrade, take private current configuration and stopped-database backups. Prefer pausing diversity through the native control for behavior rollback; this preserves the queue and newer state. A binary downgrade never requires restoring an older database or deleting ledger/coverage/archives. Older binaries can ignore/drop unknown Diversity JSON members when rewriting coverage; preserve the current backup for audit and expect conservative rediscovery after re-upgrade, rather than restoring historical user state. Protect active playback/recordings during any approved Emby restart. The beta/isolated QA plugin must never be reused as a production artifact. diff --git a/plugin.json b/plugin.json index bf1e74d..bca739c 100755 --- a/plugin.json +++ b/plugin.json @@ -5,7 +5,7 @@ "overview": "InfiniteDrive discovers your AIOStreams catalog, writes .strm files, and resolves debrid URLs on demand. Like the Infinite Improbability Drive: a stream will appear. Probably. No external processes needed.", "owner": "InfiniteDrive", "category": "General", - "version": "0.42.13.0", + "version": "0.42.14.0", "targetAbi": "4.10.0.40", "framework": "net8.0", "imageUrl": "" diff --git a/tools/ImportQa/ImportLabService.cs b/tools/ImportQa/ImportLabService.cs index 2a5c51f..e8ae8f6 100644 --- a/tools/ImportQa/ImportLabService.cs +++ b/tools/ImportQa/ImportLabService.cs @@ -39,6 +39,7 @@ public async Task Post(ImportQaRequest request) var fake = new LabInventory(real); var now = DateTimeOffset.UtcNow; var mode = ImportMode.Repair; + ImportWorkBudget? diversityBudget = null; switch (request.Step) { case "prepare": @@ -110,10 +111,19 @@ public async Task Post(ImportQaRequest request) var view = new InfiniteDrive.UI.Settings.SyncAndMarvinTabView(p.Id.ToString(), new()); await view.RunCommand("", "ImportRefresh", ""); return new { Dashboard = view.ContentData, CurrentRun = ImportRunTelemetry.Current }; + case "diversity-same": + case "diversity-add": + now = now.AddDays(request.Step == "diversity-same" ? 1 : 3); + fake.AddProvider = request.Step == "diversity-add"; + diversityBudget = new(120,20,5,5,200,now.AddDays(-6),now.AddDays(1)); + var ready = new ImportSourceHealth(); + for (var i=0;i<3;i++) ready.Record("",now); + await db.SaveImportSourceHealthAsync(ready, CancellationToken.None); + break; case "status": return await Status(db, fake); default: throw new ArgumentException("Unknown QA step"); } - await new ImportReconciliationService(db, fake, () => mode, TimeZoneInfo.Utc, () => now) + await new ImportReconciliationService(db, fake, () => mode, TimeZoneInfo.Utc, () => now, workBudget: () => diversityBudget ?? ImportWorkBudget.Normal) .RunAsync(CancellationToken.None, new[] { item }); return await Status(db, fake); } @@ -126,7 +136,8 @@ private static async Task Status(DatabaseManager db, LabInventory fake) private sealed class LabInventory : IImportInventory { private readonly ImportInventory _real; - public bool FailSecond, PartialNumbering; public int Resolutions, Published; + public bool FailSecond, PartialNumbering, AddProvider; + public string ProfileFingerprint => "offline-native-fixture"; public int Resolutions, Published; public List Episodes => Enumerable.Range(1, 2).Select(n => new ImportEpisode { Key = $"aired:1:{n}", Season = 1, Episode = n, Released = DateTimeOffset.Parse("2020-01-01T00:00:00Z") }).ToList(); public LabInventory(ImportInventory real) { _real = real; } @@ -142,11 +153,19 @@ public Task> ResolveAsync(CatalogItem item, ImportEpisode { Resolutions++; return Task.FromResult(FailSecond && ep.Episode == 2 ? new List() : new List - { new() { Stream = new() { Url = "https://example.invalid/qa/" + ep.Episode } }, - new() { Stream = new() { Url = "https://example.invalid/qa/" + ep.Episode + "/hd" }, VersionLabel = "1080p" } }); + { new() { Stream = new() { Url = "https://example.invalid/qa/" + ep.Episode, ServiceLabel="TorBox", SizeBytes=100 } }, + new() { Stream = new() { Url = "https://example.invalid/qa/" + ep.Episode + "/hd", ServiceLabel="TorBox", SizeBytes=200 }, VersionLabel = "1080p" } }); } public async Task> PublishAsync(CatalogItem item, ImportEpisode ep, List versions, CancellationToken ct) { Published++; return await _real.PublishAsync(item, ep, versions, ct); } + public Task> ResolveDiversityAsync(CatalogItem item, ImportEpisode ep, CancellationToken ct) + { + Resolutions++; + return Task.FromResult(new List { new() { VersionLabel=AddProvider ? "1080p - Usenet" : "1080p - TorBox", + Stream=new() { Url="https://example.invalid/diversity/"+ep.Episode, ServiceLabel=AddProvider ? "Usenet" : "TorBox", SizeBytes=300 } } }); + } + public async Task> PublishAdditionAsync(CatalogItem item, ImportEpisode ep,List versions,CancellationToken ct) + { Published++;return await _real.PublishAdditionAsync(item,ep,versions,ct); } public void Notify(CatalogItem item) => _real.Notify(item); } } diff --git a/tools/ImportQa/README.md b/tools/ImportQa/README.md index 8de8733..8fd96f4 100644 --- a/tools/ImportQa/README.md +++ b/tools/ImportQa/README.md @@ -45,3 +45,20 @@ episode 2 stays excluded. `discover` observes native indexing. `recheck-numbering` advances the fixture clock eight hours and restores the complete source inventory; episode 2 becomes eligible without a database reset. These are synthetic targets and do not establish playback. + +Provider-diversity regression (fresh fixture): `prepare`, register the offline TV +library as above, `retry`, then `discover` establish two indexed episodes with +two TorBox variants each. `diversity-same` advances only the injected fixture +clock and supplies synthetic healthy catch-up observations: two same-source +checks defer without changing any file. Restart only the isolated lab and +compare durable Diversity fields. `diversity-add` advances the fixture clock +again and supplies validated synthetic Usenet variants: one addition per episode, +all original hashes preserved. `discover` verifies native indexing. Inspect all +six retained URL hashes and positive size evidence. These steps use no real +provider requests; they do not test playback. + +Verified October 3, 2026 against isolated Emby 4.10.0.40 / candidate 0.42.14.0: +normal QA administrator controls passed, anonymous mutation was denied, census +created two entries, same-source results preserved files, state survived the lab +restart, two additive publications passed six-file hash/size-evidence checks and +both episodes indexed. The separate production artifact excludes this helper.