RED: ExpireRegistrationWorker loads the registration a RegistratieVerlopen job correlates to and expires it (idempotent on redelivery, throws on unknown so the job is redelivered); RegistratieVerlopenProcessor drains the parked jobs and completes each, leaving a failing one un-completed (§8.6). Mirrors the OpenZaak and escalation worker/processor pairs. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
80 lines
3.4 KiB
C#
80 lines
3.4 KiB
C#
using Big.Application;
|
|
using Big.Infrastructure;
|
|
using Microsoft.Extensions.Logging;
|
|
using Microsoft.Extensions.Logging.Abstractions;
|
|
|
|
namespace Big.Tests;
|
|
|
|
// S-10a (#102): the document-timeout drain loop. Mirrors BeoordelingEscalatieProcessor — acquire the
|
|
// parked RegistratieVerlopen jobs (the tokens the 30-day boundary timer on WachtOpDocumenten spawns),
|
|
// expire each correlated registration via the ExpireRegistrationWorker, then complete the job. A job
|
|
// whose expiry fails is logged and left un-completed for Flowable to redeliver (§8.6).
|
|
public class RegistratieVerlopenProcessorTests
|
|
{
|
|
/// <summary>A fake client scripting the jobs to acquire and recording completions.</summary>
|
|
private sealed class FakeVerlopenClient(params RegistratieVerlopenJob[] jobs) : IRegistratieVerlopenClient
|
|
{
|
|
public int AcquireCount { get; private set; }
|
|
public List<string> Completed { get; } = [];
|
|
|
|
public Task<IReadOnlyList<RegistratieVerlopenJob>> AcquireRegistratieVerlopenJobsAsync(int maxJobs, CancellationToken ct = default)
|
|
{
|
|
AcquireCount++;
|
|
return Task.FromResult<IReadOnlyList<RegistratieVerlopenJob>>(jobs.Take(maxJobs).ToList());
|
|
}
|
|
|
|
public Task CompleteRegistratieVerlopenJobAsync(string jobId, CancellationToken ct = default)
|
|
{
|
|
Completed.Add(jobId);
|
|
return Task.CompletedTask;
|
|
}
|
|
}
|
|
|
|
private static ExpireRegistrationWorker Worker(FakeRegistrationStore store) => new(store);
|
|
|
|
[Fact]
|
|
public async Task Acquires_a_job_expires_the_registration_and_completes_the_job()
|
|
{
|
|
var store = new FakeRegistrationStore();
|
|
var registration = Domain.Registration.Submit("123456782");
|
|
store.Seed(registration);
|
|
var client = new FakeVerlopenClient(new RegistratieVerlopenJob("job-9", registration.Id));
|
|
|
|
var acquired = await new RegistratieVerlopenProcessor(
|
|
client, Worker(store), NullLogger<RegistratieVerlopenProcessor>.Instance).PumpOnceAsync(5);
|
|
|
|
Assert.Equal(1, acquired);
|
|
Assert.Equal(Domain.RegistrationStatus.Verlopen, (await store.GetAsync(registration.Id))!.Status);
|
|
Assert.Equal("job-9", Assert.Single(client.Completed));
|
|
}
|
|
|
|
[Fact]
|
|
public async Task A_failing_expiry_is_left_uncompleted_for_flowable_to_redeliver()
|
|
{
|
|
// Unknown registration → the worker throws → the job is left for redelivery, error logged.
|
|
var store = new FakeRegistrationStore();
|
|
var client = new FakeVerlopenClient(new RegistratieVerlopenJob("job-9", Domain.RegistrationId.New()));
|
|
var logger = new CapturingLogger<RegistratieVerlopenProcessor>();
|
|
|
|
var acquired = await new RegistratieVerlopenProcessor(client, Worker(store), logger).PumpOnceAsync(5);
|
|
|
|
Assert.Equal(1, acquired);
|
|
Assert.Empty(client.Completed);
|
|
var error = Assert.Single(logger.Entries, e => e.Level == LogLevel.Error);
|
|
Assert.Contains("job-9", error.Message);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Does_nothing_but_poll_when_there_are_no_jobs()
|
|
{
|
|
var client = new FakeVerlopenClient();
|
|
|
|
var acquired = await new RegistratieVerlopenProcessor(
|
|
client, Worker(new FakeRegistrationStore()), NullLogger<RegistratieVerlopenProcessor>.Instance).PumpOnceAsync(5);
|
|
|
|
Assert.Equal(0, acquired);
|
|
Assert.Equal(1, client.AcquireCount);
|
|
Assert.Empty(client.Completed);
|
|
}
|
|
}
|