diff --git a/BACKLOG.md b/BACKLOG.md index 41b1b38..db28395 100644 --- a/BACKLOG.md +++ b/BACKLOG.md @@ -199,9 +199,19 @@ _Split from the original S-09 — scoped to the portal only; the approval flow i ### S-10 · Document upload + boundary timer for document timeout (Flow 2) -**Outcome:** BPMN extended with a "wacht op documenten" user task with a 30-day boundary timer. Self-service portal supports diploma upload. On timeout the case is cancelled. +Split (issue #11 closed) into two independently-demoable slices per §13 — the original spanned six net-new surfaces including a new ZGW boundary: -**Acceptance:** BDD scenarios for both branches; integration tests for the timer firing. +#### S-10a · Document-wait task + 30-day timeout cancellation + provision trigger — #102 + +**Outcome:** BPMN gains a `WachtOpDocumenten` user task with a 30-day (P30D) interrupting boundary timer. On timeout the case is cancelled — the timer runs to a dedicated cancel end-event and the domain aggregate moves to a new terminal status `Verlopen` via an external-worker (mirrors S-14 escalation / S-11 withdrawal). "Documents received" is wired end-to-end (domain endpoint + BFF + a "Documenten aanleveren" button on the self-service page) so the walking-skeleton e2e stays green — but the document is **not yet stored** in ZGW; that is S-10b. + +**Acceptance:** BDD both branches (documents-in-time vs timeout-cancel); live timer-fire via the management-API "move" idiom; the registration e2e provides documents before the behandelaar step. + +#### S-10b · Real diploma upload stored via the ACL Documenten API — #103 + +**Outcome:** the self-service "Documenten aanleveren" action becomes a real file upload; the document is stored in the ZGW Documenten (DRC) API and related to the zaak, with all document calls routed through the ACL (§8.1), and the zaak is set to a cancellation status on timeout expiry. Builds on the S-10a trigger/wait. Depends on #102. + +**Acceptance:** ACL Documenten gateway integration test; Playwright e2e uploads a real document; the openbaar/zaak reflects the stored document. ### S-11 · Withdrawal (Flow 3) diff --git a/apps/self-service/src/app/registration/registration-page.html b/apps/self-service/src/app/registration/registration-page.html index d4739e5..6910880 100644 --- a/apps/self-service/src/app/registration/registration-page.html +++ b/apps/self-service/src/app/registration/registration-page.html @@ -11,6 +11,24 @@

Uw registratie is ontvangen. Referentie: {{ reference() }}.

