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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 7 additions & 4 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down
19 changes: 16 additions & 3 deletions Services/ImportReconciliationService.cs
Original file line number Diff line number Diff line change
Expand Up @@ -57,7 +57,6 @@ public async Task RunAsync(CancellationToken ct, IReadOnlyList<CatalogItem>? sel
var runId = Guid.NewGuid().ToString("N");
var status = "success";
var pending = new List<ResolutionWork>();
var gapAttempts = 0;

async Task<ResolutionResult> ResolveAsync(CatalogItem item, ImportEpisode episode)
{
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -166,6 +176,7 @@ async Task DrainOneAsync()
var seen = new HashSet<string>(StringComparer.Ordinal);
foreach (var sourceItem in page)
{
await DrainCompletedAsync();
var item = sourceItem;
token.ThrowIfCancellationRequested();
if (_mode() == ImportMode.Off) break;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -283,16 +295,17 @@ 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++;
await _db.RecordImportAttemptAsync(lease, now, token);
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++;
Expand Down
51 changes: 45 additions & 6 deletions Tests/ImportReconciliationTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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<bool>(TaskCreationOptions.RunContinuationsAsynchronously);
var release = new TaskCompletionSource<bool>(TaskCreationOptions.RunContinuationsAsynchronously);
var active = 0; var maximum = 0;
Expand Down Expand Up @@ -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<bool>(TaskCreationOptions.RunContinuationsAsynchronously);
var release = new TaskCompletionSource<bool>(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()
{
Expand Down Expand Up @@ -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<int> Files = new(); public Func<Task>? OnResolve; public Func<CancellationToken, Task>? OnResolveWithToken;
public bool Owned, BadSnapshot, Paused, Indexed = true; public bool ProviderPaused => Paused; public HashSet<int> Files = new(); public Func<Task>? OnResolve; public Func<CancellationToken, Task>? OnResolveWithToken; public Func<ImportEpisode, Task>? OnObserve;
public Task<ImportSnapshot> FetchAsync(CatalogItem item, CancellationToken ct) => Task.FromResult(new ImportSnapshot(BadSnapshot ? new() : Enumerable.Range(1, Count).Select(Ep).ToList(), "success"));
public Task<ImportObservation> 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<ImportObservation> 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<List<SelectedVersion>> 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" } } }; }
Expand Down
13 changes: 8 additions & 5 deletions docs/import-reconciliation.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading