Files
Cleanuparr/code/backend/Cleanuparr.Infrastructure.Tests/Features/Jobs/SeekerCommandMonitorTests.cs
T

1060 lines
42 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;
using Cleanuparr.Infrastructure.Features.Arr.Interfaces;
using Cleanuparr.Infrastructure.Features.Jobs;
using Cleanuparr.Infrastructure.Tests.Features.Jobs.TestHelpers;
using Cleanuparr.Persistence;
using Cleanuparr.Persistence.Models.Configuration.Arr;
using Cleanuparr.Persistence.Models.State;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Time.Testing;
using Microsoft.EntityFrameworkCore;
using NSubstitute;
using NSubstitute.ExceptionExtensions;
using Shouldly;
using Xunit;
namespace Cleanuparr.Infrastructure.Tests.Features.Jobs;
public class SeekerCommandMonitorTests : IAsyncDisposable
{
private readonly DataContext _dataContext;
private readonly EventsContext _eventsContext;
private readonly FakeTimeProvider _timeProvider;
private readonly IArrClient _arrClient;
private readonly IEventPublisher _eventPublisher;
private readonly SeekerCommandMonitor _sut;
private readonly CancellationTokenSource _cts;
public SeekerCommandMonitorTests()
{
_dataContext = TestDataContextFactory.Create();
_eventsContext = TestDataContextFactory.CreateEvents();
_timeProvider = new FakeTimeProvider();
_arrClient = Substitute.For<IArrClient>();
_eventPublisher = Substitute.For<IEventPublisher>();
_cts = new CancellationTokenSource();
var logger = Substitute.For<ILogger<SeekerCommandMonitor>>();
var arrClientFactory = Substitute.For<IArrClientFactory>();
var serviceProvider = Substitute.For<IServiceProvider>();
serviceProvider.GetService(typeof(DataContext)).Returns(_dataContext);
serviceProvider.GetService(typeof(EventsContext)).Returns(_eventsContext);
serviceProvider.GetService(typeof(IArrClientFactory)).Returns(arrClientFactory);
serviceProvider.GetService(typeof(IArrQueueIterator))
.Returns(new ArrQueueIterator(Substitute.For<ILogger<ArrQueueIterator>>()));
serviceProvider.GetService(typeof(IEventPublisher)).Returns(_eventPublisher);
var scope = Substitute.For<IServiceScope>();
scope.ServiceProvider.Returns(serviceProvider);
var scopeFactory = Substitute.For<IServiceScopeFactory>();
scopeFactory.CreateScope().Returns(scope);
arrClientFactory.GetClient(Arg.Any<InstanceType>(), Arg.Any<float>()).Returns(_arrClient);
_sut = new SeekerCommandMonitor(logger, scopeFactory, _timeProvider);
}
public async ValueTask DisposeAsync()
{
await _cts.CancelAsync();
try { await _sut.StopAsync(CancellationToken.None); }
catch { /* expected during teardown */ }
_sut.Dispose();
_dataContext.Dispose();
_eventsContext.Dispose();
_cts.Dispose();
GC.SuppressFinalize(this);
}
[Fact]
public async Task Deduplicates_grabbed_items_by_download_id_for_sonarr_season_packs()
{
// Arrange
var sonarrInstance = TestDataContextFactory.AddSonarrInstance(_dataContext);
var eventId = Guid.NewGuid();
_eventsContext.SeekerCommandTrackers.Add(new SeekerCommandTracker
{
ArrInstanceId = sonarrInstance.Id,
CommandId = 1,
EventId = eventId,
ExternalItemId = 100,
ItemTitle = "Test Series - Season 1",
SeasonNumber = 1,
Status = SearchCommandStatus.Pending,
CreatedAt = _timeProvider.GetUtcNow().UtcDateTime
});
await _dataContext.SaveChangesAsync();
await _eventsContext.SaveChangesAsync();
StubCommandState(ArrCommandState.Completed);
// 3 episodes from same season pack share the same DownloadId
_arrClient.GetQueueItemsAsync(Arg.Any<ArrInstance>(), Arg.Any<int>())
.Returns(new QueueListResponse
{
TotalRecords = 3,
Records =
[
new QueueRecord { Id = 1, SeriesId = 100, SeasonNumber = 1, Title = "Test.Series.S01.1080p", DownloadId = "ABC123", Protocol = "torrent", Status = "downloading" },
new QueueRecord { Id = 2, SeriesId = 100, SeasonNumber = 1, Title = "Test.Series.S01.1080p", DownloadId = "ABC123", Protocol = "torrent", Status = "downloading" },
new QueueRecord { Id = 3, SeriesId = 100, SeasonNumber = 1, Title = "Test.Series.S01.1080p", DownloadId = "ABC123", Protocol = "torrent", Status = "downloading" },
]
});
var publishTcs = new TaskCompletionSource<List<string>?>();
_eventPublisher.PublishSearchCompleted(
Arg.Any<Guid>(), Arg.Any<SearchCommandStatus>(), Arg.Any<InstanceType>(), Arg.Any<string>(), Arg.Any<List<string>?>())
.Returns(Task.CompletedTask)
.AndDoes(ci => publishTcs.TrySetResult(ci.ArgAt<List<string>?>(4)));
// Act
await _sut.StartAsync(_cts.Token);
_timeProvider.Advance(TimeSpan.FromSeconds(11));
var resultData = await publishTcs.Task.WaitAsync(TimeSpan.FromSeconds(5));
// Assert
await _eventPublisher.Received(1).PublishSearchCompleted(
eventId, SearchCommandStatus.Completed, Arg.Any<InstanceType>(), Arg.Any<string>(), Arg.Any<List<string>?>());
resultData.ShouldNotBeNull();
resultData!.Count.ShouldBe(1);
}
[Fact]
public async Task Filters_out_records_with_empty_download_id()
{
// Arrange
var sonarrInstance = TestDataContextFactory.AddSonarrInstance(_dataContext);
var eventId = Guid.NewGuid();
_eventsContext.SeekerCommandTrackers.Add(new SeekerCommandTracker
{
ArrInstanceId = sonarrInstance.Id,
CommandId = 1,
EventId = eventId,
ExternalItemId = 100,
ItemTitle = "Test Series - Season 1",
SeasonNumber = 1,
Status = SearchCommandStatus.Pending,
CreatedAt = _timeProvider.GetUtcNow().UtcDateTime
});
await _dataContext.SaveChangesAsync();
await _eventsContext.SaveChangesAsync();
StubCommandState(ArrCommandState.Completed);
// Queue has records with empty DownloadId and one valid record
_arrClient.GetQueueItemsAsync(Arg.Any<ArrInstance>(), Arg.Any<int>())
.Returns(new QueueListResponse
{
TotalRecords = 3,
Records =
[
new QueueRecord { Id = 1, SeriesId = 100, SeasonNumber = 1, Title = "Empty DL 1", DownloadId = "", Protocol = "torrent", Status = "downloading" },
new QueueRecord { Id = 2, SeriesId = 100, SeasonNumber = 1, Title = "Empty DL 2", DownloadId = "", Protocol = "torrent", Status = "downloading" },
new QueueRecord { Id = 3, SeriesId = 100, SeasonNumber = 1, Title = "Valid Download", DownloadId = "VALID123", Protocol = "torrent", Status = "downloading" },
]
});
var publishTcs = new TaskCompletionSource<List<string>?>();
_eventPublisher.PublishSearchCompleted(
Arg.Any<Guid>(), Arg.Any<SearchCommandStatus>(), Arg.Any<InstanceType>(), Arg.Any<string>(), Arg.Any<List<string>?>())
.Returns(Task.CompletedTask)
.AndDoes(ci => publishTcs.TrySetResult(ci.ArgAt<List<string>?>(4)));
// Act
await _sut.StartAsync(_cts.Token);
_timeProvider.Advance(TimeSpan.FromSeconds(11));
var resultData = await publishTcs.Task.WaitAsync(TimeSpan.FromSeconds(5));
// Assert
resultData.ShouldNotBeNull();
resultData!.Count.ShouldBe(1);
resultData[0].ShouldBe("Valid Download");
}
[Fact]
public async Task Reports_multiple_grabbed_items_with_different_download_ids()
{
// Arrange
var radarrInstance = TestDataContextFactory.AddRadarrInstance(_dataContext);
var eventId = Guid.NewGuid();
_eventsContext.SeekerCommandTrackers.Add(new SeekerCommandTracker
{
ArrInstanceId = radarrInstance.Id,
CommandId = 1,
EventId = eventId,
ExternalItemId = 200,
ItemTitle = "Test Movie",
SeasonNumber = 0,
Status = SearchCommandStatus.Pending,
CreatedAt = _timeProvider.GetUtcNow().UtcDateTime
});
await _dataContext.SaveChangesAsync();
await _eventsContext.SaveChangesAsync();
StubCommandState(ArrCommandState.Completed);
// Two different downloads for the same movie
_arrClient.GetQueueItemsAsync(Arg.Any<ArrInstance>(), Arg.Any<int>())
.Returns(new QueueListResponse
{
TotalRecords = 2,
Records =
[
new QueueRecord { Id = 1, MovieId = 200, Title = "Test.Movie.720p", DownloadId = "HASH1", Protocol = "torrent", Status = "downloading" },
new QueueRecord { Id = 2, MovieId = 200, Title = "Test.Movie.1080p", DownloadId = "HASH2", Protocol = "usenet", Status = "downloading" },
]
});
var publishTcs = new TaskCompletionSource<List<string>?>();
_eventPublisher.PublishSearchCompleted(
Arg.Any<Guid>(), Arg.Any<SearchCommandStatus>(), Arg.Any<InstanceType>(), Arg.Any<string>(), Arg.Any<List<string>?>())
.Returns(Task.CompletedTask)
.AndDoes(ci => publishTcs.TrySetResult(ci.ArgAt<List<string>?>(4)));
// Act
await _sut.StartAsync(_cts.Token);
_timeProvider.Advance(TimeSpan.FromSeconds(11));
var resultData = await publishTcs.Task.WaitAsync(TimeSpan.FromSeconds(5));
// Assert
resultData.ShouldNotBeNull();
resultData!.Count.ShouldBe(2);
}
[Fact]
public async Task Publishes_timed_out_status_when_command_exceeds_timeout()
{
// Arrange
ArrInstance radarrInstance = TestDataContextFactory.AddRadarrInstance(_dataContext);
Guid eventId = Guid.NewGuid();
_eventsContext.SeekerCommandTrackers.Add(new SeekerCommandTracker
{
ArrInstanceId = radarrInstance.Id,
CommandId = 7,
EventId = eventId,
ExternalItemId = 300,
ItemTitle = "Stuck Movie",
SeasonNumber = 0,
Status = SearchCommandStatus.Pending,
CreatedAt = _timeProvider.GetUtcNow().UtcDateTime - TimeSpan.FromMinutes(31)
});
await _dataContext.SaveChangesAsync();
await _eventsContext.SaveChangesAsync();
StubCommandState(ArrCommandState.Started, commandId: 7);
Task<SearchCommandStatus> publishTask = CaptureNextPublishedStatus();
// Act
await _sut.StartAsync(_cts.Token);
_timeProvider.Advance(TimeSpan.FromSeconds(11));
SearchCommandStatus publishedStatus = await publishTask.WaitAsync(TimeSpan.FromSeconds(5));
// Assert
publishedStatus.ShouldBe(SearchCommandStatus.TimedOut);
await _eventPublisher.Received(1).PublishSearchCompleted(
eventId, SearchCommandStatus.TimedOut, Arg.Any<InstanceType>(), Arg.Any<string>(), Arg.Any<List<string>?>());
await _arrClient.Received(1).GetCommandsAsync(Arg.Any<ArrInstance>());
await _arrClient.DidNotReceive().GetCommandStatusAsync(Arg.Any<ArrInstance>(), Arg.Any<long>());
await _arrClient.DidNotReceive().GetQueueItemsAsync(Arg.Any<ArrInstance>(), Arg.Any<int>());
}
[Fact]
public async Task Completes_a_command_that_finished_just_after_the_timeout()
{
// Arrange
ArrInstance radarrInstance = TestDataContextFactory.AddRadarrInstance(_dataContext);
Guid eventId = Guid.NewGuid();
AddTracker(radarrInstance.Id, eventId, commandId: 7, age: TimeSpan.FromMinutes(31));
await _dataContext.SaveChangesAsync();
await _eventsContext.SaveChangesAsync();
StubCommandState(ArrCommandState.Completed, commandId: 7);
_arrClient.GetQueueItemsAsync(Arg.Any<ArrInstance>(), Arg.Any<int>())
.Returns(new QueueListResponse());
Task<SearchCommandStatus> publishTask = CaptureNextPublishedStatus();
// Act
await _sut.StartAsync(_cts.Token);
_timeProvider.Advance(TimeSpan.FromSeconds(11));
SearchCommandStatus publishedStatus = await publishTask.WaitAsync(TimeSpan.FromSeconds(5));
// Assert
publishedStatus.ShouldBe(SearchCommandStatus.Completed);
}
[Fact]
public async Task Does_not_time_out_a_command_that_is_still_within_the_timeout()
{
// Arrange
ArrInstance radarrInstance = TestDataContextFactory.AddRadarrInstance(_dataContext);
AddTracker(radarrInstance.Id, Guid.NewGuid(), age: TimeSpan.FromMinutes(10));
await _dataContext.SaveChangesAsync();
await _eventsContext.SaveChangesAsync();
StubCommandState(ArrCommandState.Started);
Task secondPoll = CaptureNthPoll(2);
CaptureNextPublishedStatus();
// Act
await _sut.StartAsync(_cts.Token);
await AdvanceUntilAsync(secondPoll);
// Assert
await _eventPublisher.DidNotReceive().PublishSearchCompleted(
Arg.Any<Guid>(), Arg.Any<SearchCommandStatus>(), Arg.Any<InstanceType>(), Arg.Any<string>(), Arg.Any<List<string>?>());
SeekerCommandTracker tracker = _eventsContext.SeekerCommandTrackers.AsNoTracking().Single();
tracker.Status.ShouldBe(SearchCommandStatus.Started);
}
[Fact]
public async Task Publishes_failed_status_when_the_command_reports_failed()
{
// Arrange
ArrInstance radarrInstance = TestDataContextFactory.AddRadarrInstance(_dataContext);
Guid eventId = Guid.NewGuid();
AddTracker(radarrInstance.Id, eventId);
await _dataContext.SaveChangesAsync();
await _eventsContext.SaveChangesAsync();
StubCommandState(ArrCommandState.Failed);
Task<SearchCommandStatus> publishTask = CaptureNextPublishedStatus();
// Act
await _sut.StartAsync(_cts.Token);
_timeProvider.Advance(TimeSpan.FromSeconds(11));
SearchCommandStatus publishedStatus = await publishTask.WaitAsync(TimeSpan.FromSeconds(5));
// Assert
publishedStatus.ShouldBe(SearchCommandStatus.Failed);
await _arrClient.DidNotReceive().GetQueueItemsAsync(Arg.Any<ArrInstance>(), Arg.Any<int>());
}
[Theory]
[InlineData(ArrCommandState.Aborted)]
[InlineData(ArrCommandState.Cancelled)]
[InlineData(ArrCommandState.Orphaned)]
public async Task Publishes_failed_status_for_every_unsuccessful_arr_command_state(ArrCommandState state)
{
// Arrange
ArrInstance radarrInstance = TestDataContextFactory.AddRadarrInstance(_dataContext);
Guid eventId = Guid.NewGuid();
AddTracker(radarrInstance.Id, eventId);
await _dataContext.SaveChangesAsync();
await _eventsContext.SaveChangesAsync();
StubCommandState(state);
Task<SearchCommandStatus> publishTask = CaptureNextPublishedStatus();
// Act
await _sut.StartAsync(_cts.Token);
_timeProvider.Advance(TimeSpan.FromSeconds(11));
SearchCommandStatus publishedStatus = await publishTask.WaitAsync(TimeSpan.FromSeconds(5));
// Assert
publishedStatus.ShouldBe(SearchCommandStatus.Failed);
}
[Fact]
public async Task Publishes_completed_status_when_arr_no_longer_knows_the_command()
{
// Arrange
ArrInstance radarrInstance = TestDataContextFactory.AddRadarrInstance(_dataContext);
Guid eventId = Guid.NewGuid();
AddTracker(radarrInstance.Id, eventId);
await _dataContext.SaveChangesAsync();
await _eventsContext.SaveChangesAsync();
_arrClient.GetCommandsAsync(Arg.Any<ArrInstance>())
.ThrowsAsync(new HttpRequestException("Not found", null, HttpStatusCode.NotFound));
_arrClient.GetCommandStatusAsync(Arg.Any<ArrInstance>(), Arg.Any<long>())
.ThrowsAsync(new HttpRequestException("Not found", null, HttpStatusCode.NotFound));
_arrClient.GetQueueItemsAsync(Arg.Any<ArrInstance>(), Arg.Any<int>())
.Returns(new QueueListResponse());
Task<SearchCommandStatus> publishTask = CaptureNextPublishedStatus();
// Act
await _sut.StartAsync(_cts.Token);
_timeProvider.Advance(TimeSpan.FromSeconds(11));
SearchCommandStatus publishedStatus = await publishTask.WaitAsync(TimeSpan.FromSeconds(5));
// Assert
publishedStatus.ShouldBe(SearchCommandStatus.Completed);
}
[Fact]
public async Task Keeps_polling_when_the_arr_command_state_is_not_recognized()
{
// Arrange
ArrInstance radarrInstance = TestDataContextFactory.AddRadarrInstance(_dataContext);
Guid eventId = Guid.NewGuid();
AddTracker(radarrInstance.Id, eventId);
await _dataContext.SaveChangesAsync();
await _eventsContext.SaveChangesAsync();
StubCommandState(ArrCommandState.Unknown);
Task secondPoll = CaptureNthPoll(2);
CaptureNextPublishedStatus();
// Act
await _sut.StartAsync(_cts.Token);
await AdvanceUntilAsync(secondPoll);
// Assert
await _eventPublisher.DidNotReceive().PublishSearchCompleted(
Arg.Any<Guid>(), Arg.Any<SearchCommandStatus>(), Arg.Any<InstanceType>(), Arg.Any<string>(), Arg.Any<List<string>?>());
SeekerCommandTracker tracker = _eventsContext.SeekerCommandTrackers.AsNoTracking().Single();
tracker.Status.ShouldBe(SearchCommandStatus.Pending);
}
[Fact]
public async Task Removes_the_tracker_after_publishing_the_outcome()
{
// Arrange
ArrInstance radarrInstance = TestDataContextFactory.AddRadarrInstance(_dataContext);
AddTracker(radarrInstance.Id, Guid.NewGuid());
await _dataContext.SaveChangesAsync();
await _eventsContext.SaveChangesAsync();
StubCommandState(ArrCommandState.Completed);
_arrClient.GetQueueItemsAsync(Arg.Any<ArrInstance>(), Arg.Any<int>())
.Returns(new QueueListResponse());
Task<SearchCommandStatus> publishTask = CaptureNextPublishedStatus();
// Act
await _sut.StartAsync(_cts.Token);
_timeProvider.Advance(TimeSpan.FromSeconds(11));
await publishTask.WaitAsync(TimeSpan.FromSeconds(5));
// Assert
await WaitForTrackerCountAsync(0);
}
[Fact]
public async Task Keeps_the_tracker_when_publishing_the_outcome_fails()
{
// Arrange
ArrInstance radarrInstance = TestDataContextFactory.AddRadarrInstance(_dataContext);
AddTracker(radarrInstance.Id, Guid.NewGuid());
await _dataContext.SaveChangesAsync();
await _eventsContext.SaveChangesAsync();
StubCommandState(ArrCommandState.Completed);
_arrClient.GetQueueItemsAsync(Arg.Any<ArrInstance>(), Arg.Any<int>())
.Returns(new QueueListResponse());
Task publishAttempt = FailNextPublish();
// Act
await _sut.StartAsync(_cts.Token);
_timeProvider.Advance(TimeSpan.FromSeconds(11));
await publishAttempt.WaitAsync(TimeSpan.FromSeconds(5));
// Assert
await WaitForTrackerCountAsync(1);
}
[Fact]
public async Task Publishes_a_terminal_tracker_left_over_from_a_previous_cycle()
{
// Arrange
ArrInstance radarrInstance = TestDataContextFactory.AddRadarrInstance(_dataContext);
Guid eventId = Guid.NewGuid();
AddTracker(radarrInstance.Id, eventId, status: SearchCommandStatus.Completed);
await _dataContext.SaveChangesAsync();
await _eventsContext.SaveChangesAsync();
_arrClient.GetQueueItemsAsync(Arg.Any<ArrInstance>(), Arg.Any<int>())
.Returns(new QueueListResponse());
Task<SearchCommandStatus> publishTask = CaptureNextPublishedStatus();
// Act
await _sut.StartAsync(_cts.Token);
_timeProvider.Advance(TimeSpan.FromSeconds(11));
SearchCommandStatus publishedStatus = await publishTask.WaitAsync(TimeSpan.FromSeconds(5));
// Assert
publishedStatus.ShouldBe(SearchCommandStatus.Completed);
await _arrClient.DidNotReceive().GetCommandStatusAsync(Arg.Any<ArrInstance>(), Arg.Any<long>());
}
[Fact]
public async Task Abandons_a_tracker_that_keeps_failing_to_publish()
{
// Arrange
ArrInstance radarrInstance = TestDataContextFactory.AddRadarrInstance(_dataContext);
AddTracker(radarrInstance.Id, Guid.NewGuid(), status: SearchCommandStatus.Completed, age: TimeSpan.FromMinutes(91));
await _dataContext.SaveChangesAsync();
await _eventsContext.SaveChangesAsync();
_arrClient.GetQueueItemsAsync(Arg.Any<ArrInstance>(), Arg.Any<int>())
.Returns(new QueueListResponse());
Task publishAttempt = FailNextPublish();
// Act
await _sut.StartAsync(_cts.Token);
_timeProvider.Advance(TimeSpan.FromSeconds(11));
await publishAttempt.WaitAsync(TimeSpan.FromSeconds(5));
// Assert
await WaitForTrackerCountAsync(0);
}
[Fact]
public async Task Fails_the_event_when_the_arr_instance_no_longer_exists()
{
// Arrange
Guid eventId = Guid.NewGuid();
AddTracker(Guid.NewGuid(), eventId, status: SearchCommandStatus.Completed);
await _eventsContext.SaveChangesAsync();
Task<SearchCommandStatus> publishTask = CaptureNextPublishedStatus();
// Act
await _sut.StartAsync(_cts.Token);
_timeProvider.Advance(TimeSpan.FromSeconds(11));
SearchCommandStatus publishedStatus = await publishTask.WaitAsync(TimeSpan.FromSeconds(5));
// Assert
publishedStatus.ShouldBe(SearchCommandStatus.Failed);
await _eventPublisher.Received(1).PublishSearchCompleted(
eventId, SearchCommandStatus.Failed, Arg.Any<InstanceType>(), Arg.Any<string>(), Arg.Any<List<string>?>());
await WaitForTrackerCountAsync(0);
}
[Fact]
public async Task Fails_a_pending_tracker_whose_arr_instance_no_longer_exists()
{
// Arrange
Guid eventId = Guid.NewGuid();
AddTracker(Guid.NewGuid(), eventId);
await _eventsContext.SaveChangesAsync();
Task<SearchCommandStatus> pendingPublishTask = CaptureNextPublishedStatus();
// Act
await _sut.StartAsync(_cts.Token);
_timeProvider.Advance(TimeSpan.FromSeconds(11));
SearchCommandStatus pendingPublishedStatus = await pendingPublishTask.WaitAsync(TimeSpan.FromSeconds(5));
// Assert
pendingPublishedStatus.ShouldBe(SearchCommandStatus.Failed);
await _eventPublisher.Received(1).PublishSearchCompleted(
eventId, SearchCommandStatus.Failed, Arg.Any<InstanceType>(), Arg.Any<string>(), Arg.Any<List<string>?>());
await _arrClient.DidNotReceive().GetCommandsAsync(Arg.Any<ArrInstance>());
await WaitForTrackerCountAsync(0);
}
[Fact]
public async Task Polls_every_command_of_an_instance_in_a_single_request()
{
// Arrange
ArrInstance radarrInstance = TestDataContextFactory.AddRadarrInstance(_dataContext);
AddTracker(radarrInstance.Id, Guid.NewGuid(), commandId: 1);
AddTracker(radarrInstance.Id, Guid.NewGuid(), commandId: 2);
AddTracker(radarrInstance.Id, Guid.NewGuid(), commandId: 3);
await _dataContext.SaveChangesAsync();
await _eventsContext.SaveChangesAsync();
_arrClient.GetCommandsAsync(Arg.Any<ArrInstance>())
.Returns(
[
new ArrCommandStatus(1, ArrCommandState.Completed, null),
new ArrCommandStatus(2, ArrCommandState.Completed, null),
new ArrCommandStatus(3, ArrCommandState.Completed, null),
]);
_arrClient.GetQueueItemsAsync(Arg.Any<ArrInstance>(), Arg.Any<int>())
.Returns(new QueueListResponse());
CaptureNextPublishedStatus();
// Act
await _sut.StartAsync(_cts.Token);
_timeProvider.Advance(TimeSpan.FromSeconds(11));
await WaitForTrackerCountAsync(0);
// Assert
await _arrClient.Received(1).GetCommandsAsync(Arg.Any<ArrInstance>());
await _arrClient.DidNotReceive().GetCommandStatusAsync(Arg.Any<ArrInstance>(), Arg.Any<long>());
await _arrClient.Received(1).GetQueueItemsAsync(Arg.Any<ArrInstance>(), Arg.Any<int>());
}
[Fact]
public async Task Completes_a_command_that_is_missing_from_the_command_list()
{
// Arrange
ArrInstance radarrInstance = TestDataContextFactory.AddRadarrInstance(_dataContext);
AddTracker(radarrInstance.Id, Guid.NewGuid(), commandId: 42);
await _dataContext.SaveChangesAsync();
await _eventsContext.SaveChangesAsync();
_arrClient.GetCommandsAsync(Arg.Any<ArrInstance>())
.Returns([new ArrCommandStatus(7, ArrCommandState.Started, null)]);
_arrClient.GetQueueItemsAsync(Arg.Any<ArrInstance>(), Arg.Any<int>())
.Returns(new QueueListResponse());
Task<SearchCommandStatus> publishTask = CaptureNextPublishedStatus();
// Act
await _sut.StartAsync(_cts.Token);
_timeProvider.Advance(TimeSpan.FromSeconds(11));
SearchCommandStatus publishedStatus = await publishTask.WaitAsync(TimeSpan.FromSeconds(5));
// Assert
publishedStatus.ShouldBe(SearchCommandStatus.Completed);
}
[Fact]
public async Task Falls_back_to_individual_command_checks_when_the_command_list_fails()
{
// Arrange
ArrInstance radarrInstance = TestDataContextFactory.AddRadarrInstance(_dataContext);
AddTracker(radarrInstance.Id, Guid.NewGuid());
await _dataContext.SaveChangesAsync();
await _eventsContext.SaveChangesAsync();
_arrClient.GetCommandsAsync(Arg.Any<ArrInstance>())
.ThrowsAsync(new HttpRequestException("boom"));
_arrClient.GetCommandStatusAsync(Arg.Any<ArrInstance>(), Arg.Any<long>())
.Returns(new ArrCommandStatus(1, ArrCommandState.Failed, null));
Task<SearchCommandStatus> publishTask = CaptureNextPublishedStatus();
// Act
await _sut.StartAsync(_cts.Token);
_timeProvider.Advance(TimeSpan.FromSeconds(11));
SearchCommandStatus publishedStatus = await publishTask.WaitAsync(TimeSpan.FromSeconds(5));
// Assert
publishedStatus.ShouldBe(SearchCommandStatus.Failed);
await _arrClient.Received(1).GetCommandStatusAsync(Arg.Any<ArrInstance>(), Arg.Any<long>());
}
[Fact]
public async Task Reports_grabbed_items_found_beyond_the_first_queue_page()
{
// Arrange
ArrInstance radarrInstance = TestDataContextFactory.AddRadarrInstance(_dataContext);
AddTracker(radarrInstance.Id, Guid.NewGuid());
await _dataContext.SaveChangesAsync();
await _eventsContext.SaveChangesAsync();
StubCommandState(ArrCommandState.Completed);
_arrClient.GetQueueItemsAsync(Arg.Any<ArrInstance>(), 1)
.Returns(new QueueListResponse
{
TotalRecords = 2,
Records = [new QueueRecord { Id = 1, MovieId = 999, Title = "Other.Movie", DownloadId = "OTHER", Protocol = "torrent", Status = "downloading" }]
});
_arrClient.GetQueueItemsAsync(Arg.Any<ArrInstance>(), 2)
.Returns(new QueueListResponse
{
TotalRecords = 2,
Records = [new QueueRecord { Id = 2, MovieId = 300, Title = "Wanted.Movie", DownloadId = "WANTED", Protocol = "torrent", Status = "downloading" }]
});
TaskCompletionSource<List<string>?> grabbedTcs = new();
_eventPublisher.PublishSearchCompleted(
Arg.Any<Guid>(), Arg.Any<SearchCommandStatus>(), Arg.Any<InstanceType>(), Arg.Any<string>(), Arg.Any<List<string>?>())
.Returns(Task.CompletedTask)
.AndDoes(ci => grabbedTcs.TrySetResult(ci.ArgAt<List<string>?>(4)));
// Act
await _sut.StartAsync(_cts.Token);
_timeProvider.Advance(TimeSpan.FromSeconds(11));
List<string>? grabbedItems = await grabbedTcs.Task.WaitAsync(TimeSpan.FromSeconds(5));
// Assert
grabbedItems.ShouldNotBeNull();
grabbedItems!.ShouldBe(["Wanted.Movie"]);
}
[Fact]
public async Task Publishes_completed_status_when_the_queue_cannot_be_inspected()
{
// Arrange
ArrInstance radarrInstance = TestDataContextFactory.AddRadarrInstance(_dataContext);
Guid eventId = Guid.NewGuid();
AddTracker(radarrInstance.Id, eventId);
await _dataContext.SaveChangesAsync();
await _eventsContext.SaveChangesAsync();
StubCommandState(ArrCommandState.Completed);
_arrClient.GetQueueItemsAsync(Arg.Any<ArrInstance>(), Arg.Any<int>())
.ThrowsAsync(new HttpRequestException("queue unavailable"));
Task<SearchCommandStatus> publishTask = CaptureNextPublishedStatus();
// Act
await _sut.StartAsync(_cts.Token);
_timeProvider.Advance(TimeSpan.FromSeconds(11));
SearchCommandStatus publishedStatus = await publishTask.WaitAsync(TimeSpan.FromSeconds(5));
// Assert
publishedStatus.ShouldBe(SearchCommandStatus.Completed);
await _eventPublisher.Received(1).PublishSearchCompleted(
eventId, SearchCommandStatus.Completed, Arg.Any<InstanceType>(), Arg.Any<string>(), null);
await WaitForTrackerCountAsync(0);
}
[Fact]
public async Task Keeps_polling_when_an_individual_command_check_fails()
{
// Arrange
ArrInstance radarrInstance = TestDataContextFactory.AddRadarrInstance(_dataContext);
AddTracker(radarrInstance.Id, Guid.NewGuid());
await _dataContext.SaveChangesAsync();
await _eventsContext.SaveChangesAsync();
_arrClient.GetCommandsAsync(Arg.Any<ArrInstance>())
.ThrowsAsync(new HttpRequestException("command list unavailable"));
_arrClient.GetCommandStatusAsync(Arg.Any<ArrInstance>(), Arg.Any<long>())
.ThrowsAsync(new HttpRequestException("command status unavailable"));
Task secondPoll = CaptureNthPoll(2);
CaptureNextPublishedStatus();
// Act
await _sut.StartAsync(_cts.Token);
await AdvanceUntilAsync(secondPoll);
// Assert
await _eventPublisher.DidNotReceive().PublishSearchCompleted(
Arg.Any<Guid>(), Arg.Any<SearchCommandStatus>(), Arg.Any<InstanceType>(), Arg.Any<string>(), Arg.Any<List<string>?>());
SeekerCommandTracker tracker = _eventsContext.SeekerCommandTrackers.AsNoTracking().Single();
tracker.Status.ShouldBe(SearchCommandStatus.Pending);
}
[Fact]
public async Task Propagates_the_started_status_to_the_search_event()
{
// Arrange
ArrInstance radarrInstance = TestDataContextFactory.AddRadarrInstance(_dataContext);
Guid eventId = Guid.NewGuid();
AddTracker(radarrInstance.Id, eventId);
await _dataContext.SaveChangesAsync();
await _eventsContext.SaveChangesAsync();
StubCommandState(ArrCommandState.Started);
TaskCompletionSource progressTcs = new();
_eventPublisher.PublishSearchStarted(Arg.Any<Guid>())
.Returns(Task.CompletedTask)
.AndDoes(_ => progressTcs.TrySetResult());
// Act
await _sut.StartAsync(_cts.Token);
_timeProvider.Advance(TimeSpan.FromSeconds(11));
await progressTcs.Task.WaitAsync(TimeSpan.FromSeconds(5));
// Assert
await _eventPublisher.Received(1).PublishSearchStarted(eventId);
await _eventPublisher.DidNotReceive().PublishSearchCompleted(
Arg.Any<Guid>(), Arg.Any<SearchCommandStatus>(), Arg.Any<InstanceType>(), Arg.Any<string>(), Arg.Any<List<string>?>());
}
[Fact]
public async Task Sweeps_search_events_that_never_got_a_tracker()
{
// Arrange
TaskCompletionSource sweepTcs = new();
_eventPublisher.FailAbandonedSearchEvents(Arg.Any<DateTimeOffset>())
.Returns(1)
.AndDoes(_ => sweepTcs.TrySetResult());
// Act
await _sut.StartAsync(_cts.Token);
_timeProvider.Advance(TimeSpan.FromSeconds(11));
await sweepTcs.Task.WaitAsync(TimeSpan.FromSeconds(5));
// Assert
await _eventPublisher.Received(1).FailAbandonedSearchEvents(
_timeProvider.GetUtcNow() - TimeSpan.FromMinutes(60));
}
[Fact]
public async Task Reports_the_grab_of_a_sonarr_episode_search()
{
// Arrange
ArrInstance sonarrInstance = TestDataContextFactory.AddSonarrInstance(_dataContext);
// Act
List<string>? grabbedItems = await CaptureGrabbedItemsAsync(
Tracker(sonarrInstance.Id, externalItemId: 42, episodeId: 9001, itemTitle: "Example Show - S03E04"),
new QueueRecord { Id = 1, SeriesId = 42, EpisodeId = 9001, SeasonNumber = 3, Title = "Example.Show.S03E04.1080p", DownloadId = "EP9001", Protocol = "torrent", Status = "downloading" });
// Assert
grabbedItems.ShouldNotBeNull();
grabbedItems!.ShouldBe(["Example.Show.S03E04.1080p"]);
}
[Fact]
public async Task Ignores_another_episode_of_the_same_series()
{
// Arrange
ArrInstance sonarrInstance = TestDataContextFactory.AddSonarrInstance(_dataContext);
// Act
List<string>? grabbedItems = await CaptureGrabbedItemsAsync(
Tracker(sonarrInstance.Id, externalItemId: 42, episodeId: 9001, itemTitle: "Example Show - S03E04"),
new QueueRecord { Id = 1, SeriesId = 42, EpisodeId = 9001, SeasonNumber = 3, Title = "Example.Show.S03E04.1080p", DownloadId = "EP9001", Protocol = "torrent", Status = "downloading" },
new QueueRecord { Id = 2, SeriesId = 42, EpisodeId = 9002, SeasonNumber = 3, Title = "Example.Show.S03E05.1080p", DownloadId = "EP9002", Protocol = "torrent", Status = "downloading" });
// Assert
grabbedItems.ShouldNotBeNull();
grabbedItems!.ShouldBe(["Example.Show.S03E04.1080p"]);
}
[Fact]
public async Task Ignores_another_season_of_a_specials_search()
{
// Arrange
ArrInstance sonarrInstance = TestDataContextFactory.AddSonarrInstance(_dataContext);
// Act
List<string>? grabbedItems = await CaptureGrabbedItemsAsync(
Tracker(sonarrInstance.Id, externalItemId: 42, seasonNumber: 0, itemTitle: "Example Show S00"),
new QueueRecord { Id = 1, SeriesId = 42, EpisodeId = 9001, SeasonNumber = 3, Title = "Example.Show.S03E04.1080p", DownloadId = "EP9001", Protocol = "torrent", Status = "downloading" });
// Assert
grabbedItems.ShouldBeNull();
}
[Fact]
public async Task Reports_the_grab_of_a_whisparr_v3_search()
{
// Arrange
ArrInstance whisparrInstance = TestDataContextFactory.AddWhisparrInstance(_dataContext, version: 3);
// Act
List<string>? grabbedItems = await CaptureGrabbedItemsAsync(
Tracker(whisparrInstance.Id, externalItemId: 55, itemTitle: "Example Scene"),
new QueueRecord { Id = 1, MovieId = 55, Title = "Example.Scene.1080p", DownloadId = "MV55", Protocol = "torrent", Status = "downloading" });
// Assert
grabbedItems.ShouldNotBeNull();
grabbedItems!.ShouldBe(["Example.Scene.1080p"]);
}
[Fact]
public async Task Reports_the_grab_of_a_lidarr_search()
{
// Arrange
ArrInstance lidarrInstance = TestDataContextFactory.AddLidarrInstance(_dataContext);
// Act
List<string>? grabbedItems = await CaptureGrabbedItemsAsync(
Tracker(lidarrInstance.Id, externalItemId: 77, itemTitle: "Example Album"),
new QueueRecord { Id = 1, AlbumId = 77, Title = "Example.Album.FLAC", DownloadId = "AL77", Protocol = "torrent", Status = "downloading" });
// Assert
grabbedItems.ShouldNotBeNull();
grabbedItems!.ShouldBe(["Example.Album.FLAC"]);
}
[Fact]
public async Task Reports_the_grab_of_a_readarr_search()
{
// Arrange
ArrInstance readarrInstance = TestDataContextFactory.AddReadarrInstance(_dataContext);
// Act
List<string>? grabbedItems = await CaptureGrabbedItemsAsync(
Tracker(readarrInstance.Id, externalItemId: 88, itemTitle: "Example Book"),
new QueueRecord { Id = 1, BookId = 88, Title = "Example.Book.EPUB", DownloadId = "BK88", Protocol = "torrent", Status = "downloading" });
// Assert
grabbedItems.ShouldNotBeNull();
grabbedItems!.ShouldBe(["Example.Book.EPUB"]);
}
private SeekerCommandTracker Tracker(
Guid arrInstanceId,
long externalItemId,
long episodeId = 0,
int seasonNumber = 0,
string itemTitle = "Test Item") =>
new()
{
ArrInstanceId = arrInstanceId,
CommandId = 1,
EventId = Guid.NewGuid(),
ExternalItemId = externalItemId,
EpisodeId = episodeId,
ItemTitle = itemTitle,
SeasonNumber = seasonNumber,
Status = SearchCommandStatus.Pending,
CreatedAt = _timeProvider.GetUtcNow().UtcDateTime,
};
/// <summary>
/// Runs one monitor cycle over a completed command and returns the titles it attributed.
/// </summary>
private async Task<List<string>?> CaptureGrabbedItemsAsync(SeekerCommandTracker tracker, params QueueRecord[] records)
{
_eventsContext.SeekerCommandTrackers.Add(tracker);
await _dataContext.SaveChangesAsync();
await _eventsContext.SaveChangesAsync();
StubCommandState(ArrCommandState.Completed, tracker.CommandId);
_arrClient.GetQueueItemsAsync(Arg.Any<ArrInstance>(), Arg.Any<int>())
.Returns(new QueueListResponse { TotalRecords = records.Length, Records = [.. records] });
TaskCompletionSource<List<string>?> publishTcs = new();
_eventPublisher.PublishSearchCompleted(
Arg.Any<Guid>(), Arg.Any<SearchCommandStatus>(), Arg.Any<InstanceType>(), Arg.Any<string>(), Arg.Any<List<string>?>())
.Returns(Task.CompletedTask)
.AndDoes(ci => publishTcs.TrySetResult(ci.ArgAt<List<string>?>(4)));
await _sut.StartAsync(_cts.Token);
_timeProvider.Advance(TimeSpan.FromSeconds(11));
return await publishTcs.Task.WaitAsync(TimeSpan.FromSeconds(5));
}
private void StubCommandState(ArrCommandState state, long commandId = 1)
{
_arrClient.GetCommandsAsync(Arg.Any<ArrInstance>())
.Returns([new ArrCommandStatus(commandId, state, null)]);
_arrClient.GetCommandStatusAsync(Arg.Any<ArrInstance>(), Arg.Any<long>())
.Returns(new ArrCommandStatus(commandId, state, null));
}
private void AddTracker(
Guid arrInstanceId,
Guid eventId,
long commandId = 1,
TimeSpan? age = null,
SearchCommandStatus status = SearchCommandStatus.Pending)
{
_eventsContext.SeekerCommandTrackers.Add(new SeekerCommandTracker
{
ArrInstanceId = arrInstanceId,
CommandId = commandId,
EventId = eventId,
ExternalItemId = 300,
ItemTitle = "Test Item",
SeasonNumber = 0,
Status = status,
CreatedAt = _timeProvider.GetUtcNow().UtcDateTime - (age ?? TimeSpan.Zero)
});
}
private Task FailNextPublish()
{
TaskCompletionSource tcs = new();
_eventPublisher.PublishSearchCompleted(
Arg.Any<Guid>(), Arg.Any<SearchCommandStatus>(), Arg.Any<InstanceType>(), Arg.Any<string>(), Arg.Any<List<string>?>())
.Returns<Task>(_ =>
{
tcs.TrySetResult();
throw new InvalidOperationException("publish failed");
});
return tcs.Task;
}
private async Task WaitForTrackerCountAsync(int expected)
{
for (int attempt = 0; attempt < 200; attempt++)
{
try
{
if (await _eventsContext.SeekerCommandTrackers.CountAsync() == expected)
{
return;
}
}
catch (InvalidOperationException)
{
// the monitor is mid-cycle on the shared context
}
await Task.Delay(10);
}
(await _eventsContext.SeekerCommandTrackers.CountAsync()).ShouldBe(expected);
}
private Task<SearchCommandStatus> CaptureNextPublishedStatus()
{
TaskCompletionSource<SearchCommandStatus> tcs = new();
_eventPublisher.PublishSearchCompleted(
Arg.Any<Guid>(), Arg.Any<SearchCommandStatus>(), Arg.Any<InstanceType>(), Arg.Any<string>(), Arg.Any<List<string>?>())
.Returns(Task.CompletedTask)
.AndDoes(ci => tcs.TrySetResult(ci.ArgAt<SearchCommandStatus>(1)));
return tcs.Task;
}
private Task CaptureNthPoll(int count)
{
TaskCompletionSource tcs = new();
int seen = 0;
_arrClient
.When(client => client.GetCommandsAsync(Arg.Any<ArrInstance>()))
.Do(_ =>
{
if (Interlocked.Increment(ref seen) >= count)
{
tcs.TrySetResult();
}
});
return tcs.Task;
}
private async Task AdvanceUntilAsync(Task signal)
{
_timeProvider.Advance(TimeSpan.FromSeconds(11));
for (int attempt = 0; attempt < 100 && !signal.IsCompleted; attempt++)
{
await Task.Delay(10);
_timeProvider.Advance(TimeSpan.FromSeconds(60));
}
await signal.WaitAsync(TimeSpan.FromSeconds(5));
}
}