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
9 changes: 9 additions & 0 deletions Data/DatabaseManager.ImportCoverage.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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<List<CatalogItem>> 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<List<CatalogItem>> GetCatchUpImportCatalogAsync(DateTimeOffset started, DateTimeOffset now) => QueryListAsync(
Expand Down
6 changes: 3 additions & 3 deletions InfiniteDrive.csproj
Original file line number Diff line number Diff line change
Expand Up @@ -2,9 +2,9 @@

<PropertyGroup>
<TargetFramework>net8.0</TargetFramework>
<AssemblyVersion>0.42.13.0</AssemblyVersion>
<FileVersion>0.42.13.0</FileVersion>
<Version>0.42.13.0</Version>
<AssemblyVersion>0.42.14.0</AssemblyVersion>
<FileVersion>0.42.14.0</FileVersion>
<Version>0.42.14.0</Version>
<Nullable>enable</Nullable>
<RootNamespace>InfiniteDrive</RootNamespace>
<AssemblyName>InfiniteDrive</AssemblyName>
Expand Down
2 changes: 2 additions & 0 deletions Models/ImportCoverage.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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; }
}

/// <summary>Fresh per-file metadata bound to its current URL hash, without storing the URL.</summary>
Expand Down
78 changes: 78 additions & 0 deletions Models/ImportDiversityEntry.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,78 @@
using System;
using System.Collections.Generic;
using System.Linq;

namespace InfiniteDrive.Models;

/// <summary>One durable queue entry per coverage identity/episode, never a playback claim.</summary>
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<string>? Providers(ImportEpisode ep)
{
if (ep.Paths.Count == 0) return null;
var result = new HashSet<string>(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); }
Comment on lines +59 to +61

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Preserve the due time across routine refresh indexing

When a single-source item receives its normal hourly refresh, CompleteAsync temporarily changes its state to awaiting_indexing, so reconciliation retires the diversity entry; once Emby reports it indexed again, this branch resets NextAttempt to another six hours. Because ImportWorkBudget.NeedsRefresh makes the same item refreshable again after only one hour, small libraries that revisit items within six hours repeat this cycle indefinitely and never dispatch the new diversity lookup. Preserve the existing due time across this transient refresh/indexing state, or only restart the window when the saved provider set actually changes.

Useful? React with 👍 / 👎.

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);
}
}
9 changes: 5 additions & 4 deletions Models/ImportRunTelemetry.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, ImportStageTiming> Timings,
int HttpRequests = 0, int HttpRetries = 0);
int HttpRequests = 0, int HttpRetries = 0, int DiversityAttempts = 0, int DiversityPublished = 0);

/// <summary>One bounded, thread-safe run observation. Contains no target URLs or exception messages.</summary>
public sealed class ImportRunTelemetry
Expand All @@ -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)
{
Expand All @@ -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)
{
Expand All @@ -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)
{
Expand Down Expand Up @@ -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
Expand Down
1 change: 1 addition & 0 deletions PluginConfiguration.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down
16 changes: 14 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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**.
Expand Down
10 changes: 9 additions & 1 deletion Services/Api/ImportHealthService.cs
Original file line number Diff line number Diff line change
Expand Up @@ -45,12 +45,14 @@ public async Task<object> 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<object> Post(ImportActionRequest request)
Expand All @@ -71,6 +73,12 @@ internal static async Task<object> 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<ImportMode>(request.Mode, false, out var mode) || !Enum.IsDefined(mode))
return new { Status = "invalid_mode" };
Expand Down
Loading
Loading