+ @if (documentsProvided()) { +

Uw documenten zijn aangeleverd.

+ } @else { + @if (provideDocumentsFailed()) { +

+ Het aanleveren van uw documenten is niet gelukt. Probeer het opnieuw. +

+ } + + } @if (withdrawFailed()) {

Het intrekken van uw registratie is niet gelukt. Probeer het opnieuw. diff --git a/apps/self-service/src/app/registration/registration-page.spec.ts b/apps/self-service/src/app/registration/registration-page.spec.ts index 1119762..f17ff06 100644 --- a/apps/self-service/src/app/registration/registration-page.spec.ts +++ b/apps/self-service/src/app/registration/registration-page.spec.ts @@ -20,10 +20,12 @@ class FakeAuth extends AuthService { function providers( post = vi.fn().mockReturnValue(of({ registrationId: 'reg-9', status: 'Ingediend' })), withdraw = vi.fn().mockReturnValue(of(undefined)), + provideDocuments = vi.fn().mockReturnValue(of(undefined)), ) { return { post, withdraw, + provideDocuments, providers: [ { provide: AuthService, useClass: FakeAuth }, { @@ -31,6 +33,7 @@ function providers( useValue: { postSelfServiceRegistrations: post, postSelfServiceRegistrationsIdWithdraw: withdraw, + postSelfServiceRegistrationsIdDocuments: provideDocuments, }, }, ], @@ -80,6 +83,37 @@ describe('RegistrationPage', () => { expect(await screen.findByText(/ingetrokken/i)).toBeTruthy(); }); + it('offers to provide documents after submitting, and doing so confirms', async () => { + const { provideDocuments, providers: p } = providers(); + await render(RegistrationPage, { providers: p }); + + fireEvent.click(screen.getByRole('button', { name: /indienen/i })); + await screen.findByText(/ontvangen/i); + + fireEvent.click(await screen.findByRole('button', { name: /documenten aanleveren/i })); + + // The provide-documents call is keyed by the reference the submit returned, and the page confirms. + expect(provideDocuments).toHaveBeenCalledWith('reg-9'); + expect(await screen.findByText(/documenten.*aangeleverd/i)).toBeTruthy(); + }); + + it('surfaces a provide-documents failure and keeps the action available', async () => { + const { providers: p } = providers( + vi.fn().mockReturnValue(of({ registrationId: 'reg-9', status: 'Ingediend' })), + vi.fn().mockReturnValue(of(undefined)), + vi.fn().mockReturnValue(throwError(() => new Error('documents rejected'))), + ); + await render(RegistrationPage, { providers: p }); + + fireEvent.click(screen.getByRole('button', { name: /indienen/i })); + await screen.findByText(/ontvangen/i); + fireEvent.click(await screen.findByRole('button', { name: /documenten aanleveren/i })); + + expect(await screen.findByRole('alert')).toBeTruthy(); + expect(screen.queryByText(/aangeleverd/i)).toBeNull(); + expect(screen.getByRole('button', { name: /documenten aanleveren/i })).toBeTruthy(); + }); + it('surfaces a withdraw failure and keeps the action available', async () => { const { providers: p } = providers( vi.fn().mockReturnValue(of({ registrationId: 'reg-9', status: 'Ingediend' })), diff --git a/apps/self-service/src/app/registration/registration-page.ts b/apps/self-service/src/app/registration/registration-page.ts index a5380de..8c93c48 100644 --- a/apps/self-service/src/app/registration/registration-page.ts +++ b/apps/self-service/src/app/registration/registration-page.ts @@ -26,6 +26,9 @@ export class RegistrationPage { protected readonly withdrawing = signal(false); protected readonly withdrawn = signal(false); protected readonly withdrawFailed = signal(false); + protected readonly providingDocuments = signal(false); + protected readonly documentsProvided = signal(false); + protected readonly provideDocumentsFailed = signal(false); submit(): void { this.submitting.set(true); @@ -44,6 +47,26 @@ export class RegistrationPage { }); } + provideDocuments(): void { + const reference = this.reference(); + if (!reference) { + return; + } + this.providingDocuments.set(true); + this.provideDocumentsFailed.set(false); + this.bff.postSelfServiceRegistrationsIdDocuments(reference).subscribe({ + next: () => { + this.documentsProvided.set(true); + this.providingDocuments.set(false); + }, + // Surface the failure instead of swallowing it: keep the action so the user can retry. + error: () => { + this.provideDocumentsFailed.set(true); + this.providingDocuments.set(false); + }, + }); + } + withdraw(): void { const reference = this.reference(); if (!reference) { diff --git a/docs/architecture/adr-0017-document-wait-timeout-cancellation.md b/docs/architecture/adr-0017-document-wait-timeout-cancellation.md new file mode 100644 index 0000000..78ce9d6 --- /dev/null +++ b/docs/architecture/adr-0017-document-wait-timeout-cancellation.md @@ -0,0 +1,90 @@ +# ADR-0017: A document-wait task with a 30-day interrupting timer cancels the registration + +- **Status:** Accepted +- **Date:** 2026-07-20 +- **Deciders:** Respellion engineering +- **Relates to:** S-10a (#102); proposal #104; split from S-10 (#11). Builds on ADR-0009 (external-task + worker / Workflow Client), ADR-0014 (withdrawal cancels the process), ADR-0015 (beoordeling + escalation — the boundary-timer + external-worker pattern), ADR-0016 (diploma-eligibility DMN). + +## Context + +Flow 2 (PRD §5) requires the citizen to supply documents (their diploma) after submitting. The +registratie process must park waiting for those documents and, if they do not arrive within 30 days, +cancel the case. S-10 was split (§13): **S-10a** is this workflow/timeout spine (backend only); +**S-10b** wires the actual upload (portal → BFF → domain → ACL → Documenten API) that completes the +wait. This ADR records the spine: where the wait sits, how the timeout cancels, and how the domain +aggregate stays in sync. + +## Decision + +**A `WachtOpDocumenten` user task is inserted immediately after `OpenZaakAanmaken`, carrying an +`cancelActivity="true"` (interrupting) `P30D` boundary timer. "Documents received" completes the task +and the process continues into the diploma-eligibility routing; on timeout the timer cancels the task, +runs a `RegistratieVerlopen` external-worker task, and ends the process at `endVerlopen`. A domain +worker expires the correlated aggregate to a new terminal status `Verlopen`.** + +- **Where the wait sits.** Right after the zaak is opened, before the diploma-eligibility DMN: the zaak + exists, then the process waits for documents; on receipt it continues to the DMN routing → Beoordelen + (ADR-0016). The wait gates the whole assessment, so it precedes the routing rather than sitting + between the gateway and Beoordelen. +- **Interrupting timer, mirroring the existing constructs.** Unlike the S-14 escalation timer + (non-interrupting — the Beoordelen task stays open), this timer is interrupting: when it fires the + wait token is consumed and the case is cancelled, like the S-11 withdrawal boundary (ADR-0014). The + timeout branch runs a `RegistratieVerlopen` external-worker task (topic mirrors + `OpenZaakAanmaken`/`BeoordelingEscaleren`) → `endVerlopen`. +- **The domain stays authoritative.** The `RegistratieVerlopen` job carries the `registrationId`; the + `RegistratieVerlopenProcessor` drains it and the `ExpireRegistrationWorker` loads the aggregate and + calls `Registration.Expire()`, moving it to the new terminal status `Verlopen`. This keeps the + aggregate — which the projection/openbaar view reads — the source of truth, exactly as escalation and + withdrawal do. Idempotent per §8.6: a redelivered job whose aggregate is already `Verlopen` completes + without persisting again; an unknown registration throws so the job is redelivered. +- **Documents-in-time transition.** `IWorkflowClient.CompleteDocumentWaitAsync(processInstanceId)` + completes the `WachtOpDocumenten` task (the Workflow Client remains the only code that talks to + Flowable, §8.2). It is best-effort — a no-op if the instance already left the wait (continued, or + timed out). The trigger is wired end-to-end in S-10a: a `ProvideDocuments` application use case behind + an owner-scoped domain endpoint `POST /registrations/{id}/documents`, a BFF passthrough + `POST /self-service/registrations/{id}/documents` (bsn from the DigiD token), and a "Documenten + aanleveren" action on the self-service page — so the walking-skeleton e2e stays green (a registration + can still reach the behandelaar). **S-10b replaces the stub trigger with a real file upload stored in + the ZGW Documenten (DRC) API via the ACL**; the completion of the wait is unchanged. + - *Why the trigger lives here, not in S-10b:* inserting the `WachtOpDocumenten` gate without any way + to pass it breaks the submit→beoordeling e2e (a merge gate). Splitting "gate" from "means to pass + the gate" across slices would leave `main` red, so S-10a owns both; S-10b is purely the ZGW storage + behind the same action. + +## Consequences + +**Positive** + +- The wait/timeout is a first-class workflow construct that reuses the boundary-timer + external-worker + pattern already proven by S-14, so the domain change is small and additive: one terminal status, one + worker trio (worker + processor + pump), one Workflow Client method. +- §8 stays clean: the Workflow Client is still the only Flowable caller, and no new ZGW boundary is + introduced in S-10a. +- The timeout is verified live (verify-domain fires the P30D timer via the management-API "move" idiom + and asserts the domain reaches `Verlopen`), consistent with ADR-0009/0014/0015. + +**Negative / costs** + +- Every registration now parks at `WachtOpDocumenten` before Beoordelen, so the other flows must supply + documents first: the live-check blocks (S-11/S-12b/S-13/S-14) complete the task via Flowable, and the + registration e2e clicks "Documenten aanleveren". A small, explicit step, but it touches every path + through the process. +- On expiry S-10a cancels the *process* and marks the aggregate `Verlopen` but does **not** set the ZGW + *zaak* to a cancellation status — that needs a new ACL method + statustype seeding, which overlaps + S-10b's ACL/infra work. Deferred to S-10b (or a follow-up); noted here as the S-10a/S-10b boundary. +- Withdrawing while parked at `WachtOpDocumenten` marks the aggregate `Ingetrokken` but does not cancel + the process (the withdrawal message boundary is on `Beoordelen`); the timeout worker tolerates this + by no-op'ing on an already-resolved aggregate. Extending withdrawal to the wait state is a follow-up. + +## Alternatives considered + +- **Pure-BPMN cancellation (timer → end event, no worker).** Rejected: the domain aggregate would then + be out of sync with the cancelled process, and the openbaar/projection view reads the aggregate's + status — the case would still look open. +- **Wait task between the gateway and Beoordelen.** Rejected: documents gate the whole assessment + (including the CBGV-advies routing), so the wait belongs before the DMN, not after it. +- **A dedicated timeout status per branch vs. reusing an open-state guard.** `Expire()` reuses the same + `RequireOpenForDecision` guard as withdrawal/decision, so only an `INGEDIEND`/`IN_BEHANDELING` + registration can lapse and the terminal states stay mutually exclusive — no new guard logic. diff --git a/docs/demo-script.md b/docs/demo-script.md index 436fbcb..7b48e0c 100644 --- a/docs/demo-script.md +++ b/docs/demo-script.md @@ -361,7 +361,8 @@ DOM=http://localhost:8080 # domain service curl -s -i -X POST "$DOM/registrations" -H 'Content-Type: application/json' \ -d '{"bsn":"123456782","diplomaOrigin":"Buitenlands"}' | grep -i '^location:' # -# 2. Once the zaak is opened, the process parks at the CBGV-advies task (NOT Beoordelen). In Flowable: +# 2. Once the zaak is opened, the process first parks at WachtOpDocumenten (S-10a); complete that task +# (documents received) — then it parks at the CBGV-advies task (NOT Beoordelen). In Flowable: FL=http://localhost:8090/flowable-rest/service curl -s -u rest-admin:test -X POST "$FL/query/tasks" -H 'Content-Type: application/json' \ -d '{"processDefinitionKey":"registratie","taskDefinitionKey":"CBGVAdvies"}' | python3 -m json.tool @@ -382,3 +383,54 @@ domestic: `Beoordelen` directly (§8.2, ADR-0016). > The domestic/foreign paths are covered by the `Een diploma op herkomst routeren` acceptance > scenarios and unit tests (the origin is carried into the process); the DMN decision and the > foreign→CBGV routing are asserted live by the verify-domain check. + +## S-10a — Document wait + 30-day timeout cancels the registration (#102, ADR-0017) + +After the zaak is opened the registratie process parks at a **WachtOpDocumenten** user task, waiting +for the citizen's documents (their diploma). Two things can happen: + +- **Documents arrive in time** → the task completes and the process continues to the diploma-eligibility + routing (S-13) → beoordeling. +- **30 days pass with no documents** → an interrupting `P30D` boundary timer cancels the wait, runs the + `RegistratieVerlopen` external task, and the domain expires the registration to the terminal status + **VERLOPEN** (the case is cancelled). + +The "documents received" trigger is wired end-to-end in S-10a: the self-service page shows a +**"Documenten aanleveren"** button after submit (portal → BFF → domain → completes the wait). S-10b +turns that into a real file upload stored in the ZGW Documenten API via the ACL. The timeout branch is +demonstrated by firing the 30-day timer early via the management API. + +```bash +DOM=http://localhost:8080 # domain service +FL=http://localhost:8090/flowable-rest/service # flowable-rest + +# 1. Submit a registration; once the zaak is opened it parks at WachtOpDocumenten: +curl -s -i -X POST "$DOM/registrations" -H 'Content-Type: application/json' \ + -d '{"bsn":"123456782"}' | grep -i '^location:' # note the /registrations/ reference +WQ='{"processDefinitionKey":"registratie","taskDefinitionKey":"WachtOpDocumenten"}' + +# 2a. Documents-in-time: complete the WachtOpDocumenten task → the process advances to beoordeling. +TID=$(curl -s -u rest-admin:test -X POST "$FL/query/tasks" -H 'Content-Type: application/json' \ + -d "$WQ" | python3 -c 'import sys,json;print(json.load(sys.stdin)["data"][0]["id"])') +curl -s -u rest-admin:test -X POST "$FL/runtime/tasks/$TID" \ + -H 'Content-Type: application/json' -d '{"action":"complete"}' + +# 2b. Timeout: instead of completing it, fire the 30-day timer early via the management API. Find the +# instance's timer job, "move" it to executable; the async executor fires the interrupting event. +PID=$(curl -s -u rest-admin:test -X POST "$FL/query/tasks" -H 'Content-Type: application/json' \ + -d "$WQ" | python3 -c 'import sys,json;print(json.load(sys.stdin)["data"][0]["processInstanceId"])') +JID=$(curl -s -u rest-admin:test "$FL/management/timer-jobs?processInstanceId=$PID" \ + | python3 -c 'import sys,json;print(json.load(sys.stdin)["data"][0]["id"])') +curl -s -u rest-admin:test -X POST "$FL/management/timer-jobs/$JID" \ + -H 'Content-Type: application/json' -d '{"action":"move"}' +# The RegistratieVerlopen worker then expires the aggregate — read it back as VERLOPEN: +curl -s "$DOM/registrations/" # → {"status":"Verlopen", ...} +``` + +**The path:** registratie process parks at `WachtOpDocumenten` → documents received completes it (→ +routing → `Beoordelen`), OR the `P30D` interrupting timer fires → `RegistratieVerlopen` external task +→ domain worker expires the aggregate to `Verlopen` → `endVerlopen` (§8.2, ADR-0017). + +> Both branches are covered by the `Een documenttermijn laten verlopen` acceptance scenarios (worker + +> aggregate) and unit tests; the wait completion and the 30-day timer firing are asserted live by the +> verify-domain check. diff --git a/infra/run-domain-check.sh b/infra/run-domain-check.sh index 814683d..7294c54 100755 --- a/infra/run-domain-check.sh +++ b/infra/run-domain-check.sh @@ -93,6 +93,27 @@ print(next((t['id'] for t in (d.get('data') or []) flcurl() { docker run --rm --network "$net" curlimages/curl:latest -fsS -u rest-admin:test "$@"; } query='{"processDefinitionKey":"registratie","taskDefinitionKey":"Beoordelen","includeProcessVariables":true}' +wacht_query='{"processDefinitionKey":"registratie","taskDefinitionKey":"WachtOpDocumenten","includeProcessVariables":true}' + +# S-10a: every registration now parks at WachtOpDocumenten first (interrupting P30D timer). Completing +# that task stands in for the citizen's document upload (wired for real in S-10b), letting the process +# advance to the diploma routing / Beoordelen so the checks below still hold. The 30-day timeout branch +# is exercised separately at the end. +complete_wacht() { # reg_id + local rid="$1" wid="" r + for _ in $(seq 1 30); do + r="$(flcurl -X POST "$fl_base/query/tasks" -H 'Content-Type: application/json' -d "$wacht_query" 2>/dev/null || true)" + wid="$(printf '%s' "$r" | task_for_reg "$rid")" + [ -n "$wid" ] && break + sleep 2 + done + [ -n "$wid" ] || { echo "FAIL — no WachtOpDocumenten task appeared for $rid" >&2; docker logs "$dom" 2>&1 | tail -15 >&2; exit 1; } + flcurl -X POST "$fl_base/runtime/tasks/$wid" -H 'Content-Type: application/json' -d '{"action":"complete"}' >/dev/null + echo ">> completed WachtOpDocumenten for $rid (documents received)" +} + +echo ">> completing WachtOpDocumenten so the process advances (documents received)" +complete_wacht "$reg_id" echo ">> polling Flowable for the Beoordelen user task (werkbak)" task_id="" @@ -130,6 +151,7 @@ loc2="$(docker run --rm --network "$net" curlimages/curl:latest \ [ -n "$loc2" ] || { echo "FAIL — second POST /registrations returned no Location" >&2; exit 1; } reg_id2="${loc2##*/}" echo ">> second registration $reg_id2" +complete_wacht "$reg_id2" echo ">> polling Flowable for its Beoordelen task" task_id2="" @@ -171,6 +193,7 @@ locf="$(docker run --rm --network "$net" curlimages/curl:latest \ [ -n "$locf" ] || { echo "FAIL — foreign POST /registrations returned no Location" >&2; exit 1; } reg_idf="${locf##*/}" echo ">> foreign registration $reg_idf" +complete_wacht "$reg_idf" echo ">> polling Flowable for its CBGV-advies task (foreign diplomas route here first)" cbgv_task="" @@ -241,6 +264,8 @@ except Exception: d={} print(((d.get('data') or [{}])[0]).get('id',''))"; } +complete_wacht "$reg_id3" + echo ">> polling Flowable for its Beoordelen task" task_id3=""; pid3="" for _ in $(seq 1 30); do @@ -278,4 +303,50 @@ for _ in $(seq 1 30); do done [ -n "$escalated" ] || { echo "FAIL — Beoordelen task not reassigned to teamlead (candidate groups: '$groups')" >&2; docker logs "$dom" 2>&1 | tail -15 >&2; exit 1; } echo "OK — the 14-day timer escalated the still-open Beoordelen task to the teamlead" + +# ── S-10a: document timeout. A registration parks at WachtOpDocumenten and — unlike every block above — +# its documents never arrive. We fire its 30-day boundary timer early via the management API; the +# INTERRUPTING timer cancels the wait and routes a token to the RegistratieVerlopen external task. The +# domain's timeout worker acquires it and expires the registration to VERLOPEN (ADR-0017). ──────────── +echo ">> submitting a registration to let its document term lapse" +locv="$(docker run --rm --network "$net" curlimages/curl:latest \ + -fsS -D - -o /dev/null -X POST "http://$dom_ip:8080/registrations" \ + -H 'Content-Type: application/json' -d '{"bsn":"123456782"}' \ + | sed -n 's/\r$//; s/^[Ll]ocation: //p' | head -1)" +[ -n "$locv" ] || { echo "FAIL — timeout POST /registrations returned no Location" >&2; exit 1; } +reg_idv="${locv##*/}" +echo ">> timeout registration $reg_idv" + +echo ">> polling Flowable for its WachtOpDocumenten task" +wacht_id=""; pidv="" +for _ in $(seq 1 30); do + resp="$(flcurl -X POST "$fl_base/query/tasks" -H 'Content-Type: application/json' -d "$wacht_query" 2>/dev/null || true)" + read -r wacht_id pidv <<<"$(printf '%s' "$resp" | task_and_pid_for_reg "$reg_idv")" + [ -n "$wacht_id" ] && break + sleep 2 +done +[ -n "$wacht_id" ] || { echo "FAIL — no WachtOpDocumenten task appeared for $reg_idv" >&2; docker logs "$dom" 2>&1 | tail -15 >&2; exit 1; } +echo ">> WachtOpDocumenten task $wacht_id (instance $pidv) is waiting for documents" + +echo ">> firing the 30-day document timer early via the management API" +timer_idv="$(flcurl "$fl_base/management/timer-jobs?processInstanceId=$pidv" | first_job_id)" +[ -n "$timer_idv" ] || { echo "FAIL — no timer job found for instance $pidv" >&2; exit 1; } +# Move the timer job to an executable async job; the async executor fires the interrupting boundary +# event. It may run before we look, so executing it explicitly is a best-effort nudge (as for S-14). +flcurl -X POST "$fl_base/management/timer-jobs/$timer_idv" -H 'Content-Type: application/json' -d '{"action":"move"}' >/dev/null +async_idv="$(flcurl "$fl_base/management/jobs?processInstanceId=$pidv" 2>/dev/null | first_job_id || true)" +if [ -n "$async_idv" ]; then + flcurl -X POST "$fl_base/management/jobs/$async_idv" -H 'Content-Type: application/json' -d '{"action":"execute"}' >/dev/null 2>&1 || true +fi +echo ">> timer fired; the RegistratieVerlopen token is parked for the domain worker" + +echo ">> polling the domain until the timeout worker expires the registration to VERLOPEN" +verlopen="" +for _ in $(seq 1 30); do + body="$(docker run --rm --network "$net" curlimages/curl:latest -fsS "http://$dom_ip:8080$locv" 2>/dev/null || true)" + printf '%s' "$body" | grep -qi 'verlopen' && { verlopen=1; break; } + sleep 2 +done +[ -n "$verlopen" ] || { echo "FAIL — registration $reg_idv not VERLOPEN after the document timer fired (body: $body)" >&2; docker logs "$dom" 2>&1 | tail -15 >&2; exit 1; } +echo "OK — the 30-day document timer expired the registration to VERLOPEN" exit 0 diff --git a/libs/api-client/src/lib/generated/bff-api.ts b/libs/api-client/src/lib/generated/bff-api.ts index d227c97..f5d35d7 100644 --- a/libs/api-client/src/lib/generated/bff-api.ts +++ b/libs/api-client/src/lib/generated/bff-api.ts @@ -226,6 +226,40 @@ export class BffApiV1Service { ); } + postSelfServiceRegistrationsIdDocuments(id: string, options?: HttpClientBodyOptions): Observable; + postSelfServiceRegistrationsIdDocuments(id: string, options?: HttpClientEventOptions): Observable>; + postSelfServiceRegistrationsIdDocuments(id: string, options?: HttpClientResponseOptions): Observable>; + postSelfServiceRegistrationsIdDocuments( + id: string, options?: HttpClientObserveOptions): Observable | AngularHttpResponse> { + if (options?.observe === 'events') { + return this.http.post( + `/self-service/registrations/${id}/documents`, + undefined,{ + ...(options as Omit, 'observe'>), + observe: 'events', + } + ); + } + + if (options?.observe === 'response') { + return this.http.post( + `/self-service/registrations/${id}/documents`, + undefined,{ + ...(options as Omit, 'observe'>), + observe: 'response', + } + ); + } + + return this.http.post( + `/self-service/registrations/${id}/documents`, + undefined,{ + ...(options as Omit, 'observe'>), + observe: 'body', + } + ); + } + getOpenbaarRegister(params?: GetOpenbaarRegisterParams, options?: HttpClientBodyOptions): Observable; getOpenbaarRegister(params?: GetOpenbaarRegisterParams, options?: HttpClientEventOptions): Observable>; getOpenbaarRegister(params?: GetOpenbaarRegisterParams, options?: HttpClientResponseOptions): Observable>; diff --git a/services/bff/Bff.Api/DownstreamClients.cs b/services/bff/Bff.Api/DownstreamClients.cs index 30b2429..f935247 100644 --- a/services/bff/Bff.Api/DownstreamClients.cs +++ b/services/bff/Bff.Api/DownstreamClients.cs @@ -27,6 +27,11 @@ public interface IDomainClient /// unknown or not the caller's (404), so the BFF can relay a 404 rather than a 500. Task WithdrawRegistrationAsync(string registrationId, string bsn, CancellationToken ct = default); + ///

Provide the documents the caller's own registration is waiting for ("documenten + /// aanleveren"). Owner-scoped by . Returns false when the domain + /// reports the registration is unknown or not the caller's (404), so the BFF can relay a 404. + Task ProvideDocumentsAsync(string registrationId, string bsn, CancellationToken ct = default); + /// The behandelaar's werkbak — registrations awaiting beoordeling. Task> GetWerkbakAsync(CancellationToken ct = default); @@ -63,6 +68,17 @@ public sealed class DomainClient(HttpClient http) : IDomainClient return true; } + public async Task ProvideDocumentsAsync(string registrationId, string bsn, CancellationToken ct = default) + { + using var response = await http.PostAsJsonAsync( + $"registrations/{registrationId}/documents", new { bsn }, ct); + // The domain 404s an unknown or not-owned registration; relay that rather than fail hard. + if (response.StatusCode == System.Net.HttpStatusCode.NotFound) + return false; + response.EnsureSuccessStatusCode(); + return true; + } + public async Task> GetWerkbakAsync(CancellationToken ct = default) => await http.GetFromJsonAsync>("behandel/werkbak", ct) ?? []; diff --git a/services/bff/Bff.Api/Program.cs b/services/bff/Bff.Api/Program.cs index c4683ce..940d5f2 100644 --- a/services/bff/Bff.Api/Program.cs +++ b/services/bff/Bff.Api/Program.cs @@ -104,6 +104,26 @@ app.MapPost("/self-service/registrations/{id}/withdraw", async (string id, Claim .Produces(StatusCodes.Status401Unauthorized) .Produces(StatusCodes.Status404NotFound); +// Self-service provide-documents (S-10a): the signed-in zorgprofessional supplies the documents their +// registration is waiting for ("documenten aanleveren"). The bsn comes from the DigiD token and is +// forwarded to the domain, which owner-scopes the action and completes the WachtOpDocumenten task; a +// registration that is unknown or not the caller's comes back 404. The real file upload + ZGW storage +// is S-10b — this is the trigger that unblocks the process. +app.MapPost("/self-service/registrations/{id}/documents", async (string id, ClaimsPrincipal user, IDomainClient domain, CancellationToken ct) => +{ + var bsn = user.FindFirstValue("bsn"); + if (string.IsNullOrWhiteSpace(bsn)) + return Results.BadRequest("The token carries no bsn claim."); + + var provided = await domain.ProvideDocumentsAsync(id, bsn, ct); + return provided ? Results.NoContent() : Results.NotFound(); +}) + .RequireAuthorization() + .Produces(StatusCodes.Status204NoContent) + .Produces(StatusCodes.Status400BadRequest) + .Produces(StatusCodes.Status401Unauthorized) + .Produces(StatusCodes.Status404NotFound); + // Openbaar register: an anonymous public lookup that exposes only public-safe fields (S-09). app.MapGet("/openbaar/register", async (string? q, IProjectionClient projection, CancellationToken ct) => { diff --git a/services/bff/Bff.Tests/BffFactory.cs b/services/bff/Bff.Tests/BffFactory.cs index f108da8..492ab05 100644 --- a/services/bff/Bff.Tests/BffFactory.cs +++ b/services/bff/Bff.Tests/BffFactory.cs @@ -94,6 +94,18 @@ internal sealed class FakeDomainClient : IDomainClient return Task.FromResult(WithdrawSucceeds); } + public (string RegistrationId, string Bsn)? DocumentsProvidedFor { get; private set; } + + /// Whether the fake domain reports the provide-documents as done (true → 204) or + /// not-found/not-owned (false → 404). Tests set this to exercise the relay. + public bool ProvideDocumentsSucceeds { get; set; } = true; + + public Task ProvideDocumentsAsync(string registrationId, string bsn, CancellationToken ct = default) + { + DocumentsProvidedFor = (registrationId, bsn); + return Task.FromResult(ProvideDocumentsSucceeds); + } + public (string RegistrationId, string Besluit)? Decided { get; private set; } public Task> GetWerkbakAsync(CancellationToken ct = default) diff --git a/services/bff/Bff.Tests/SelfServiceEndpointTests.cs b/services/bff/Bff.Tests/SelfServiceEndpointTests.cs index f7adaa9..c8e0ccd 100644 --- a/services/bff/Bff.Tests/SelfServiceEndpointTests.cs +++ b/services/bff/Bff.Tests/SelfServiceEndpointTests.cs @@ -112,5 +112,46 @@ public class SelfServiceEndpointTests Assert.Equal(HttpStatusCode.NotFound, response.StatusCode); } + private static HttpRequestMessage ProvideDocuments(string? bearer, string id = "reg-123") + { + var request = new HttpRequestMessage(HttpMethod.Post, $"/self-service/registrations/{id}/documents"); + if (bearer is not null) + request.Headers.Authorization = new AuthenticationHeaderValue("Bearer", bearer); + return request; + } + + [Fact] + public async Task Rejects_providing_documents_without_a_token() + { + using var factory = new BffFactory(); + + var response = await factory.CreateClient().SendAsync(ProvideDocuments(bearer: null)); + + Assert.Equal(HttpStatusCode.Unauthorized, response.StatusCode); + Assert.Null(factory.Domain.DocumentsProvidedFor); + } + + [Fact] + public async Task Provides_documents_for_the_callers_registration_forwarding_the_id_and_bsn() + { + using var factory = new BffFactory(); + + var response = await factory.CreateClient().SendAsync(ProvideDocuments(TestTokens.Valid("123456782"), "reg-9")); + + Assert.Equal(HttpStatusCode.NoContent, response.StatusCode); + Assert.Equal(("reg-9", "123456782"), factory.Domain.DocumentsProvidedFor); + } + + [Fact] + public async Task Relays_not_found_providing_documents_for_an_unknown_or_not_owned_registration() + { + using var factory = new BffFactory(); + factory.Domain.ProvideDocumentsSucceeds = false; + + var response = await factory.CreateClient().SendAsync(ProvideDocuments(TestTokens.Valid("123456782"))); + + Assert.Equal(HttpStatusCode.NotFound, response.StatusCode); + } + private sealed record SubmitAcceptedDto(string RegistrationId, string Status); } diff --git a/services/bff/openapi.json b/services/bff/openapi.json index 989a7a1..fd2101a 100644 --- a/services/bff/openapi.json +++ b/services/bff/openapi.json @@ -61,6 +61,37 @@ } } }, + "/self-service/registrations/{id}/documents": { + "post": { + "tags": [ + "Bff.Api" + ], + "parameters": [ + { + "name": "id", + "in": "path", + "required": true, + "schema": { + "type": "string" + } + } + ], + "responses": { + "204": { + "description": "No Content" + }, + "400": { + "description": "Bad Request" + }, + "401": { + "description": "Unauthorized" + }, + "404": { + "description": "Not Found" + } + } + } + }, "/openbaar/register": { "get": { "tags": [ diff --git a/services/domain/Big.Api/Program.cs b/services/domain/Big.Api/Program.cs index ac83c22..c78e79b 100644 --- a/services/domain/Big.Api/Program.cs +++ b/services/domain/Big.Api/Program.cs @@ -22,22 +22,29 @@ builder.Services.AddTransient(sp => sp.GetRequiredService(sp => sp.GetRequiredService()); builder.Services.AddTransient(sp => sp.GetRequiredService()); builder.Services.AddTransient(sp => sp.GetRequiredService()); +builder.Services.AddTransient(sp => sp.GetRequiredService()); builder.Services.AddHttpClient(); builder.Services.AddScoped(); builder.Services.AddScoped(); builder.Services.AddScoped(); builder.Services.AddScoped(); +builder.Services.AddScoped(); builder.Services.AddScoped(); builder.Services.AddScoped(); builder.Services.AddScoped(); builder.Services.AddScoped(); +builder.Services.AddScoped(); +builder.Services.AddScoped(); // The hosted external-task job worker polls Flowable and drives OpenZaakAanmaken to completion. builder.Services.AddHostedService(); // The escalation worker polls the BeoordelingEscaleren jobs the 14-day timer parks and reassigns // each overdue beoordeling to the teamlead (S-14). builder.Services.AddHostedService(); +// The document-timeout worker polls the RegistratieVerlopen jobs the 30-day timer on WachtOpDocumenten +// parks and expires each lapsed registration to VERLOPEN (S-10a, ADR-0017). +builder.Services.AddHostedService(); var app = builder.Build(); @@ -101,6 +108,23 @@ app.MapPost("/registrations/{id}/withdraw", async (string id, WithdrawRequest bo return outcome == WithdrawOutcome.Withdrawn ? Results.NoContent() : Results.NotFound(); }); +// Provide documents (S-10a): the zorgprofessional supplies the documents their registration is parked +// waiting for, completing the WachtOpDocumenten task so the process advances to beoordeling (ADR-0017). +// Owner-scoped by the caller's bsn (the BFF forwards it from the DigiD token); unknown or not-the- +// caller's is 404 (indistinguishable). Idempotent — completing an already-left wait is a no-op. The +// real file upload + ZGW storage is S-10b; this endpoint is the trigger that unblocks the process. +app.MapPost("/registrations/{id}/documents", async (string id, ProvideDocumentsRequest body, ProvideDocuments provide, CancellationToken ct) => +{ + if (!Guid.TryParse(id, out var guid)) + return Results.NotFound(); + + if (string.IsNullOrWhiteSpace(body?.Bsn)) + return Results.BadRequest(new { error = "A bsn is required to provide documents." }); + + var outcome = await provide.HandleAsync(new ProvideDocumentsCommand(new RegistrationId(guid), body.Bsn), ct); + return outcome == ProvideDocumentsOutcome.Accepted ? Results.NoContent() : Results.NotFound(); +}); + // The behandelaar's werkbak (S-12): the registrations awaiting beoordeling, read from the open // Beoordelen user tasks (§8.2) and enriched with bsn + status. The BFF proxies this behind // medewerker-realm + behandelaar-role authorization; the domain trusts its callers (§8.3). @@ -128,6 +152,8 @@ public sealed record DecideRequest(string Besluit); public sealed record WithdrawRequest(string Bsn); +public sealed record ProvideDocumentsRequest(string Bsn); + public sealed record RegistrationResponse(string RegistrationId, string Status, string? ZaakUrl); public partial class Program; diff --git a/services/domain/Big.Application/ExpireRegistrationWorker.cs b/services/domain/Big.Application/ExpireRegistrationWorker.cs new file mode 100644 index 0000000..05f8b94 --- /dev/null +++ b/services/domain/Big.Application/ExpireRegistrationWorker.cs @@ -0,0 +1,37 @@ +using Big.Domain; + +namespace Big.Application; + +/// +/// Handles one acquired RegistratieVerlopen external-worker job (S-10a, ADR-0017): load the +/// registration the job correlates to and expire it to VERLOPEN — the 30-day document-wait timer fired +/// before the documents arrived, so the case is cancelled. Pure application logic over ports; it knows +/// nothing of Flowable. The polling loop that feeds it jobs lives in Infrastructure. Mirrors +/// . +/// +public sealed class ExpireRegistrationWorker(IRegistrationStore store) +{ + /// + /// Process the job. Idempotent and tolerant of races (§8.6, at-least-once delivery): a job whose + /// registration is already resolved — a redelivered expiry (VERLOPEN), or one withdrawn/decided + /// while it waited (INGETROKKEN/INGESCHREVEN/AFGEWEZEN) — is a no-op, so the job still completes + /// rather than throwing into a redelivery loop. Only a still-open registration is expired. An + /// unknown registration is an error: it throws, leaving the job un-completed for Flowable to redeliver. + /// + public async Task HandleAsync(RegistratieVerlopenJob job, CancellationToken ct = default) + { + ArgumentNullException.ThrowIfNull(job); + + var registration = await store.GetAsync(job.RegistrationId, ct) + ?? throw new InvalidOperationException( + $"No registration {job.RegistrationId} for RegistratieVerlopen job {job.JobId}."); + + // Only a still-open registration lapses; an already-resolved one (expired, or withdrawn/decided + // while it waited) is left untouched so the job can complete without violating the aggregate. + if (registration.Status is not (RegistrationStatus.Ingediend or RegistrationStatus.InBehandeling)) + return; + + registration.Expire(); + await store.SaveAsync(registration, ct); + } +} diff --git a/services/domain/Big.Application/Ports.cs b/services/domain/Big.Application/Ports.cs index b89ac82..af12888 100644 --- a/services/domain/Big.Application/Ports.cs +++ b/services/domain/Big.Application/Ports.cs @@ -25,6 +25,14 @@ public interface IWorkflowClient /// ended, or not yet parked) it is a no-op; the aggregate is INGETROKKEN regardless. /// Task WithdrawProcessAsync(string processInstanceId, CancellationToken ct = default); + + /// + /// Signal that the required documents have arrived (S-10a): complete the WachtOpDocumenten + /// user task in the instance so the process leaves the 30-day wait state and continues to + /// beoordeling (ADR-0017). Best-effort — if the instance is not parked at that task (already + /// continued, or timed out) it is a no-op. The upload trigger that calls this is wired in S-10b. + /// + Task CompleteDocumentWaitAsync(string processInstanceId, CancellationToken ct = default); } /// @@ -95,3 +103,11 @@ public sealed record OpenZaakJob(string JobId, RegistrationId RegistrationId); /// once the 14-day boundary timer fires (ADR-0015). /// public sealed record EscalatieJob(string JobId, string ProcessInstanceId); + +/// +/// An acquired RegistratieVerlopen job (S-10a): the Flowable job id and the registration id it +/// carries as a process variable. The 30-day boundary timer on WachtOpDocumenten spawns it when +/// the required documents were not supplied in time; expiring the correlated registration to VERLOPEN +/// cancels the case (ADR-0017). +/// +public sealed record RegistratieVerlopenJob(string JobId, RegistrationId RegistrationId); diff --git a/services/domain/Big.Application/ProvideDocuments.cs b/services/domain/Big.Application/ProvideDocuments.cs new file mode 100644 index 0000000..c5cdcce --- /dev/null +++ b/services/domain/Big.Application/ProvideDocuments.cs @@ -0,0 +1,47 @@ +using Big.Domain; + +namespace Big.Application; + +/// A zorgprofessional's signal that they have supplied the documents their registration is +/// waiting for ("documenten aanleveren"). is the authenticated caller (from the +/// DigiD token, forwarded by the BFF): only the registration's own bsn may provide its documents. +public sealed record ProvideDocumentsCommand(RegistrationId RegistrationId, string Bsn); + +/// The outcome of a provide-documents request. +public enum ProvideDocumentsOutcome +{ + /// The documents were accepted; the process's document wait was completed (if any). + Accepted, + + /// No registration with that id belongs to the caller — unknown, or owned by someone else + /// (the two are deliberately indistinguishable, so the endpoint reveals neither). + NotFound, +} + +/// +/// The provide-documents use case (S-10a): a zorgprofessional supplies the documents their registration +/// is parked waiting for, completing the WachtOpDocumenten task so the registratie process leaves the +/// 30-day wait and continues to beoordeling (ADR-0017). Owner-scoped by bsn. Completing the wait is +/// best-effort: if the registration never started a process (or already left the wait), the request +/// still stands, mirroring how cancels best-effort. The actual file +/// upload and its ZGW storage via the ACL is S-10b; this is the trigger that unblocks the process. +/// +public sealed class ProvideDocuments(IRegistrationStore store, IWorkflowClient workflow) +{ + public async Task HandleAsync(ProvideDocumentsCommand command, CancellationToken ct = default) + { + ArgumentNullException.ThrowIfNull(command); + + var registration = await store.GetAsync(command.RegistrationId, ct); + + // Unknown, or not the caller's registration: report NotFound either way (don't reveal which). + if (registration is null || registration.Bsn != command.Bsn) + return ProvideDocumentsOutcome.NotFound; + + // Complete the document wait (if a process is running) so beoordeling can proceed. + if (registration.ProcessInstanceId is not null) + await workflow.CompleteDocumentWaitAsync(registration.ProcessInstanceId, ct); + + return ProvideDocumentsOutcome.Accepted; + } +} diff --git a/services/domain/Big.Domain/Registration.cs b/services/domain/Big.Domain/Registration.cs index 51f20cf..a3b3503 100644 --- a/services/domain/Big.Domain/Registration.cs +++ b/services/domain/Big.Domain/Registration.cs @@ -133,8 +133,24 @@ public sealed class Registration Status = RegistrationStatus.Ingetrokken; } - // A decision (or withdrawal) is only valid while the registration is still open (INGEDIEND or - // IN_BEHANDELING). + /// + /// Expire the registration — the 30-day document-wait timer fired before the required documents + /// were supplied, so the registratie process cancels the case (S-10a). Allowed while it is still + /// open (INGEDIEND or IN_BEHANDELING) and needs no zaak; a decided (INGESCHREVEN/AFGEWEZEN) or + /// withdrawn (INGETROKKEN) registration can no longer expire. Re-expiring one already + /// is a no-op — the worker job may be redelivered (§8.6). + /// + public void Expire() + { + if (Status == RegistrationStatus.Verlopen) + return; + + RequireOpenForDecision(nameof(Expire)); + Status = RegistrationStatus.Verlopen; + } + + // A decision (or withdrawal, or expiry) is only valid while the registration is still open + // (INGEDIEND or IN_BEHANDELING). private void RequireOpenForDecision(string decision) { if (Status is not (RegistrationStatus.Ingediend or RegistrationStatus.InBehandeling)) diff --git a/services/domain/Big.Domain/RegistrationStatus.cs b/services/domain/Big.Domain/RegistrationStatus.cs index fccfac7..5e8a28c 100644 --- a/services/domain/Big.Domain/RegistrationStatus.cs +++ b/services/domain/Big.Domain/RegistrationStatus.cs @@ -21,4 +21,8 @@ public enum RegistrationStatus /// Withdrawn by the zorgprofessional before a decision (S-11). Terminal. Ingetrokken, + + /// Lapsed: the required documents were not supplied within the 30-day window, so the + /// registratie process cancelled the case (S-10a). Terminal. + Verlopen, } diff --git a/services/domain/Big.Infrastructure/FlowableWorkflowClient.cs b/services/domain/Big.Infrastructure/FlowableWorkflowClient.cs index edf379c..b3938b4 100644 --- a/services/domain/Big.Infrastructure/FlowableWorkflowClient.cs +++ b/services/domain/Big.Infrastructure/FlowableWorkflowClient.cs @@ -15,12 +15,14 @@ namespace Big.Infrastructure; /// The REST contract here is the one verified against a live flowable-rest engine (ADR-0009). /// 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( + "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> 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> AcquireRegistratieVerlopenJobsAsync(int maxJobs, CancellationToken ct = default) + { + var request = new AcquireJobsRequest(VerlopenTopic, options.LockDuration, maxJobs, options.WorkerId); + + var jobs = await PostAsync>( + "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 GetAsync(string path, CancellationToken ct) { var message = new HttpRequestMessage(HttpMethod.Get, new Uri(options.BaseUrl, path)); diff --git a/services/domain/Big.Infrastructure/IExternalWorkerClient.cs b/services/domain/Big.Infrastructure/IExternalWorkerClient.cs index ade8872..912ff59 100644 --- a/services/domain/Big.Infrastructure/IExternalWorkerClient.cs +++ b/services/domain/Big.Infrastructure/IExternalWorkerClient.cs @@ -36,3 +36,19 @@ public interface IBeoordelingEscalatieClient /// Complete an acquired escalation job so its token reaches the escalation end event. Task CompleteBeoordelingEscalatieJobAsync(string jobId, CancellationToken ct = default); } + +/// +/// The document-timeout side of the Workflow Client (S-10a): the RegistratieVerlopen +/// external-worker jobs parked by the 30-day boundary timer on WachtOpDocumenten. Kept separate +/// from the other worker ports (interface segregation) so neither the OpenZaak nor escalation worker +/// sees expiry. Implemented by — the only code that talks to +/// Flowable (§8.2, ADR-0017). +/// +public interface IRegistratieVerlopenClient +{ + /// Acquire and lock up to RegistratieVerlopen jobs. + Task> AcquireRegistratieVerlopenJobsAsync(int maxJobs, CancellationToken ct = default); + + /// Complete an acquired expiry job so its token reaches the endVerlopen end event. + Task CompleteRegistratieVerlopenJobAsync(string jobId, CancellationToken ct = default); +} diff --git a/services/domain/Big.Infrastructure/RegistratieVerlopenProcessor.cs b/services/domain/Big.Infrastructure/RegistratieVerlopenProcessor.cs new file mode 100644 index 0000000..681e765 --- /dev/null +++ b/services/domain/Big.Infrastructure/RegistratieVerlopenProcessor.cs @@ -0,0 +1,41 @@ +using Big.Application; +using Microsoft.Extensions.Logging; + +namespace Big.Infrastructure; + +/// +/// One poll tick of the document-timeout worker (S-10a, ADR-0017): acquire the parked +/// RegistratieVerlopen jobs — the tokens the 30-day boundary timer on WachtOpDocumenten +/// spawns — expire each correlated registration via the , and +/// complete the job so its token reaches endVerlopen. 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 +/// and . +/// +public sealed class RegistratieVerlopenProcessor( + IRegistratieVerlopenClient client, + ExpireRegistrationWorker worker, + ILogger logger) +{ + /// Acquire and process up to jobs. Returns the number acquired. + public async Task 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; + } +} diff --git a/services/domain/Big.Infrastructure/RegistratieVerlopenPump.cs b/services/domain/Big.Infrastructure/RegistratieVerlopenPump.cs new file mode 100644 index 0000000..4b2ba82 --- /dev/null +++ b/services/domain/Big.Infrastructure/RegistratieVerlopenPump.cs @@ -0,0 +1,49 @@ +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Logging; + +namespace Big.Infrastructure; + +/// +/// The hosted polling loop of the document-timeout worker (S-10a, ADR-0017): on an interval it +/// resolves a scoped and asks it to drain the parked +/// RegistratieVerlopen 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 . +/// +public sealed class RegistratieVerlopenPump( + IServiceScopeFactory scopeFactory, + FlowableOptions options, + ILogger logger) : BackgroundService +{ + protected override async Task ExecuteAsync(CancellationToken stoppingToken) + { + while (!stoppingToken.IsCancellationRequested) + { + try + { + using var scope = scopeFactory.CreateScope(); + var processor = scope.ServiceProvider.GetRequiredService(); + 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; + } + } + } +} diff --git a/services/domain/Big.Tests/ExpireRegistrationWorkerTests.cs b/services/domain/Big.Tests/ExpireRegistrationWorkerTests.cs new file mode 100644 index 0000000..dafc963 --- /dev/null +++ b/services/domain/Big.Tests/ExpireRegistrationWorkerTests.cs @@ -0,0 +1,84 @@ +using Big.Application; +using Big.Domain; + +namespace Big.Tests; + +// S-10a (#102): the application handler behind the RegistratieVerlopen external-worker job. The 30-day +// document-wait timer fired, so the correlated registration is expired to VERLOPEN. Mirrors +// OpenZaakWorker — pure application logic over ports, idempotent under at-least-once delivery (§8.6). +public class ExpireRegistrationWorkerTests +{ + private const string Bsn = "123456782"; + + private static Registration Submitted(string processInstanceId = "proc-1") + { + var registration = Registration.Submit(Bsn); + registration.RecordProcessStarted(processInstanceId); + return registration; + } + + [Fact] + public async Task Expires_the_registration_the_job_correlates_to() + { + var store = new FakeRegistrationStore(); + var registration = Submitted(); + store.Seed(registration); + + await new ExpireRegistrationWorker(store).HandleAsync( + new RegistratieVerlopenJob("job-7", registration.Id)); + + var saved = await store.GetAsync(registration.Id); + Assert.Equal(RegistrationStatus.Verlopen, saved!.Status); + Assert.Equal(1, store.SaveCount); + } + + [Fact] + public async Task An_already_verlopen_registration_is_not_persisted_again() + { + // A redelivered job (§8.6) finds the aggregate already VERLOPEN: a no-op, not saved again. + var store = new FakeRegistrationStore(); + var registration = Submitted(); + registration.Expire(); + store.Seed(registration); + + await new ExpireRegistrationWorker(store).HandleAsync( + new RegistratieVerlopenJob("job-7", registration.Id)); + + Assert.Equal(0, store.SaveCount); + Assert.Equal(RegistrationStatus.Verlopen, (await store.GetAsync(registration.Id))!.Status); + } + + [Fact] + public async Task An_already_resolved_registration_is_left_alone_and_the_job_completes() + { + // Race with S-11: the citizen withdrew while parked at WachtOpDocumenten, so the aggregate is + // already terminal (INGETROKKEN) when the timer's job arrives. Expiring it would violate the + // aggregate's invariant; the worker must instead no-op (and let the job complete), not throw + // into a redelivery loop. + var store = new FakeRegistrationStore(); + var registration = Submitted(); + registration.Withdraw(); + store.Seed(registration); + + await new ExpireRegistrationWorker(store).HandleAsync( + new RegistratieVerlopenJob("job-7", registration.Id)); + + Assert.Equal(0, store.SaveCount); + Assert.Equal(RegistrationStatus.Ingetrokken, (await store.GetAsync(registration.Id))!.Status); + } + + [Fact] + public async Task An_unknown_registration_throws_so_the_job_is_redelivered() + { + var store = new FakeRegistrationStore(); + + await Assert.ThrowsAsync(() => + new ExpireRegistrationWorker(store).HandleAsync( + new RegistratieVerlopenJob("job-7", RegistrationId.New()))); + } + + [Fact] + public async Task Rejects_a_null_job() + => await Assert.ThrowsAsync(() => + new ExpireRegistrationWorker(new FakeRegistrationStore()).HandleAsync(null!)); +} diff --git a/services/domain/Big.Tests/Fakes.cs b/services/domain/Big.Tests/Fakes.cs index e9c6953..f183fda 100644 --- a/services/domain/Big.Tests/Fakes.cs +++ b/services/domain/Big.Tests/Fakes.cs @@ -34,6 +34,7 @@ internal sealed class FakeWorkflowClient(string processInstanceId = "proc-1", Ac public RegistrationId? StartedFor { get; private set; } public DiplomaOrigin? StartedWithOrigin { get; private set; } public string? WithdrawnProcessInstanceId { get; private set; } + public string? CompletedDocumentWaitFor { get; private set; } public Task StartRegistrationProcessAsync( RegistrationId registrationId, DiplomaOrigin diplomaOrigin, CancellationToken ct = default) @@ -49,6 +50,12 @@ internal sealed class FakeWorkflowClient(string processInstanceId = "proc-1", Ac WithdrawnProcessInstanceId = processInstanceId; return Task.CompletedTask; } + + public Task CompleteDocumentWaitAsync(string processInstanceId, CancellationToken ct = default) + { + CompletedDocumentWaitFor = processInstanceId; + return Task.CompletedTask; + } } /// A fake user-task client for the werkbak/decision use cases: returns a scripted set of diff --git a/services/domain/Big.Tests/FlowableWorkflowClientTests.cs b/services/domain/Big.Tests/FlowableWorkflowClientTests.cs index b67d9d7..9f7b184 100644 --- a/services/domain/Big.Tests/FlowableWorkflowClientTests.cs +++ b/services/domain/Big.Tests/FlowableWorkflowClientTests.cs @@ -426,4 +426,106 @@ public class FlowableWorkflowClientTests capture.Seen.RequestUri!.ToString()); Assert.Contains("\"workerId\":\"worker-x\"", capture.Body); } + + // ── S-10a (#102): document-wait timeout → RegistratieVerlopen (ADR-0017) ────────────────────── + // A 30-day interrupting boundary timer on WachtOpDocumenten spawns a RegistratieVerlopen + // external-worker job carrying the registration id; the worker expires the registration and + // completes the job. Separately, "documents received" completes the WachtOpDocumenten user task. + + [Fact] + public async Task Acquire_verlopen_jobs_posts_the_topic_and_parses_jobs_with_their_registration_id() + { + var rid = RegistrationId.New(); + var capture = new RequestCapture(); + var client = Client(capture.Responds(HttpStatusCode.OK, + $$"""[{"id":"job-9","variables":[{"name":"registrationId","type":"string","value":"{{rid}}"}]}]""")); + + var jobs = await client.AcquireRegistratieVerlopenJobsAsync(3); + + var job = Assert.Single(jobs); + Assert.Equal("job-9", job.JobId); + Assert.Equal(rid, job.RegistrationId); + Assert.Equal("http://flowable/flowable-rest/external-job-api/acquire/jobs", + capture.Seen!.RequestUri!.ToString()); + Assert.Contains("\"topic\":\"RegistratieVerlopen\"", capture.Body); + Assert.Contains("\"numberOfTasks\":3", capture.Body); + Assert.Contains("\"workerId\":\"worker-x\"", capture.Body); + } + + [Theory] + [InlineData("[]")] + [InlineData("null")] + public async Task Acquire_verlopen_jobs_returns_empty_when_none_are_parked(string body) + { + var capture = new RequestCapture(); + var client = Client(capture.Responds(HttpStatusCode.OK, body)); + + Assert.Empty(await client.AcquireRegistratieVerlopenJobsAsync(1)); + Assert.NotNull(capture.Seen); + } + + [Fact] + public async Task Complete_verlopen_job_posts_to_the_job_complete_endpoint() + { + var capture = new RequestCapture(); + var client = Client(capture.Responds(HttpStatusCode.NoContent)); + + await client.CompleteRegistratieVerlopenJobAsync("job-9"); + + Assert.Equal(HttpMethod.Post, capture.Seen!.Method); + Assert.Equal("http://flowable/flowable-rest/external-job-api/acquire/jobs/job-9/complete", + capture.Seen.RequestUri!.ToString()); + Assert.Contains("\"workerId\":\"worker-x\"", capture.Body); + } + + [Fact] + public async Task Provide_documents_completes_the_wacht_op_documenten_task_in_the_instance() + { + var requests = new List<(HttpMethod Method, string Url, string? Body)>(); + var client = Client(new StubHandler(async req => + { + requests.Add((req.Method, req.RequestUri!.ToString(), + req.Content is null ? null : await req.Content.ReadAsStringAsync())); + return req.RequestUri!.AbsoluteUri.EndsWith("service/query/tasks") + ? new HttpResponseMessage(HttpStatusCode.OK) + { + Content = new StringContent("""{"data":[{"id":"task-3"}],"total":1}""", + Encoding.UTF8, "application/json"), + } + : new HttpResponseMessage(HttpStatusCode.OK); + })); + + await client.CompleteDocumentWaitAsync("pi-1"); + + // 1. Find the still-open WachtOpDocumenten task in this process instance. + var query = requests.Single(r => r.Url.EndsWith("service/query/tasks")); + Assert.Equal(HttpMethod.Post, query.Method); + Assert.Contains("\"processInstanceId\":\"pi-1\"", query.Body); + Assert.Contains("\"taskDefinitionKey\":\"WachtOpDocumenten\"", query.Body); + // 2. Complete that task so the process leaves the wait state. + var complete = requests.Single(r => r.Url.EndsWith("service/runtime/tasks/task-3")); + Assert.Equal(HttpMethod.Post, complete.Method); + Assert.Contains("\"action\":\"complete\"", complete.Body); + } + + [Fact] + public async Task Provide_documents_is_a_no_op_when_the_wait_task_is_no_longer_open() + { + // The process already left WachtOpDocumenten (e.g. timed out): nothing to complete, no throw. + var methods = new List(); + var client = Client(new StubHandler(req => + { + methods.Add(req.Method); + return Task.FromResult(new HttpResponseMessage(HttpStatusCode.OK) + { + Content = new StringContent("""{"data":[],"total":0}""", Encoding.UTF8, "application/json"), + }); + })); + + await client.CompleteDocumentWaitAsync("pi-1"); + + // Only the query ran; no task-completion POST followed. + Assert.DoesNotContain(methods, m => m == HttpMethod.Put || m == HttpMethod.Delete); + Assert.Single(methods); + } } diff --git a/services/domain/Big.Tests/ProvideDocumentsTests.cs b/services/domain/Big.Tests/ProvideDocumentsTests.cs new file mode 100644 index 0000000..308d863 --- /dev/null +++ b/services/domain/Big.Tests/ProvideDocumentsTests.cs @@ -0,0 +1,85 @@ +using Big.Application; +using Big.Domain; + +namespace Big.Tests; + +// S-10a (#102): the "documents received" use case. A zorgprofessional supplies the documents their +// registration is waiting for; the handler completes the WachtOpDocumenten task via the Workflow Client +// so the process leaves the 30-day wait and continues to beoordeling. Owner-scoped by the caller's bsn, +// like WithdrawRegistration. (The real file upload + ZGW storage is S-10b; this is the trigger path.) +public class ProvideDocumentsTests +{ + private const string Bsn = "123456782"; + + private static Registration Submitted(string processInstanceId = "proc-1") + { + var registration = Registration.Submit(Bsn); + registration.RecordProcessStarted(processInstanceId); + return registration; + } + + private static ProvideDocumentsCommand Command(RegistrationId id, string bsn = Bsn) => new(id, bsn); + + [Fact] + public async Task Providing_documents_completes_the_document_wait() + { + var store = new FakeRegistrationStore(); + var registration = Submitted("proc-42"); + store.Seed(registration); + var workflow = new FakeWorkflowClient(); + var handler = new ProvideDocuments(store, workflow); + + var outcome = await handler.HandleAsync(Command(registration.Id)); + + Assert.Equal(ProvideDocumentsOutcome.Accepted, outcome); + Assert.Equal("proc-42", workflow.CompletedDocumentWaitFor); + } + + [Fact] + public async Task A_different_bsn_cannot_provide_documents() + { + // Owner-scoping: only the registration's own bsn may supply its documents. Another bsn is told + // NotFound (existence not revealed) and the wait is not completed. + var store = new FakeRegistrationStore(); + var registration = Submitted(); + store.Seed(registration); + var workflow = new FakeWorkflowClient(); + var handler = new ProvideDocuments(store, workflow); + + var outcome = await handler.HandleAsync(Command(registration.Id, bsn: "999999990")); + + Assert.Equal(ProvideDocumentsOutcome.NotFound, outcome); + Assert.Null(workflow.CompletedDocumentWaitFor); + } + + [Fact] + public async Task Providing_for_an_unknown_registration_is_not_found() + { + var store = new FakeRegistrationStore(); + var handler = new ProvideDocuments(store, new FakeWorkflowClient()); + + Assert.Equal(ProvideDocumentsOutcome.NotFound, await handler.HandleAsync(Command(RegistrationId.New()))); + } + + [Fact] + public async Task Providing_before_a_process_started_is_accepted_without_calling_the_workflow() + { + // No process yet → no wait task to complete; the request still stands (best-effort, mirroring + // WithdrawRegistration) and the Workflow Client is not called. + var store = new FakeRegistrationStore(); + var registration = Registration.Submit(Bsn); // no RecordProcessStarted + store.Seed(registration); + var workflow = new FakeWorkflowClient(); + var handler = new ProvideDocuments(store, workflow); + + var outcome = await handler.HandleAsync(Command(registration.Id)); + + Assert.Equal(ProvideDocumentsOutcome.Accepted, outcome); + Assert.Null(workflow.CompletedDocumentWaitFor); + } + + [Fact] + public async Task Rejects_a_null_command() + => await Assert.ThrowsAsync(() => + new ProvideDocuments(new FakeRegistrationStore(), new FakeWorkflowClient()).HandleAsync(null!)); +} diff --git a/services/domain/Big.Tests/RegistratieVerlopenProcessorTests.cs b/services/domain/Big.Tests/RegistratieVerlopenProcessorTests.cs new file mode 100644 index 0000000..e13c950 --- /dev/null +++ b/services/domain/Big.Tests/RegistratieVerlopenProcessorTests.cs @@ -0,0 +1,79 @@ +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 +{ + /// A fake client scripting the jobs to acquire and recording completions. + private sealed class FakeVerlopenClient(params RegistratieVerlopenJob[] jobs) : IRegistratieVerlopenClient + { + public int AcquireCount { get; private set; } + public List Completed { get; } = []; + + public Task> AcquireRegistratieVerlopenJobsAsync(int maxJobs, CancellationToken ct = default) + { + AcquireCount++; + return Task.FromResult>(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.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(); + + 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.Instance).PumpOnceAsync(5); + + Assert.Equal(0, acquired); + Assert.Equal(1, client.AcquireCount); + Assert.Empty(client.Completed); + } +} diff --git a/services/domain/Big.Tests/RegistrationTests.cs b/services/domain/Big.Tests/RegistrationTests.cs index b166831..003021e 100644 --- a/services/domain/Big.Tests/RegistrationTests.cs +++ b/services/domain/Big.Tests/RegistrationTests.cs @@ -294,4 +294,63 @@ public class RegistrationTests Assert.Contains("only an INGEDIEND", ex.Message); Assert.Equal(RegistrationStatus.Afgewezen, registration.Status); } + + [Fact] + public void Expiring_an_ingediend_registration_sets_it_verlopen() + { + // The 30-day document-wait timer fired before the documents arrived (S-10a): the case is + // cancelled and the aggregate becomes terminal VERLOPEN. + var registration = Registration.Submit("123456782"); + + registration.Expire(); + + Assert.Equal(RegistrationStatus.Verlopen, registration.Status); + } + + [Fact] + public void Expiring_needs_no_zaak() + { + // The timer fires on a purely time-based boundary; expiry does not depend on the zaak. + var registration = Registration.Submit("123456782"); + + registration.Expire(); + + Assert.Equal(RegistrationStatus.Verlopen, registration.Status); + Assert.Null(registration.ZaakUrl); + } + + [Fact] + public void Re_expiring_an_already_verlopen_registration_is_idempotent() + { + // The RegistratieVerlopen worker job may be redelivered (§8.6); re-expiring is a no-op. + var registration = Registration.Submit("123456782"); + registration.Expire(); + + registration.Expire(); + + Assert.Equal(RegistrationStatus.Verlopen, registration.Status); + } + + [Fact] + public void Expiring_an_approved_registration_is_rejected() + { + var registration = Registration.Submit("123456782"); + registration.AttachZaak(new Uri("http://openzaak/zaken/api/v1/zaken/abc")); + registration.Approve(); + + var ex = Assert.Throws(() => registration.Expire()); + Assert.Contains("only an INGEDIEND", ex.Message); + Assert.Equal(RegistrationStatus.Ingeschreven, registration.Status); + } + + [Fact] + public void Expiring_a_withdrawn_registration_is_rejected() + { + var registration = Registration.Submit("123456782"); + registration.Withdraw(); + + var ex = Assert.Throws(() => registration.Expire()); + Assert.Contains("only an INGEDIEND", ex.Message); + Assert.Equal(RegistrationStatus.Ingetrokken, registration.Status); + } } diff --git a/services/domain/stryker-config.json b/services/domain/stryker-config.json index 0391844..ac22608 100644 --- a/services/domain/stryker-config.json +++ b/services/domain/stryker-config.json @@ -5,7 +5,8 @@ "reporters": ["progress", "html"], "mutate": [ "!**/OpenZaakJobPump.cs", - "!**/BeoordelingEscalatiePump.cs" + "!**/BeoordelingEscalatiePump.cs", + "!**/RegistratieVerlopenPump.cs" ], "thresholds": { "high": 95, diff --git a/tests/acceptance/Features/EenDocumentTermijnVerlopen.feature b/tests/acceptance/Features/EenDocumentTermijnVerlopen.feature new file mode 100644 index 0000000..7ec02bf --- /dev/null +++ b/tests/acceptance/Features/EenDocumentTermijnVerlopen.feature @@ -0,0 +1,23 @@ +# language: en +# Drives S-10a (#102). After the zaak is opened the process parks at WachtOpDocumenten with an +# INTERRUPTING 30-day boundary timer. If the documents do not arrive in time the timer cancels the +# task and parks a RegistratieVerlopen job (ADR-0017) which the timeout worker drains, expiring the +# registration to VERLOPEN. Documents received before the timer fires close the wait, so no expiry +# happens. This scenario exercises the timeout worker against an in-memory Flowable stand-in; the timer +# firing live is verify-domain. +Feature: Een documenttermijn laten verlopen + Als registerbeheerder wil ik dat een aanvraag waarvoor de documenten niet binnen 30 dagen binnen zijn + automatisch vervalt zodat onvolledige aanvragen niet blijven liggen. + + Scenario: Zonder documenten binnen 30 dagen vervalt de registratie + Given a registration parked at the WachtOpDocumenten task + When the 30-day document timer fires + And the document-timeout worker runs + Then the registration is verlopen + + Scenario: Tijdig aangeleverde documenten laten de registratie niet vervallen + Given a registration parked at the WachtOpDocumenten task + When the documents arrive before the timer fires + And the 30-day document timer fires + And the document-timeout worker runs + Then the registration is not verlopen diff --git a/tests/acceptance/Steps/EenDocumentTermijnVerlopenSteps.cs b/tests/acceptance/Steps/EenDocumentTermijnVerlopenSteps.cs new file mode 100644 index 0000000..402cfde --- /dev/null +++ b/tests/acceptance/Steps/EenDocumentTermijnVerlopenSteps.cs @@ -0,0 +1,52 @@ +using Acceptance.Support; +using Big.Application; +using Big.Domain; +using Big.Infrastructure; +using Microsoft.Extensions.Logging.Abstractions; +using Reqnroll; +using Xunit; + +namespace Acceptance.Steps; + +/// Bindings for EenDocumentTermijnVerlopen.feature (S-10a). Drives the timeout worker +/// ( over the ) against +/// an in-memory Flowable stand-in and a shared registration store; one instance per scenario. The +/// interrupting 30-day timer either cancels the wait and expires the registration, or — if the +/// documents arrived first — never fires; the scenario asserts on the aggregate's status. +[Binding] +[Scope(Feature = "Een documenttermijn laten verlopen")] +public sealed class EenDocumentTermijnVerlopenSteps +{ + private readonly InMemoryDocumentTimeoutClient _flowable = new(); + private readonly Support.InMemoryRegistrationStore _store = new(); + private Registration _registration = null!; + private string _processInstanceId = ""; + + [Given("a registration parked at the WachtOpDocumenten task")] + public async Task GivenARegistrationParkedAtWachtOpDocumenten() + { + _registration = Registration.Submit("123456782"); + await _store.SaveAsync(_registration); + _processInstanceId = _flowable.ParkWaitingForDocuments(_registration.Id); + } + + [When("the 30-day document timer fires")] + public void WhenTheDocumentTimerFires() => _flowable.FireDocumentTimer(_processInstanceId); + + [When("the documents arrive before the timer fires")] + public void WhenTheDocumentsArriveBeforeTheTimer() => _flowable.ReceiveDocuments(_processInstanceId); + + [When("the document-timeout worker runs")] + public async Task WhenTheTimeoutWorkerRuns() + => await new RegistratieVerlopenProcessor( + _flowable, new ExpireRegistrationWorker(_store), + NullLogger.Instance).PumpOnceAsync(5); + + [Then("the registration is verlopen")] + public async Task ThenTheRegistrationIsVerlopen() + => Assert.Equal(RegistrationStatus.Verlopen, (await _store.GetAsync(_registration.Id))!.Status); + + [Then("the registration is not verlopen")] + public async Task ThenTheRegistrationIsNotVerlopen() + => Assert.Equal(RegistrationStatus.Ingediend, (await _store.GetAsync(_registration.Id))!.Status); +} diff --git a/tests/acceptance/Support/BffAcceptanceHost.cs b/tests/acceptance/Support/BffAcceptanceHost.cs index 5374b6f..ee4cd8c 100644 --- a/tests/acceptance/Support/BffAcceptanceHost.cs +++ b/tests/acceptance/Support/BffAcceptanceHost.cs @@ -72,6 +72,9 @@ public sealed class CapturingDomainClient : IDomainClient public Task WithdrawRegistrationAsync(string registrationId, string bsn, CancellationToken ct = default) => Task.FromResult(true); + public Task ProvideDocumentsAsync(string registrationId, string bsn, CancellationToken ct = default) + => Task.FromResult(true); + public Task> GetWerkbakAsync(CancellationToken ct = default) => Task.FromResult>([]); diff --git a/tests/acceptance/Support/InMemoryDomainPorts.cs b/tests/acceptance/Support/InMemoryDomainPorts.cs index af07de0..f480320 100644 --- a/tests/acceptance/Support/InMemoryDomainPorts.cs +++ b/tests/acceptance/Support/InMemoryDomainPorts.cs @@ -14,6 +14,7 @@ public sealed class InMemoryWorkflowClient : IWorkflowClient public RegistrationId? StartedFor { get; private set; } public DiplomaOrigin? StartedWithOrigin { get; private set; } public string? WithdrawnProcessInstanceId { get; private set; } + public string? CompletedDocumentWaitFor { get; private set; } public Task StartRegistrationProcessAsync( RegistrationId registrationId, DiplomaOrigin diplomaOrigin, CancellationToken ct = default) @@ -28,6 +29,12 @@ public sealed class InMemoryWorkflowClient : IWorkflowClient WithdrawnProcessInstanceId = processInstanceId; return Task.CompletedTask; } + + public Task CompleteDocumentWaitAsync(string processInstanceId, CancellationToken ct = default) + { + CompletedDocumentWaitFor = processInstanceId; + return Task.CompletedTask; + } } /// An in-memory ACL stand-in: records the bsn it opened a zaak for and returns a fixed URL, @@ -130,6 +137,57 @@ public sealed class InMemoryEscalatieClient : IBeoordelingEscalatieClient } } +/// An in-memory Flowable stand-in for the document-timeout scenario (S-10a): it models one +/// WachtOpDocumenten wait per process instance — whether it is still open and the registration it +/// correlates to — and the RegistratieVerlopen jobs the interrupting 30-day boundary timer parks. It +/// drives the timeout worker's behaviour without a running Flowable; the timer firing live is the +/// verify-domain check. +public sealed class InMemoryDocumentTimeoutClient : IRegistratieVerlopenClient +{ + private sealed class Wait + { + public required RegistrationId RegistrationId { get; init; } + public bool IsWaiting { get; set; } = true; + } + + private readonly Dictionary _waits = []; + private readonly List _parked = []; + private int _seq; + + /// A registration parks at WachtOpDocumenten, waiting for the citizen's documents. + public string ParkWaitingForDocuments(RegistrationId registrationId) + { + var pid = $"pi-{++_seq}"; + _waits[pid] = new Wait { RegistrationId = registrationId }; + return pid; + } + + /// The documents arrive before the timer fires: the wait task closes, so the interrupting + /// timer no longer fires (mirrors the Workflow Client completing WachtOpDocumenten). + public void ReceiveDocuments(string processInstanceId) => _waits[processInstanceId].IsWaiting = false; + + /// The 30-day interrupting boundary timer fires: if still waiting, it cancels the wait and + /// parks a RegistratieVerlopen job carrying the correlated registration id. A no-op if the documents + /// already arrived (the wait/timer race, §8.6). + public void FireDocumentTimer(string processInstanceId) + { + var wait = _waits[processInstanceId]; + if (!wait.IsWaiting) + return; + wait.IsWaiting = false; + _parked.Add(new RegistratieVerlopenJob($"job-{++_seq}", wait.RegistrationId)); + } + + public Task> AcquireRegistratieVerlopenJobsAsync(int maxJobs, CancellationToken ct = default) + => Task.FromResult>(_parked.Take(maxJobs).ToList()); + + public Task CompleteRegistratieVerlopenJobAsync(string jobId, CancellationToken ct = default) + { + _parked.RemoveAll(j => j.JobId == jobId); + return Task.CompletedTask; + } +} + /// An in-memory registration store for the domain acceptance scenario. public sealed class InMemoryRegistrationStore : IRegistrationStore { diff --git a/tests/e2e/registration.spec.ts b/tests/e2e/registration.spec.ts index c6395ff..b508352 100644 --- a/tests/e2e/registration.spec.ts +++ b/tests/e2e/registration.spec.ts @@ -1,11 +1,15 @@ import { expect, test } from '@playwright/test'; -// Walking-skeleton happy path (S-08d + S-09 + S-09b + S-12): a zorgprofessional logs in via mock -// DigiD and submits through the self-service portal → BFF → domain; the entry appears in the openbaar -// register as INGEDIEND; a behandelaar then logs in to the behandel portal, finds the registration in -// the werkbak, and approves it (goedkeuren); the decision completes the Flowable Beoordelen task and -// flows via the ACL → NRC → event-subscriber → projection, and the openbaar register shows INGESCHREVEN. -test('DigiD submit → public INGEDIEND → behandelaar goedkeurt → public INGESCHREVEN', async ({ page }) => { +// Walking-skeleton happy path (S-08d + S-09 + S-09b + S-12 + S-10a): a zorgprofessional logs in via +// mock DigiD and submits through the self-service portal → BFF → domain; the entry appears in the +// openbaar register as INGEDIEND; the citizen supplies the documents the process is waiting for +// (S-10a); a behandelaar then logs in to the behandel portal, finds the registration in the werkbak, +// and approves it (goedkeuren); the decision completes the Flowable Beoordelen task and flows via the +// ACL → NRC → event-subscriber → projection, and the openbaar register shows INGESCHREVEN. +test('DigiD submit → public INGEDIEND → documenten → behandelaar goedkeurt → public INGESCHREVEN', async ({ + page, + context, +}) => { // Visiting the guarded page redirects to the Keycloak (mock DigiD) login. await page.goto('/'); @@ -26,42 +30,53 @@ test('DigiD submit → public INGEDIEND → behandelaar goedkeurt → public ING expect(reference, 'the confirmation shows a registration reference').toBeTruthy(); // The openbaar register (anonymous, its own origin) shows the submitted entry once the projection - // catches up. The projection updates asynchronously (NRC → event-subscriber), and the register loads - // on open, so reload until *this* submission's row appears. We poll on the reference cell (not a - // generic INGEDIEND cell): the shared verify stack already holds INGEDIEND rows from earlier checks, - // so a status-only poll would short-circuit on a stale row before our row is projected. - await page.goto('http://openbaar/'); - await expect(page.getByRole('heading', { name: /Openbaar BIG-register/i })).toBeVisible(); + // catches up. We check it on a SEPARATE page so the self-service tab keeps its (in-memory) submitted + // state — the "Documenten aanleveren" action below acts on that same session. The projection updates + // asynchronously (NRC → event-subscriber), so reload until *this* submission's row appears. We poll + // on the reference cell (not a generic INGEDIEND cell): the shared verify stack already holds + // INGEDIEND rows from earlier checks, so a status-only poll would short-circuit on a stale row. + const staff = await context.newPage(); + await staff.goto('http://openbaar/'); + await expect(staff.getByRole('heading', { name: /Openbaar BIG-register/i })).toBeVisible(); // #78: the reference shown in the public register must be the exact one the citizen saw on the // submit confirmation — no mismatch between the two portals. await expect .poll(async () => { - await page.reload(); - return page.getByRole('cell', { name: reference }).count(); + await staff.reload(); + return staff.getByRole('cell', { name: reference }).count(); }, { timeout: 30_000, intervals: [1_000, 2_000, 3_000, 5_000] }) .toBeGreaterThan(0); - await expect(page.getByRole('row', { name: reference }).getByRole('cell', { name: 'INGEDIEND' })) + await expect(staff.getByRole('row', { name: reference }).getByRole('cell', { name: 'INGEDIEND' })) .toBeVisible(); - // A behandelaar picks the registration up in the behandel-portal werkbak and approves it - // (goedkeuren) — the S-12 flow that replaces the temporary admin endpoint. Navigating here switches - // to the medewerker realm (a different Keycloak realm than the citizen's digid session). - await page.goto('http://behandel/'); - await page.locator('#username').fill('merel-behandelaar'); - await page.locator('#password').fill('test123'); - await page.locator('#kc-login').click(); + // Provide the documents the registration is waiting for (S-10a), on the still-open self-service tab. + // The process parks at WachtOpDocumenten only after the zaak is opened; the INGEDIEND row above proves + // the zaak exists — so the OpenZaak worker has completed and the process is now at the wait — which is + // why we supply the documents here rather than right after submit, when the trigger would race the + // wait and no-op. (S-10b turns this into a real file upload; here it is the trigger that unblocks + // beoordeling.) + await page.getByRole('button', { name: /documenten aanleveren/i }).click(); + await expect(page.getByText(/documenten zijn aangeleverd/i)).toBeVisible(); - await expect(page.getByRole('heading', { name: /Werkbak/i })).toBeVisible(); + // A behandelaar picks the registration up in the behandel-portal werkbak and approves it (goedkeuren) + // — the S-12 flow that replaces the temporary admin endpoint. The staff tab switches to the + // medewerker realm (a different Keycloak realm than the citizen's digid session). + await staff.goto('http://behandel/'); + await staff.locator('#username').fill('merel-behandelaar'); + await staff.locator('#password').fill('test123'); + await staff.locator('#kc-login').click(); - // The registration parks at the Beoordelen user task only after the worker has opened its zaak, so + await expect(staff.getByRole('heading', { name: /Werkbak/i })).toBeVisible(); + + // The registration reaches the Beoordelen user task only after its documents are provided (above), so // it appears in the werkbak asynchronously — reload until this reference's row shows up. Target the // decide button by reference (not a generic "Goedkeuren"): the shared verify stack holds other open // tasks, so a positional match could act on someone else's registration. - const goedkeuren = page.getByRole('button', { name: `Goedkeuren ${reference}` }); + const goedkeuren = staff.getByRole('button', { name: `Goedkeuren ${reference}` }); await expect .poll(async () => { - await page.reload(); + await staff.reload(); return goedkeuren.count(); }, { timeout: 30_000, intervals: [1_000, 2_000, 3_000, 5_000] }) .toBeGreaterThan(0); @@ -69,7 +84,7 @@ test('DigiD submit → public INGEDIEND → behandelaar goedkeurt → public ING // Click and wait for the decide POST to finish (204) BEFORE leaving the page. `click()` only // dispatches the request; navigating away immediately cancels it in flight (nginx logs a 499) and // the decision never reaches the domain — so the registration would stay INGEDIEND. - const decided = page.waitForResponse( + const decided = staff.waitForResponse( (r) => r.url().includes(`/behandel/registrations/${reference}/decide`) && r.request().method() === 'POST', @@ -79,11 +94,11 @@ test('DigiD submit → public INGEDIEND → behandelaar goedkeurt → public ING // The approval flows back to the projection; back on the openbaar register *our* row (matched by // its reference) now shows INGESCHREVEN. - await page.goto('http://openbaar/'); + await staff.goto('http://openbaar/'); await expect .poll(async () => { - await page.reload(); - return page.getByRole('row', { name: reference }).getByRole('cell', { name: 'INGESCHREVEN' }).count(); + await staff.reload(); + return staff.getByRole('row', { name: reference }).getByRole('cell', { name: 'INGESCHREVEN' }).count(); }, { timeout: 30_000, intervals: [1_000, 2_000, 3_000, 5_000] }) .toBeGreaterThan(0); }); diff --git a/workflows/registratie.bpmn b/workflows/registratie.bpmn index 360297e..669a6ab 100644 --- a/workflows/registratie.bpmn +++ b/workflows/registratie.bpmn @@ -22,10 +22,16 @@ (BeoordelingEscaleren); the Workflow Client reassigns the still-open Beoordelen task from the behandelaar group to teamlead (ADR-0015). The Beoordelen task stays open throughout — the timer only changes who may claim it. - S-13 adds diploma-eligibility routing: between OpenZaakAanmaken and Beoordelen a DMN service + S-13 adds diploma-eligibility routing: between the document wait and Beoordelen a DMN service task (flowable:type="dmn") evaluates the `diploma-eligibility` decision on the diplomaOrigin start variable; an exclusive gateway routes a foreign diploma through the CBGV-advies user task - before Beoordelen, a domestic one straight there (ADR-0016). --> + before Beoordelen, a domestic one straight there (ADR-0016). + S-10a adds the document wait: right after the zaak is opened the process parks at a + WachtOpDocumenten user task with an INTERRUPTING P30D boundary timer. "Documents received" + (the S-10b upload, via the Workflow Client) completes the task and the process continues to the + diploma routing; if the 30 days lapse first the timer cancels the task and runs the + RegistratieVerlopen external-worker task, whose worker expires the registration to VERLOPEN, + ending the process as "verlopen" (ADR-0017). --> @@ -38,7 +44,32 @@ flowable:type="external-worker" flowable:topic="OpenZaakAanmaken"/> - + + + + + + + + + + P30D + + + + + + + + + + +