diff --git a/README.md b/README.md index 27c6d2d..97b8639 100755 --- a/README.md +++ b/README.md @@ -31,8 +31,9 @@ connections. It has its own catalogs, filters, sorting and labels. Other clients use separate profiles. AIOMetadata supplies Popular, Top Rated and Trending feeds for movies and TV. -We cap each of those six feeds at 4,000 entries. They overlap; this is neither -24,000 unique titles nor a complete TMDB import. Feed order depends on the +We cap each movie feed at 200 titles and each TV feed at 50 series. They overlap; +750 memberships is the combined ceiling, not a unique-title count or a complete +TMDB import. Feed order depends on the provider. We skip search and calendar feeds during bulk import. InfiniteDrive handles our MDBList subscriptions directly. We disable their @@ -201,8 +202,10 @@ Our Marvin task runs every ten minutes; a large library takes many passes. For a big refresh, engage **Infinite Improbability Drive** on the Marvin page. The switch runs up to 64 lookups together, paced to two starts a second. Each -pass gets eight minutes, up to 4,096 attempts and a rolling daily ceiling of -40,000. Existing STRMs get priority. It remembers refreshed episodes, +pass gets eight minutes, up to 4,096 attempts shared by missing-file repairs and +existing-file refreshes, and a rolling daily ceiling of 40,000. Existing managed +titles get priority. Completed lookups publish while inventory checks continue. +It remembers refreshed episodes, keeps provider backoff and media protections, and switches itself off after seven days. Disengage it any time: **Normality has been restored.** Actual speed depends on your sources; the ceiling is not a completion estimate. diff --git a/Services/ImportReconciliationService.cs b/Services/ImportReconciliationService.cs index 1f537bb..7e0ad50 100644 --- a/Services/ImportReconciliationService.cs +++ b/Services/ImportReconciliationService.cs @@ -57,7 +57,6 @@ public async Task RunAsync(CancellationToken ct, IReadOnlyList? sel var runId = Guid.NewGuid().ToString("N"); var status = "success"; var pending = new List(); - var gapAttempts = 0; async Task ResolveAsync(CatalogItem item, ImportEpisode episode) { @@ -134,6 +133,17 @@ async Task DrainOneAsync() pending.Remove(work); await CompleteAsync(work); } + + async Task DrainCompletedAsync() + { + // Publish ready results before more inventory work. Waiting for 64 + // queued requests (or the end of the page) strands small batches. + while (pending.FirstOrDefault(x => x.Resolution.IsCompleted) is { } work) + { + pending.Remove(work); + await CompleteAsync(work); + } + } try { await _db.EnsureImportCoverageAsync(token); @@ -166,6 +176,7 @@ async Task DrainOneAsync() var seen = new HashSet(StringComparer.Ordinal); foreach (var sourceItem in page) { + await DrainCompletedAsync(); var item = sourceItem; token.ThrowIfCancellationRequested(); if (_mode() == ImportMode.Off) break; @@ -240,6 +251,7 @@ async Task DrainOneAsync() if (episodes.Count == 0) { coverage.Cursor = ""; episodes = coverage.Items.OrderBy(x => x.Key, StringComparer.Ordinal).Take(200).ToList(); } foreach (var episode in episodes) { + await DrainCompletedAsync(); token.ThrowIfCancellationRequested(); now = _clock(); ImportObservation observation; @@ -283,9 +295,9 @@ async Task DrainOneAsync() // Adopting an existing file is not a successful refresh against the current profile. if ((episode.State != "missing" && !upgrade) || coverage.SnapshotStatus != "success" || _inventory.ProviderPaused || attempts >= allowance.AttemptsPerSlice || - (allowance.IsCatchUp && !upgrade && gapAttempts >= ImportWorkBudget.Normal.AttemptsPerSlice) || _cooldown() > now || await _db.GetRecentImportAttemptsAsync(now) >= allowance.AttemptsPerDay) continue; + _cooldown() > now || await _db.GetRecentImportAttemptsAsync(now) >= allowance.AttemptsPerDay) continue; attempts++; - if (upgrade) upgrades++; else gapAttempts++; + if (upgrade) upgrades++; var lease = Guid.NewGuid().ToString("N"); episode.Lease = lease; episode.LeaseUntil = now.AddMinutes(10); if (!upgrade) episode.InitialFailure = true; episode.Attempts++; @@ -293,6 +305,7 @@ async Task DrainOneAsync() await SaveObservedAsync(coverage, token); pending.Add(new(coverage, item, episode, coverage.Generation, lease, upgrade, episode.State, ResolveAsync(item, episode))); + await DrainCompletedAsync(); if (pending.Count >= allowance.Parallelism) await DrainOneAsync(); } scanned++; diff --git a/Tests/ImportReconciliationTests.cs b/Tests/ImportReconciliationTests.cs index c5c014a..cfdf2e1 100644 --- a/Tests/ImportReconciliationTests.cs +++ b/Tests/ImportReconciliationTests.cs @@ -348,10 +348,12 @@ [Fact] public async Task ImprobabilityHonorsProviderPauseAndOwnership() Assert.Equal(0, h.Inventory.Published); } - [Fact] public async Task CatchUpLookupsOverlapButPublicationIsSerialAndCheckpointsAllSiblings() + [Theory] [InlineData(false)] [InlineData(true)] + public async Task CatchUpLookupsOverlapButPublicationIsSerialAndCheckpointsAllSiblings(bool existingFiles) { using var h = new Harness(); h.Inventory.Count = 80; - h.Inventory.Files.UnionWith(Enumerable.Range(1, 80)); await h.Seed(); h.Engage(); + if (existingFiles) h.Inventory.Files.UnionWith(Enumerable.Range(1, 80)); + await h.Seed(); h.Engage(); var full = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); var release = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); var active = 0; var maximum = 0; @@ -391,10 +393,43 @@ [Fact] public async Task ExistingFilesRefreshWhileNativeIndexingIsPending() Assert.All((await h.State()).Items, e => Assert.Equal("awaiting_indexing", e.State)); await h.Run(); Assert.Equal(2, h.Inventory.Published); } - [Fact] public async Task CatchUpDoesNotSpendTheBacklogBudgetCreatingThousandsOfMissingEpisodes() + [Fact] public async Task CatchUpRepairsMissingEpisodesBeyondTheNormalTwentyAttemptLimit() + { + using var h = new Harness(); h.Inventory.Count = 80; await h.Seed(); h.Engage(); + await h.Run(); Assert.Equal(80, h.Inventory.Resolutions); + Assert.Equal(80, h.Inventory.Published); + await h.Run(); Assert.Equal(80, h.Inventory.Resolutions); + } + [Fact] public async Task CatchUpMissingFillsStillShareTheRunAttemptCeiling() { using var h = new Harness(); h.Inventory.Count = 80; await h.Seed(); h.Engage(); - await h.Run(); Assert.Equal(20, h.Inventory.Resolutions); + var limit = ImportWorkBudget.For(h.Config, h.Now) with { AttemptsPerSlice = 37 }; + var worker = new ImportReconciliationService(h.Db, h.Inventory, () => ImportMode.Repair, + TimeZoneInfo.Utc, () => h.Now, workBudget: () => limit); + await worker.RunAsync(default, new[] { h.Item }); + Assert.Equal(37, h.Inventory.Resolutions); Assert.Equal(37, h.Inventory.Published); + await worker.RunAsync(default, new[] { h.Item }); + Assert.Equal(74, h.Inventory.Resolutions); Assert.Equal(74, h.Inventory.Published); + } + [Fact] public async Task ReadyResultsPublishBeforeTheNextInventoryObservationFinishes() + { + using var h = new Harness(); await h.Seed(); h.Engage(); + var laterObservation = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var release = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + h.Inventory.OnObserve = async ep => + { + if (ep.Episode != 2) return; + laterObservation.TrySetResult(true); + await release.Task; + }; + var run = h.Run(); + try + { + await laterObservation.Task.WaitAsync(TimeSpan.FromSeconds(10)); + Assert.Equal(1, h.Inventory.Published); + } + finally { release.TrySetResult(true); await run; } + Assert.Equal(2, h.Inventory.Published); } [Fact] public async Task CancellationJoinsAllLookupsBeforeReturningAndKeepsOldFiles() { @@ -482,9 +517,13 @@ private sealed class Harness : IDisposable private sealed class FakeInventory : IImportInventory { public int Resolutions, Published, Notifications, FailEpisode, Count = 2; - public bool Owned, BadSnapshot, Paused, Indexed = true; public bool ProviderPaused => Paused; public HashSet Files = new(); public Func? OnResolve; public Func? OnResolveWithToken; + public bool Owned, BadSnapshot, Paused, Indexed = true; public bool ProviderPaused => Paused; public HashSet Files = new(); public Func? OnResolve; public Func? OnResolveWithToken; public Func? OnObserve; public Task FetchAsync(CatalogItem item, CancellationToken ct) => Task.FromResult(new ImportSnapshot(BadSnapshot ? new() : Enumerable.Range(1, Count).Select(Ep).ToList(), "success")); - public Task ObserveAsync(CatalogItem item, ImportEpisode ep, CancellationToken ct) => Task.FromResult(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)); + 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); + } public bool IsOwned(CatalogItem item) => Owned; public async Task> ResolveAsync(CatalogItem item, ImportEpisode ep, CancellationToken ct) { Resolutions++; 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" } } }; } diff --git a/docs/import-reconciliation.md b/docs/import-reconciliation.md index 1d733ad..eaabc53 100644 --- a/docs/import-reconciliation.md +++ b/docs/import-reconciliation.md @@ -99,21 +99,24 @@ again does not extend its window or reset its progress. | Stream attempts per rolling 24 hours | 200 | 40,000 | | Concurrent recovery lookups | 1 | 64 | | Starts per second in maintenance lane | — | 2 | -| Missing-file fills per run | 20 | 20 | +| Missing-file fills per run | 20 | Up to 4,096, sharing the total attempt budget | These are ceilings, not promised throughput. Upstream latency, unavailable sources, provider backoff, native indexing and the rest of Marvin's work determine actual progress. The existing task schedule stays in place. Catch-up has a separate 64-slot AIO -maintenance lane, paced to two starts a second. Ordinary lookups retain their +maintenance lane, paced to two starts a second. Missing-file repairs and existing +version refreshes share its larger attempt allowance. Ordinary lookups retain their shared two-slot lane. File publication and state changes remain serial. No extra -scheduler is started. Queued lookups are joined before a cancelled run releases +scheduler is started. Completed lookups publish between inventory checks, rather +than waiting for a full request batch or the end of the catalog page. Queued lookups are joined before a cancelled run releases its lock; late results cannot publish after disengagement or a generation change. A successfully refreshed movie/episode is skipped for the remainder of that drive window. Its checkpoint survives a restart. Unfinished existing versions are revisited alongside bounded missing-file repairs. Catch-up scans catalog rows with -existing managed paths first, using its own window/cursor; it does not spend the -backlog budget creating thousands of newly catalogued episodes. The normal catalog +existing managed paths first, using its own window/cursor; it also repairs eligible +gaps for these titles and due authorized work. It does not sweep the entire +unmaterialized catalog. The normal catalog sweep resumes when catch-up ends. The cursor also revisits episode inventories larger than a 200-key page. An unavailable episode does not prevent refreshing its eligible siblings. Valid existing managed