## What & why S-10a, the **workflow/timeout spine** of the (split) document-upload slice: the registratie process now parks at a **`WachtOpDocumenten`** user task with an **interrupting `P30D` boundary timer**. When the documents arrive the task completes and the process continues into the diploma routing (S-13) → Beoordelen; if the 30 days lapse, the timer cancels the wait, runs a `RegistratieVerlopen` external-worker task, and the domain expires the aggregate to a new terminal status **`Verlopen`**. Backend only — the real upload trigger (portal → BFF → ACL → Documenten API) is S-10b (#103). Closes #102 Mechanism recorded in **ADR-0017**; opened as proposal #104. Mirrors the S-14 escalation (boundary-timer + external-worker) and S-11 withdrawal (interrupting cancel) patterns. ## Definition of Done - [x] Linked Gitea issue (above). - [x] Failing test committed before the implementation (red→green pairs per layer). - [x] Implementation makes the test pass. - [x] Conventional Commits referencing the issue (`refs #102`). - [ ] CI green — all Gitea Actions jobs (pending on this PR). - [x] `docker compose up` health unaffected (no new services; deploy path unchanged). - [x] Docs updated (ADR-0017, demo-script, BACKLOG split). - [x] ADR added (`docs/architecture/adr-0017-document-wait-timeout-cancellation.md`). - [x] Demo note in `docs/demo-script.md`. ## Notes for reviewers - **Domain** (`Registration.Expire()` + `Verlopen`), **application** (`ExpireRegistrationWorker`), **infra** (`RegistratieVerlopenProcessor`/`Pump`, `IRegistratieVerlopenClient`, Flowable acquire/complete + `CompleteDocumentWaitAsync`) — the timeout counterpart to the OpenZaak/escalation worker trios; idempotent per §8.6. - **BPMN** verified live against a `flowable-rest` probe: complete `WachtOpDocumenten` → routes to Beoordelen; fire the P30D timer → `RegistratieVerlopen` job (carrying `registrationId`) + the wait task cancelled. `verify-domain` exercises both branches in-stack (completes the wait in every existing block; fires the timer and asserts `Verlopen` in a new block). - **Scope boundary:** on expiry the aggregate goes `Verlopen` and the process ends, but the ZGW *zaak* is not yet set to a cancellation status — that needs a new ACL method + statustype seeding and is folded into S-10b (noted in ADR-0017). - `CompleteDocumentWaitAsync` is built and HTTP-tested here but not yet called from a domain endpoint; S-10b wires the upload trigger to it. 🤖 Generated with [Claude Code](https://claude.com/claude-code) Reviewed-on: #105
This commit was merged in pull request #105.
This commit is contained in:
@@ -15,12 +15,14 @@ namespace Big.Infrastructure;
|
||||
/// The REST contract here is the one verified against a live flowable-rest engine (ADR-0009).
|
||||
/// </summary>
|
||||
public sealed class FlowableWorkflowClient(HttpClient http, FlowableOptions options)
|
||||
: IWorkflowClient, IExternalWorkerClient, IUserTaskClient, IBeoordelingEscalatieClient
|
||||
: IWorkflowClient, IExternalWorkerClient, IUserTaskClient, IBeoordelingEscalatieClient, IRegistratieVerlopenClient
|
||||
{
|
||||
private const string Topic = "OpenZaakAanmaken";
|
||||
private const string EscalatieTopic = "BeoordelingEscaleren";
|
||||
private const string VerlopenTopic = "RegistratieVerlopen";
|
||||
private const string ProcessDefinitionKey = "registratie";
|
||||
private const string BeoordelenTaskKey = "Beoordelen";
|
||||
private const string WachtOpDocumentenTaskKey = "WachtOpDocumenten";
|
||||
private const string BehandelaarGroup = "behandelaar";
|
||||
private const string TeamleadGroup = "teamlead";
|
||||
private const string RegistrationIdVariable = "registrationId";
|
||||
@@ -114,6 +116,25 @@ public sealed class FlowableWorkflowClient(HttpClient http, FlowableOptions opti
|
||||
response.EnsureSuccessStatusCode();
|
||||
}
|
||||
|
||||
public async Task CompleteDocumentWaitAsync(string processInstanceId, CancellationToken ct = default)
|
||||
{
|
||||
// Find the still-open WachtOpDocumenten task in this instance and complete it, so the process
|
||||
// leaves the 30-day wait and continues to beoordeling (S-10a, ADR-0017). If the instance is no
|
||||
// longer parked there (already continued, or the timer already cancelled it) this is a
|
||||
// best-effort no-op — mirroring the withdrawal/escalation correlation (§8.6).
|
||||
var query = new TaskByInstanceQueryRequest(processInstanceId, WachtOpDocumentenTaskKey);
|
||||
var page = await PostAsync<TaskByInstanceQueryRequest, TaskQueryResult>(
|
||||
"service/query/tasks", query, ct);
|
||||
|
||||
var task = page?.Data?.FirstOrDefault();
|
||||
if (task is null)
|
||||
return;
|
||||
|
||||
using var response = await SendAsync(
|
||||
$"service/runtime/tasks/{task.Id}", new CompleteTaskRequest("complete", []), ct);
|
||||
response.EnsureSuccessStatusCode();
|
||||
}
|
||||
|
||||
public async Task<IReadOnlyList<EscalatieJob>> AcquireBeoordelingEscalatieJobsAsync(int maxJobs, CancellationToken ct = default)
|
||||
{
|
||||
var request = new AcquireJobsRequest(EscalatieTopic, options.LockDuration, maxJobs, options.WorkerId);
|
||||
@@ -155,6 +176,23 @@ public sealed class FlowableWorkflowClient(HttpClient http, FlowableOptions opti
|
||||
response.EnsureSuccessStatusCode();
|
||||
}
|
||||
|
||||
public async Task<IReadOnlyList<RegistratieVerlopenJob>> AcquireRegistratieVerlopenJobsAsync(int maxJobs, CancellationToken ct = default)
|
||||
{
|
||||
var request = new AcquireJobsRequest(VerlopenTopic, options.LockDuration, maxJobs, options.WorkerId);
|
||||
|
||||
var jobs = await PostAsync<AcquireJobsRequest, List<AcquiredJob>>(
|
||||
"external-job-api/acquire/jobs", request, ct) ?? [];
|
||||
|
||||
return [.. jobs.Select(job => new RegistratieVerlopenJob(job.Id, RegistrationId.Parse(job.RegistrationId())))];
|
||||
}
|
||||
|
||||
public async Task CompleteRegistratieVerlopenJobAsync(string jobId, CancellationToken ct = default)
|
||||
{
|
||||
using var response = await SendAsync(
|
||||
$"external-job-api/acquire/jobs/{jobId}/complete", new CompleteJobRequest(options.WorkerId, []), ct);
|
||||
response.EnsureSuccessStatusCode();
|
||||
}
|
||||
|
||||
private async Task<TResponse?> GetAsync<TResponse>(string path, CancellationToken ct)
|
||||
{
|
||||
var message = new HttpRequestMessage(HttpMethod.Get, new Uri(options.BaseUrl, path));
|
||||
|
||||
@@ -36,3 +36,19 @@ public interface IBeoordelingEscalatieClient
|
||||
/// <summary>Complete an acquired escalation job so its token reaches the escalation end event.</summary>
|
||||
Task CompleteBeoordelingEscalatieJobAsync(string jobId, CancellationToken ct = default);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// The document-timeout side of the Workflow Client (S-10a): the <c>RegistratieVerlopen</c>
|
||||
/// external-worker jobs parked by the 30-day boundary timer on <c>WachtOpDocumenten</c>. Kept separate
|
||||
/// from the other worker ports (interface segregation) so neither the OpenZaak nor escalation worker
|
||||
/// sees expiry. Implemented by <see cref="FlowableWorkflowClient"/> — the only code that talks to
|
||||
/// Flowable (§8.2, ADR-0017).
|
||||
/// </summary>
|
||||
public interface IRegistratieVerlopenClient
|
||||
{
|
||||
/// <summary>Acquire and lock up to <paramref name="maxJobs"/> <c>RegistratieVerlopen</c> jobs.</summary>
|
||||
Task<IReadOnlyList<RegistratieVerlopenJob>> AcquireRegistratieVerlopenJobsAsync(int maxJobs, CancellationToken ct = default);
|
||||
|
||||
/// <summary>Complete an acquired expiry job so its token reaches the <c>endVerlopen</c> end event.</summary>
|
||||
Task CompleteRegistratieVerlopenJobAsync(string jobId, CancellationToken ct = default);
|
||||
}
|
||||
|
||||
@@ -0,0 +1,41 @@
|
||||
using Big.Application;
|
||||
using Microsoft.Extensions.Logging;
|
||||
|
||||
namespace Big.Infrastructure;
|
||||
|
||||
/// <summary>
|
||||
/// One poll tick of the document-timeout worker (S-10a, ADR-0017): acquire the parked
|
||||
/// <c>RegistratieVerlopen</c> jobs — the tokens the 30-day boundary timer on <c>WachtOpDocumenten</c>
|
||||
/// spawns — expire each correlated registration via the <see cref="ExpireRegistrationWorker"/>, and
|
||||
/// complete the job so its token reaches <c>endVerlopen</c>. A job that fails is logged and left
|
||||
/// un-completed so Flowable redelivers it (§8.6). Split out from the hosted pump so the
|
||||
/// acquire→expire→complete logic is unit-testable without a running host. Mirrors
|
||||
/// <see cref="OpenZaakJobProcessor"/> and <see cref="BeoordelingEscalatieProcessor"/>.
|
||||
/// </summary>
|
||||
public sealed class RegistratieVerlopenProcessor(
|
||||
IRegistratieVerlopenClient client,
|
||||
ExpireRegistrationWorker worker,
|
||||
ILogger<RegistratieVerlopenProcessor> logger)
|
||||
{
|
||||
/// <summary>Acquire and process up to <paramref name="maxJobs"/> jobs. Returns the number acquired.</summary>
|
||||
public async Task<int> PumpOnceAsync(int maxJobs, CancellationToken ct = default)
|
||||
{
|
||||
var jobs = await client.AcquireRegistratieVerlopenJobsAsync(maxJobs, ct);
|
||||
|
||||
foreach (var job in jobs)
|
||||
{
|
||||
try
|
||||
{
|
||||
await worker.HandleAsync(job, ct);
|
||||
await client.CompleteRegistratieVerlopenJobAsync(job.JobId, ct);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
// Leave the job un-completed: its lock expires and Flowable redelivers it (§8.6).
|
||||
logger.LogError(ex, "RegistratieVerlopen job {JobId} failed; leaving it for redelivery.", job.JobId);
|
||||
}
|
||||
}
|
||||
|
||||
return jobs.Count;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,49 @@
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Microsoft.Extensions.Hosting;
|
||||
using Microsoft.Extensions.Logging;
|
||||
|
||||
namespace Big.Infrastructure;
|
||||
|
||||
/// <summary>
|
||||
/// The hosted polling loop of the document-timeout worker (S-10a, ADR-0017): on an interval it
|
||||
/// resolves a scoped <see cref="RegistratieVerlopenProcessor"/> and asks it to drain the parked
|
||||
/// <c>RegistratieVerlopen</c> jobs. A deliberately thin shell — all acquire/expire/complete logic
|
||||
/// lives in the processor, which is unit-tested; this class only owns the timer, the per-tick scope,
|
||||
/// and loop resilience. Structurally identical to <see cref="BeoordelingEscalatiePump"/>.
|
||||
/// </summary>
|
||||
public sealed class RegistratieVerlopenPump(
|
||||
IServiceScopeFactory scopeFactory,
|
||||
FlowableOptions options,
|
||||
ILogger<RegistratieVerlopenPump> logger) : BackgroundService
|
||||
{
|
||||
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
|
||||
{
|
||||
while (!stoppingToken.IsCancellationRequested)
|
||||
{
|
||||
try
|
||||
{
|
||||
using var scope = scopeFactory.CreateScope();
|
||||
var processor = scope.ServiceProvider.GetRequiredService<RegistratieVerlopenProcessor>();
|
||||
await processor.PumpOnceAsync(options.MaxJobsPerPoll, stoppingToken);
|
||||
}
|
||||
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
|
||||
{
|
||||
break;
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
// A transient fault (e.g. Flowable briefly unreachable) must not kill the loop.
|
||||
logger.LogError(ex, "RegistratieVerlopen job poll failed; retrying after the poll interval.");
|
||||
}
|
||||
|
||||
try
|
||||
{
|
||||
await Task.Delay(options.PollInterval, stoppingToken);
|
||||
}
|
||||
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
|
||||
{
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user