mirror of
https://github.com/Cleanuparr/Cleanuparr.git
synced 2026-09-14 22:37:40 -04:00
438 lines
16 KiB
C#
438 lines
16 KiB
C#
using System.Net;
|
|
using Cleanuparr.Domain.Entities.Arr;
|
|
using Cleanuparr.Domain.Entities.Arr.Queue;
|
|
using Cleanuparr.Domain.Enums;
|
|
using Cleanuparr.Infrastructure.Events.Interfaces;
|
|
using Cleanuparr.Infrastructure.Features.Arr.Interfaces;
|
|
using Cleanuparr.Persistence;
|
|
using Cleanuparr.Persistence.Models.Configuration.Arr;
|
|
using Cleanuparr.Persistence.Models.State;
|
|
using Microsoft.EntityFrameworkCore;
|
|
using Microsoft.Extensions.DependencyInjection;
|
|
using Microsoft.Extensions.Hosting;
|
|
using Microsoft.Extensions.Logging;
|
|
|
|
namespace Cleanuparr.Infrastructure.Features.Jobs;
|
|
|
|
/// <summary>
|
|
/// Background service that polls arr command status for pending search commands
|
|
/// and inspects the download queue for grabbed items after completion.
|
|
/// </summary>
|
|
public class SeekerCommandMonitor : BackgroundService
|
|
{
|
|
private static readonly TimeSpan PollInterval = TimeSpan.FromSeconds(60);
|
|
private static readonly TimeSpan IdleInterval = TimeSpan.FromSeconds(60);
|
|
private static readonly TimeSpan CommandTimeout = TimeSpan.FromMinutes(30);
|
|
private static readonly TimeSpan AbandonAfter = TimeSpan.FromMinutes(90);
|
|
|
|
private readonly ILogger<SeekerCommandMonitor> _logger;
|
|
private readonly IServiceScopeFactory _scopeFactory;
|
|
private readonly TimeProvider _timeProvider;
|
|
|
|
public SeekerCommandMonitor(
|
|
ILogger<SeekerCommandMonitor> logger,
|
|
IServiceScopeFactory scopeFactory,
|
|
TimeProvider timeProvider)
|
|
{
|
|
_logger = logger;
|
|
_scopeFactory = scopeFactory;
|
|
_timeProvider = timeProvider;
|
|
}
|
|
|
|
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
|
|
{
|
|
// Wait for app startup
|
|
await Task.Delay(TimeSpan.FromSeconds(10), _timeProvider, stoppingToken);
|
|
|
|
while (!stoppingToken.IsCancellationRequested)
|
|
{
|
|
try
|
|
{
|
|
bool hadWork = await ProcessPendingCommandsAsync(stoppingToken);
|
|
await Task.Delay(hadWork ? PollInterval : IdleInterval, _timeProvider, stoppingToken);
|
|
}
|
|
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
|
|
{
|
|
break;
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
_logger.LogError(ex, "Error in SeekerCommandMonitor");
|
|
await Task.Delay(IdleInterval, _timeProvider, stoppingToken);
|
|
}
|
|
}
|
|
}
|
|
|
|
private async Task<bool> ProcessPendingCommandsAsync(CancellationToken stoppingToken)
|
|
{
|
|
await using AsyncServiceScope scope = _scopeFactory.CreateAsyncScope();
|
|
var dataContext = scope.ServiceProvider.GetRequiredService<DataContext>();
|
|
var eventsContext = scope.ServiceProvider.GetRequiredService<EventsContext>();
|
|
var arrClientFactory = scope.ServiceProvider.GetRequiredService<IArrClientFactory>();
|
|
var queueIterator = scope.ServiceProvider.GetRequiredService<IArrQueueIterator>();
|
|
var eventPublisher = scope.ServiceProvider.GetRequiredService<IEventPublisher>();
|
|
|
|
DateTimeOffset now = _timeProvider.GetUtcNow();
|
|
|
|
int abandonedEvents = await eventPublisher.FailAbandonedSearchEvents(now - (CommandTimeout * 2));
|
|
|
|
if (abandonedEvents > 0)
|
|
{
|
|
_logger.LogWarning("Failed {Count} search events that never got a command tracker", abandonedEvents);
|
|
}
|
|
|
|
List<SeekerCommandTracker> trackers = await eventsContext.SeekerCommandTrackers
|
|
.OrderBy(t => t.CreatedAt)
|
|
.ToListAsync(stoppingToken);
|
|
|
|
if (trackers.Count == 0)
|
|
{
|
|
return false;
|
|
}
|
|
|
|
Dictionary<Guid, ArrInstance> instancesById = await LoadInstancesAsync(dataContext, trackers, stoppingToken);
|
|
bool didWork = false;
|
|
|
|
List<SeekerCommandTracker> toPoll = [];
|
|
List<SeekerCommandTracker> expired = [];
|
|
|
|
foreach (SeekerCommandTracker tracker in trackers)
|
|
{
|
|
if (IsTerminal(tracker.Status))
|
|
{
|
|
continue;
|
|
}
|
|
|
|
bool timedOut = now - tracker.CreatedAt > CommandTimeout;
|
|
|
|
if (!instancesById.ContainsKey(tracker.ArrInstanceId))
|
|
{
|
|
tracker.Status = SearchCommandStatus.Failed;
|
|
continue;
|
|
}
|
|
|
|
if (timedOut)
|
|
{
|
|
expired.Add(tracker);
|
|
}
|
|
|
|
toPoll.Add(tracker);
|
|
}
|
|
|
|
foreach (IGrouping<Guid, SeekerCommandTracker> group in toPoll.GroupBy(t => t.ArrInstanceId))
|
|
{
|
|
ArrInstance arrInstance = instancesById[group.Key];
|
|
IArrClient arrClient = arrClientFactory.GetClient(arrInstance.ArrConfig.Type, arrInstance.Version);
|
|
didWork = true;
|
|
|
|
await PollInstanceCommandsAsync(arrClient, arrInstance, group.ToList(), eventPublisher);
|
|
}
|
|
|
|
foreach (SeekerCommandTracker tracker in expired.Where(tracker => !IsTerminal(tracker.Status)))
|
|
{
|
|
_logger.LogDebug(
|
|
"Command {CommandId} for '{Title}' is still running after {Timeout}, giving up",
|
|
tracker.CommandId, tracker.ItemTitle, CommandTimeout);
|
|
tracker.Status = SearchCommandStatus.TimedOut;
|
|
}
|
|
|
|
await eventsContext.SaveChangesAsync(stoppingToken);
|
|
|
|
Dictionary<Guid, IReadOnlyList<QueueRecord>> queueSnapshots = await BuildQueueSnapshotsAsync(
|
|
trackers, instancesById, arrClientFactory, queueIterator);
|
|
|
|
foreach (SeekerCommandTracker tracker in trackers.Where(t => IsTerminal(t.Status)))
|
|
{
|
|
didWork = true;
|
|
|
|
if (await TryPublishOutcomeAsync(tracker, instancesById, queueSnapshots, eventPublisher))
|
|
{
|
|
eventsContext.SeekerCommandTrackers.Remove(tracker);
|
|
continue;
|
|
}
|
|
|
|
if (now - tracker.CreatedAt > AbandonAfter)
|
|
{
|
|
_logger.LogWarning(
|
|
"Abandoning search command {CommandId} for '{Title}' after repeated publish failures (event {EventId})",
|
|
tracker.CommandId, tracker.ItemTitle, tracker.EventId);
|
|
eventsContext.SeekerCommandTrackers.Remove(tracker);
|
|
}
|
|
}
|
|
|
|
try
|
|
{
|
|
await eventsContext.SaveChangesAsync(stoppingToken);
|
|
}
|
|
catch (DbUpdateConcurrencyException ex)
|
|
{
|
|
_logger.LogWarning(ex, "Search command trackers changed while they were being processed");
|
|
}
|
|
|
|
return didWork;
|
|
}
|
|
|
|
private async Task PollInstanceCommandsAsync(
|
|
IArrClient arrClient,
|
|
ArrInstance arrInstance,
|
|
List<SeekerCommandTracker> trackers,
|
|
IEventPublisher eventPublisher)
|
|
{
|
|
Dictionary<long, ArrCommandStatus>? commands = await TryListCommandsAsync(arrClient, arrInstance);
|
|
|
|
if (commands is null)
|
|
{
|
|
foreach (SeekerCommandTracker tracker in trackers)
|
|
{
|
|
await PollSingleCommandAsync(arrClient, arrInstance, tracker, eventPublisher);
|
|
}
|
|
|
|
return;
|
|
}
|
|
|
|
foreach (SeekerCommandTracker tracker in trackers)
|
|
{
|
|
if (commands.TryGetValue(tracker.CommandId, out ArrCommandStatus? status))
|
|
{
|
|
await ApplyCommandStatusAsync(tracker, status, eventPublisher);
|
|
continue;
|
|
}
|
|
|
|
MarkForgottenCommandAsCompleted(tracker, arrInstance);
|
|
}
|
|
}
|
|
|
|
private async Task<Dictionary<long, ArrCommandStatus>?> TryListCommandsAsync(IArrClient arrClient, ArrInstance arrInstance)
|
|
{
|
|
try
|
|
{
|
|
List<ArrCommandStatus> commands = await arrClient.GetCommandsAsync(arrInstance);
|
|
|
|
return commands
|
|
.GroupBy(command => command.Id)
|
|
.ToDictionary(group => group.Key, group => group.First());
|
|
}
|
|
catch (OperationCanceledException)
|
|
{
|
|
throw;
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
_logger.LogWarning(ex,
|
|
"Failed to list commands on {Instance}, falling back to checking each command individually",
|
|
arrInstance.Name);
|
|
return null;
|
|
}
|
|
}
|
|
|
|
private async Task PollSingleCommandAsync(
|
|
IArrClient arrClient,
|
|
ArrInstance arrInstance,
|
|
SeekerCommandTracker tracker,
|
|
IEventPublisher eventPublisher)
|
|
{
|
|
try
|
|
{
|
|
ArrCommandStatus status = await arrClient.GetCommandStatusAsync(arrInstance, tracker.CommandId);
|
|
await ApplyCommandStatusAsync(tracker, status, eventPublisher);
|
|
}
|
|
catch (HttpRequestException ex) when (ex.StatusCode is HttpStatusCode.NotFound)
|
|
{
|
|
MarkForgottenCommandAsCompleted(tracker, arrInstance);
|
|
}
|
|
catch (OperationCanceledException)
|
|
{
|
|
throw;
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
_logger.LogWarning(ex, "Failed to check command {CommandId} status on {Instance}",
|
|
tracker.CommandId, arrInstance.Name);
|
|
}
|
|
}
|
|
|
|
private static async Task ApplyCommandStatusAsync(
|
|
SeekerCommandTracker tracker,
|
|
ArrCommandStatus status,
|
|
IEventPublisher eventPublisher)
|
|
{
|
|
SearchCommandStatus previousStatus = tracker.Status;
|
|
UpdateTrackerStatus(tracker, status);
|
|
|
|
if (tracker.Status is SearchCommandStatus.Started && previousStatus is not SearchCommandStatus.Started)
|
|
{
|
|
await eventPublisher.PublishSearchStarted(tracker.EventId);
|
|
}
|
|
}
|
|
|
|
private void MarkForgottenCommandAsCompleted(SeekerCommandTracker tracker, ArrInstance arrInstance)
|
|
{
|
|
_logger.LogDebug(
|
|
"Command {CommandId} is no longer known to {Instance}, treating '{Title}' as completed",
|
|
tracker.CommandId, arrInstance.Name, tracker.ItemTitle);
|
|
tracker.Status = SearchCommandStatus.Completed;
|
|
}
|
|
|
|
private async Task<Dictionary<Guid, IReadOnlyList<QueueRecord>>> BuildQueueSnapshotsAsync(
|
|
List<SeekerCommandTracker> trackers,
|
|
Dictionary<Guid, ArrInstance> instancesById,
|
|
IArrClientFactory arrClientFactory,
|
|
IArrQueueIterator queueIterator)
|
|
{
|
|
Dictionary<Guid, IReadOnlyList<QueueRecord>> snapshots = [];
|
|
|
|
List<Guid> instanceIds = trackers
|
|
.Where(tracker => tracker.Status is SearchCommandStatus.Completed)
|
|
.Select(tracker => tracker.ArrInstanceId)
|
|
.Distinct()
|
|
.Where(instancesById.ContainsKey)
|
|
.ToList();
|
|
|
|
foreach (Guid instanceId in instanceIds)
|
|
{
|
|
ArrInstance arrInstance = instancesById[instanceId];
|
|
List<QueueRecord> records = [];
|
|
|
|
try
|
|
{
|
|
IArrClient arrClient = arrClientFactory.GetClient(arrInstance.ArrConfig.Type, arrInstance.Version);
|
|
|
|
await queueIterator.Iterate(arrClient, arrInstance, pageRecords =>
|
|
{
|
|
records.AddRange(pageRecords);
|
|
return Task.CompletedTask;
|
|
});
|
|
}
|
|
catch (OperationCanceledException)
|
|
{
|
|
throw;
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
_logger.LogWarning(ex, "Failed to inspect the download queue on {Instance}", arrInstance.Name);
|
|
}
|
|
|
|
snapshots[instanceId] = records;
|
|
}
|
|
|
|
return snapshots;
|
|
}
|
|
|
|
private async Task<bool> TryPublishOutcomeAsync(
|
|
SeekerCommandTracker tracker,
|
|
Dictionary<Guid, ArrInstance> instancesById,
|
|
Dictionary<Guid, IReadOnlyList<QueueRecord>> queueSnapshots,
|
|
IEventPublisher eventPublisher)
|
|
{
|
|
try
|
|
{
|
|
if (!instancesById.TryGetValue(tracker.ArrInstanceId, out ArrInstance? arrInstance))
|
|
{
|
|
_logger.LogWarning(
|
|
"Failing search command {CommandId} for '{Title}': arr instance {ArrInstanceId} no longer exists (event {EventId})",
|
|
tracker.CommandId, tracker.ItemTitle, tracker.ArrInstanceId, tracker.EventId);
|
|
await eventPublisher.PublishSearchCompleted(tracker.EventId, SearchCommandStatus.Failed, default, string.Empty);
|
|
return true;
|
|
}
|
|
|
|
InstanceType instanceType = arrInstance.ArrConfig.Type;
|
|
string instanceUrl = arrInstance.ExternalOrInternalUrl.ToString();
|
|
|
|
if (tracker.Status is SearchCommandStatus.Failed or SearchCommandStatus.TimedOut)
|
|
{
|
|
await eventPublisher.PublishSearchCompleted(tracker.EventId, tracker.Status, instanceType, instanceUrl);
|
|
_logger.LogWarning(
|
|
"Search command {CommandId} for '{Title}' on {Instance} finished with status {Status} (event {EventId})",
|
|
tracker.CommandId, tracker.ItemTitle, arrInstance.Name, tracker.Status, tracker.EventId);
|
|
return true;
|
|
}
|
|
|
|
IReadOnlyList<QueueRecord> queue = queueSnapshots.GetValueOrDefault(tracker.ArrInstanceId, []);
|
|
List<string>? grabbedItems = FindGrabbedItems(tracker, arrInstance, queue);
|
|
await eventPublisher.PublishSearchCompleted(tracker.EventId, SearchCommandStatus.Completed, instanceType, instanceUrl, grabbedItems);
|
|
_logger.LogDebug("Search command completed for event {EventId}", tracker.EventId);
|
|
|
|
return true;
|
|
}
|
|
catch (OperationCanceledException)
|
|
{
|
|
throw;
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
_logger.LogError(ex,
|
|
"Failed to publish the outcome of search command {CommandId} for '{Title}' (event {EventId})",
|
|
tracker.CommandId, tracker.ItemTitle, tracker.EventId);
|
|
return false;
|
|
}
|
|
}
|
|
|
|
private static bool IsTerminal(SearchCommandStatus status) =>
|
|
status is SearchCommandStatus.Completed or SearchCommandStatus.Failed or SearchCommandStatus.TimedOut;
|
|
|
|
private static async Task<Dictionary<Guid, ArrInstance>> LoadInstancesAsync(
|
|
DataContext dataContext,
|
|
List<SeekerCommandTracker> trackers,
|
|
CancellationToken stoppingToken)
|
|
{
|
|
List<Guid> ids = trackers.Select(t => t.ArrInstanceId).Distinct().ToList();
|
|
return await dataContext.ArrInstances
|
|
.Include(a => a.ArrConfig)
|
|
.Where(a => ids.Contains(a.Id))
|
|
.ToDictionaryAsync(a => a.Id, stoppingToken);
|
|
}
|
|
|
|
private static void UpdateTrackerStatus(SeekerCommandTracker tracker, ArrCommandStatus commandStatus)
|
|
{
|
|
tracker.Status = commandStatus.Status switch
|
|
{
|
|
ArrCommandState.Completed => SearchCommandStatus.Completed,
|
|
ArrCommandState.Failed => SearchCommandStatus.Failed,
|
|
ArrCommandState.Aborted => SearchCommandStatus.Failed,
|
|
ArrCommandState.Cancelled => SearchCommandStatus.Failed,
|
|
ArrCommandState.Orphaned => SearchCommandStatus.Failed,
|
|
ArrCommandState.Started => SearchCommandStatus.Started,
|
|
_ => tracker.Status
|
|
};
|
|
}
|
|
|
|
private List<string>? FindGrabbedItems(
|
|
SeekerCommandTracker tracker,
|
|
ArrInstance arrInstance,
|
|
IReadOnlyList<QueueRecord> queue)
|
|
{
|
|
List<string> grabbedTitles = queue
|
|
.Where(r => MatchesTracker(r, tracker, arrInstance))
|
|
.Where(r => !string.IsNullOrEmpty(r.DownloadId))
|
|
.GroupBy(r => r.DownloadId)
|
|
.Select(g => g.First())
|
|
.Select(r => r.Title)
|
|
.ToList();
|
|
|
|
if (grabbedTitles.Count > 0)
|
|
{
|
|
_logger.LogInformation("Search for '{Title}' on {Instance} grabbed {Count} items: {Items}",
|
|
tracker.ItemTitle, arrInstance.Name, grabbedTitles.Count,
|
|
string.Join(", ", grabbedTitles));
|
|
}
|
|
|
|
return grabbedTitles.Count > 0 ? grabbedTitles : null;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Tells whether a queue record holds what the search asked for.
|
|
/// Each arr fills its own id on the record, so match the tracker against that id.
|
|
/// </summary>
|
|
private static bool MatchesTracker(QueueRecord record, SeekerCommandTracker tracker, ArrInstance arrInstance) =>
|
|
arrInstance.ArrConfig.Type switch
|
|
{
|
|
InstanceType.Radarr => record.MovieId == tracker.ExternalItemId,
|
|
InstanceType.Whisparr when arrInstance.Version is 3 => record.MovieId == tracker.ExternalItemId,
|
|
InstanceType.Lidarr => record.AlbumId == tracker.ExternalItemId,
|
|
InstanceType.Readarr => record.BookId == tracker.ExternalItemId,
|
|
_ => tracker.EpisodeId > 0
|
|
? record.EpisodeId == tracker.EpisodeId
|
|
: record.SeriesId == tracker.ExternalItemId && record.SeasonNumber == tracker.SeasonNumber,
|
|
};
|
|
}
|