Compare commits

..
Author SHA1 Message Date
not 2d783448b7 fix(e2e): assert the register record where a real approval happens (refs #149)
CI / build (pull_request) Successful in 1m6s
CI / lint (pull_request) Successful in 1m19s
CI / unit (pull_request) Successful in 1m24s
CI / verify-stack (pull_request) Successful in 7m49s
CI / frontend (pull_request) Successful in 2m57s
CI / mutation (pull_request) Successful in 6m7s
verify-domain was the wrong home for the assertion, and CI was right to fail it.
That check completes the Beoordelen task straight through Flowable REST — on
purpose, it exists to exercise the Workflow Client's REST contract — which
bypasses the domain `decide` path that calls the ACL. No approval reached the
ACL there, so no record was ever written.

The Playwright happy path is the only check that drives a real approval
(behandel portal → BFF → domain → ACL), and it already knows its own reference.
Assert there instead: exactly one RegisterRecord for that reference,
INGESCHREVEN, carrying nothing outside the public-safe schema. Drops
register-record-check.py and the verify-domain block.

The helper was run under real Playwright against a live Objecten before
committing — one record found, none for an unknown reference.
2026-08-14 10:43:34 +02:00
not 10b784cc05 fix(infra): reach Objecten by service name in the register-record check (refs #149)
CI / build (pull_request) Successful in 1m8s
CI / lint (pull_request) Successful in 1m23s
CI / unit (pull_request) Successful in 1m23s
CI / frontend (pull_request) Successful in 2m59s
CI / mutation (pull_request) Successful in 6m8s
CI / verify-stack (pull_request) Failing after 7m24s
CI caught my own check falling into the constraint ADR-0028 documents: it looked
Objecttypen up by container IP, so the objecttype URL came back IP-addressed and
Objecten rejected it as "not one of the available choices". Reach both by
service name — compose DNS resolves them, and neither request has OpenZaak's
URL-validity constraint that made IPs necessary elsewhere in this script.

The 400 also spent the full 60s timeout disguised as "transport:" because
HTTPError is a URLError subclass. Handle it separately: a 4xx now fails
immediately with the response body, which is where the real reason was.

Verified both ways against a live Objecten: absent record → exit 1 with the
reason, present record → exit 0.
2026-08-14 10:20:26 +02:00
not 2bb7d9c165 test(acl): ObjectenGateway integration test against live Objecten (refs #149)
CI / build (pull_request) Successful in 4m20s
CI / lint (pull_request) Successful in 4m38s
CI / unit (pull_request) Successful in 1m27s
CI / frontend (pull_request) Successful in 4m1s
CI / mutation (pull_request) Successful in 6m52s
CI / verify-stack (pull_request) Failing after 18m4s
Drives the real gateway against a running Objecten + Objecttypen pair: two
writes for the same id leave exactly one object carrying the second write's
status, and nothing outside the public-safe schema. Runs under verify-acl,
inside the compose network — which it must, because Objecttypen echoes the
request Host into the objecttype `url` and Objecten only accepts the one
matching its configured api_root. ADR-0028 records that constraint.
2026-08-14 09:44:15 +02:00
not 4a047c618c fix(infra): let Objecten actually accept the register record (refs #149)
Replaying the gateway's calls against a live Objecten + Objecttypen pair turned
up two blockers CI would only have found after the fact:

- Objecten rejects an objecttype it has not been configured with, and it
  identifies one by uuid — assigned at seed time by a one-shot that runs after
  Objecten's static setup_configuration. Pin the uuid on both sides instead.
- Objecten notifies on every write and notifications_api_common *raises* when
  that config is absent, so every POST 500'd after rolling the object back.
  Objecten → NRC has no broker, worker, kanaal or abonnement yet, so disable
  notifications rather than wire a client that drops every message; S-19b turns
  them on for real.

With both in place the full exchange verifies end to end: lookup → version →
search → create → update (still one object), and a record carrying a bsn is
rejected by the schema. ADR-0028 records both.
2026-08-14 09:41:00 +02:00
not 43b45ad756 fix(acl): read each objecttype version instead of the versions collection (refs #149)
The version resolve assumed `GET {objecttype}/versions` returns a bare list.
Every other collection in the Objecttypen API returns a paginated envelope, and
nothing in the repo exercises that endpoint, so the shape was a guess. Follow
the path infra/registerrecord-check.py already proves against the real API
instead: read the `versions` URLs off the objecttype and fetch each for its
status. Costs a request per version, once per gateway instance.

ACL mutation score 92.23% (baseline 91.37%).
2026-08-14 09:33:15 +02:00
not 5502e4c099 test(acl): raise the Objecten gateway above the mutation ratchet (refs #149)
The new gateway landed at 77.6%, dragging the ACL score under its 90 break
threshold. The gaps were all real behaviour nobody was asserting: a failed or
empty read being mistaken for "nothing there yet" and followed by a blind
write, a `results`-less response taking down the resolve with an
ArgumentNullException, the CRS headers going to Objecttypen (which is not a geo
API), and the write body being sent chunked. ACL score 86.63% → 92.08%.
2026-08-14 09:29:03 +02:00
not 3705a18e18 docs: ADR-0028 + demo note — Objecten holds the register (refs #149)
ADR-0028 records why the register record lives in Objecten rather than as zaak
eigenschappen, why the ACL owns the hop, and how two non-atomic writes are made
to converge instead. Also retires the PRD §15 out-of-scope line the slice
supersedes.
2026-08-14 09:23:41 +02:00
not 400bdcafc4 test(infra): assert the approval wrote the register record to Objecten (refs #149)
verify-domain already drives a full approval; it now also asserts Objecten holds
exactly one RegisterRecord for that registration — matched on its own reference,
because the shared verify stack carries records from earlier runs. The check
covers the three things that can silently go wrong: the record is missing (the
ACL's Objecten hop never ran), duplicated (the upsert is not idempotent), or
carries a field outside the public-safe schema.
2026-08-14 09:22:09 +02:00
not c67ee7d3f5 refactor(acl): resolve the objecttype's highest published version (refs #149)
Counting the `versions` URLs assumed a contiguous, all-published list. Read the
objecttype's versions collection instead and take the highest one whose status
is `published`, so a draft version — whose schema is still being shaped — is
never written against.
2026-08-14 09:19:50 +02:00
not d14f379358 feat(acl): write the RegisterRecord to Objecten on approval (refs #149)
ApproveZaakAsync now does two writes: the ZGW eindstatus (the process) and the
register record in Objecten (the register). The record is keyed on the zaak
UUID — the same key the read projection rows carry — and its reference is the
zaak's identificatie, so nothing personal crosses into the world-readable
register (ADR-0027).

ObjectenGateway resolves the objecttype by name (its URL and version are
assigned at seed time, as with ADR-0021), searches for an existing object by
data attribute, then POSTs or PATCHes. Resolution is lazy, so the ACL needs no
depends_on on Objecten and does not crash-loop when it boots first.
2026-08-14 09:18:34 +02:00
not 66f8322580 test(acl): approval writes the RegisterRecord to Objecten (refs #149)
Ports and failing tests for the Objecten hop, ahead of the implementation:

- IRegisterRecordGateway + RegisterRecord — the Application-side port; the
  record mirrors the objecttype schema registered in S-18c (ADR-0027).
- AclService takes the port but does not yet call it, so the approval test
  fails on an empty upsert list.
- ObjectenGateway is a shell throwing NotImplementedException; its tests pin
  the contract: resolve the objecttype by name, search by data attribute,
  POST when absent / PATCH when present, static Token auth per API, the CRS
  headers the geo API requires, and a surfaced error body.

Also splits S-19 (#20) into #149/#150 in BACKLOG.md — the approval-side write
and the projection re-sourcing are independently deployable (CLAUDE.md §13).
2026-08-14 09:16:23 +02:00
40 changed files with 270 additions and 1256 deletions
+1 -6
View File
@@ -219,9 +219,6 @@ jobs:
- name: OpenZaak → NRC → Event Subscriber → projection-api - name: OpenZaak → NRC → Event Subscriber → projection-api
id: projection id: projection
run: make verify-projection run: make verify-projection
- name: Objecten → NRC notification delivery
id: objecten_nrc
run: make verify-objecten-notifications
- name: Domain → Flowable → ACL → OpenZaak - name: Domain → Flowable → ACL → OpenZaak
id: domain id: domain
run: make verify-domain run: make verify-domain
@@ -248,7 +245,6 @@ jobs:
OBJECTTYPEN: ${{ steps.objecttypen.outcome }} OBJECTTYPEN: ${{ steps.objecttypen.outcome }}
OBJECTEN: ${{ steps.objecten.outcome }} OBJECTEN: ${{ steps.objecten.outcome }}
REGISTERRECORD: ${{ steps.registerrecord.outcome }} REGISTERRECORD: ${{ steps.registerrecord.outcome }}
OBJECTEN_NOTIFICATIONS: ${{ steps.objecten_nrc.outcome }}
ACL: ${{ steps.acl.outcome }} ACL: ${{ steps.acl.outcome }}
NRC: ${{ steps.nrc.outcome }} NRC: ${{ steps.nrc.outcome }}
PROJECTION: ${{ steps.projection.outcome }} PROJECTION: ${{ steps.projection.outcome }}
@@ -270,7 +266,6 @@ jobs:
echo "| Objecttypen API + token | $(icon "$OBJECTTYPEN") |" echo "| Objecttypen API + token | $(icon "$OBJECTTYPEN") |"
echo "| Objecten API + token | $(icon "$OBJECTEN") |" echo "| Objecten API + token | $(icon "$OBJECTEN") |"
echo "| RegisterRecord objecttype | $(icon "$REGISTERRECORD") |" echo "| RegisterRecord objecttype | $(icon "$REGISTERRECORD") |"
echo "| Objecten → NRC | $(icon "$OBJECTEN_NOTIFICATIONS") |"
echo "| ACL ↔ OpenZaak | $(icon "$ACL") |" echo "| ACL ↔ OpenZaak | $(icon "$ACL") |"
echo "| OpenZaak → NRC | $(icon "$NRC") |" echo "| OpenZaak → NRC | $(icon "$NRC") |"
echo "| NRC → Event Subscriber → projection | $(icon "$PROJECTION") |" echo "| NRC → Event Subscriber → projection | $(icon "$PROJECTION") |"
@@ -290,7 +285,7 @@ jobs:
# Log dump must precede teardown (which removes the containers). # Log dump must precede teardown (which removes the containers).
- name: Dump container logs on failure - name: Dump container logs on failure
if: failure() if: failure()
run: docker compose -f infra/docker-compose.yml logs --no-color --tail=100 oz-init openzaak nrc-init nrc-web nrc-celery nrc-beat flowable-db flowable-rest flowable-init keycloak acl bff domain projection-db event-subscriber projection-api self-service openbaar behandel beheer objecttypen-db objecttypen-redis objecttypen-init objecttypen objecten-db objecten-redis objecten-init objecten objecten-celery registerrecord-init tempo prometheus grafana 2>&1 || true run: docker compose -f infra/docker-compose.yml logs --no-color --tail=100 oz-init openzaak nrc-init nrc-web nrc-celery nrc-beat flowable-db flowable-rest flowable-init keycloak acl bff domain projection-db event-subscriber projection-api self-service openbaar behandel beheer objecttypen-db objecttypen-redis objecttypen-init objecttypen objecten-db objecten-redis objecten-init objecten registerrecord-init tempo prometheus grafana 2>&1 || true
- name: Tear down - name: Tear down
if: always() if: always()
run: make down run: make down
+1 -3
View File
@@ -296,9 +296,7 @@ Split into independently deployable sub-slices (CLAUDE.md §13):
Split into independently deployable sub-slices (CLAUDE.md §13): Split into independently deployable sub-slices (CLAUDE.md §13):
- **S-19a** (#149, ✅) · ACL writes the `RegisterRecord` to Objecten on approval, idempotently, alongside the ZGW eindstatus. Carries the ADR (ADR-0028). - **S-19a** (#149, ✅) · ACL writes the `RegisterRecord` to Objecten on approval, idempotently, alongside the ZGW eindstatus. Carries the ADR (ADR-0028).
- **S-19b** (#150, ✅) · Read projection sourced from Objecten instead of NRC zaak events. *(split — #150 closed)* - **S-19b** (#150) · Read projection sourced from Objecten instead of NRC zaak events. Depends on S-19a.
- **S-19b-1** (#152, ✅) · Objecten publishes to NRC — broker, celery worker, `objecten` kanaal, notifications config. Turns back on what ADR-0028 deliberately disabled.
- **S-19b-2** (#153, ✅) · Projection derived from `RegisterRecord` objects, rebuildable from the Objecten-derived log. The ACL also writes an INGEDIEND record on submit, so the register holds the whole lifecycle. Carries ADR-0030.
--- ---
+1 -7
View File
@@ -43,7 +43,7 @@ export DOCKER_HOST := unix://$(PODMAN_SOCK)
endif endif
endif endif
.PHONY: ci lint build unit mutation frontend integration verify verify-up verify-acl verify-nrc verify-projection verify-bff verify-domain verify-observability verify-tracing verify-metrics verify-objecttypen verify-objecten verify-registerrecord verify-objecten-notifications verify-notifications smoke up down local verify-local local-down changelog openzaak-up openzaak-smoke openzaak-seed openzaak-down stack-up stack-smoke stack-down keycloak-up keycloak-smoke keycloak-down flowable-up flowable-smoke flowable-down help .PHONY: ci lint build unit mutation frontend integration verify verify-up verify-acl verify-nrc verify-projection verify-bff verify-domain verify-observability verify-tracing verify-metrics verify-objecttypen verify-objecten verify-registerrecord verify-notifications smoke up down local verify-local local-down changelog openzaak-up openzaak-smoke openzaak-seed openzaak-down stack-up stack-smoke stack-down keycloak-up keycloak-smoke keycloak-down flowable-up flowable-smoke flowable-down help
## ci: run the full pipeline — lint, build, unit, mutation, frontend, verify (mirrors Gitea Actions) ## ci: run the full pipeline — lint, build, unit, mutation, frontend, verify (mirrors Gitea Actions)
## `verify` is the live-stack stage (full stack up once → ACL + notification checks). ## `verify` is the live-stack stage (full stack up once → ACL + notification checks).
@@ -201,11 +201,6 @@ verify-objecten:
verify-registerrecord: verify-registerrecord:
bash infra/run-registerrecord-check.sh bash infra/run-registerrecord-check.sh
## verify-objecten-notifications: assert a RegisterRecord write in Objecten is DELIVERED as an
## `objecten` notification via NRC (S-19b-1), against the already-running stack.
verify-objecten-notifications:
bash infra/run-objecten-notifications-check.sh
## verify: local mirror of the CI verify-stack job — full stack up once, all checks, ## verify: local mirror of the CI verify-stack job — full stack up once, all checks,
## tear down (always). For fast single-concern local iteration use `integration` ## tear down (always). For fast single-concern local iteration use `integration`
## (oz-only) or `verify-notifications` (oz+nrc) instead. ## (oz-only) or `verify-notifications` (oz+nrc) instead.
@@ -217,7 +212,6 @@ verify:
&& bash infra/run-acl-integration.sh \ && bash infra/run-acl-integration.sh \
&& bash infra/run-notification-check.sh \ && bash infra/run-notification-check.sh \
&& bash infra/run-projection-check.sh \ && bash infra/run-projection-check.sh \
&& bash infra/run-objecten-notifications-check.sh \
&& bash infra/run-domain-check.sh \ && bash infra/run-domain-check.sh \
&& bash infra/run-bff-check.sh \ && bash infra/run-bff-check.sh \
&& bash infra/run-e2e-check.sh || rc=$$?; \ && bash infra/run-e2e-check.sh || rc=$$?; \
@@ -119,8 +119,8 @@ every message was dropped on the floor — a delivery path that looks wired and
- ponytail ceiling: Objecten emits no notifications, so nothing downstream can react to a - ponytail ceiling: Objecten emits no notifications, so nothing downstream can react to a
register write yet. register write yet.
- **Lifted by ADR-0029** (S-19b-1, #152): broker, worker, `objecten` kanaal and - Upgrade path: S-19b (#150) needs those notifications to source the projection from
notifications config now exist, and `NOTIFICATIONS_DISABLED` is `false`. Objecten, and turns them on together with the broker, worker, kanaal and abonnement.
## Consequences ## Consequences
@@ -130,8 +130,8 @@ every message was dropped on the floor — a delivery path that looks wired and
independent of the case that produced it. independent of the case that produced it.
- The disclosure boundary is enforced by Objecten's schema validation (ADR-0027), not by - The disclosure boundary is enforced by Objecten's schema validation (ADR-0027), not by
discipline in projection code. discipline in projection code.
- The read projection can become a cache of Objecten rather than a re-derivation of ZGW - The read projection can become a cache of Objecten rather than a re-derivation of ZGW
done in S-19b-2 (#153), ADR-0030. (S-19b, #150).
**Negative / costs** **Negative / costs**
@@ -142,9 +142,8 @@ every message was dropped on the floor — a delivery path that looks wired and
(`Acl__Objecten__Token`) in compose. (`Acl__Objecten__Token`) in compose.
- Two new hand-kept constants: the pinned objecttype UUID (two files) and the objecttype - Two new hand-kept constants: the pinned objecttype UUID (two files) and the objecttype
name (compose + `register.py`). name (compose + `register.py`).
- ~~Until S-19b lands, the public register is still read from the NRC-derived projection, so - Until S-19b lands, the public register is still read from the NRC-derived projection, so
the register record is written but not yet read — the two must agree.~~ Closed by ADR-0030: the register record is written but not yet read — the two must agree.
the projection is now derived from the register, so there is only one source to agree with.
## Coupling rules touched (CLAUDE.md §8) ## Coupling rules touched (CLAUDE.md §8)
@@ -1,122 +0,0 @@
# ADR-0029: Objecten publishes register events to NRC
- **Status:** Accepted
- **Date:** 2026-08-14
- **Deciders:** Respellion engineering
- **Slice:** S-19b-1 (#152), first of the S-19b (#150) split
- **Supersedes in part:** ADR-0028's "Objecten's notifications are off for this slice"
## Context
ADR-0028 put the authoritative register record in the Objecten API and had the ACL write
it on approval. It also switched Objecten's notifications **off** — deliberately, with a
stated ceiling: there was no broker, no worker, no `objecten` kanaal and no abonnement, so
turning the client side on alone would have produced a delivery path that looks wired and
drops every message.
S-19b-2 (#153) wants the read projection sourced from register writes rather than
re-derived from ZGW zaak events. That needs the notifications to actually arrive. This ADR
builds the four missing pieces and lifts the ceiling.
## Decision
**Objecten publishes to the same NRC OpenZaak already publishes to, on the `objecten`
kanaal, delivered by its own Celery worker — provisioned declaratively on both sides,
exactly as ADR-0007 did for OpenZaak.**
- **Objecten** (`infra/objecten/setup_configuration/data.yaml`): a `zgw_consumers` service
`nrc` (api_type `nrc`) plus a `notifications_config` step naming it, and
`NOTIFICATIONS_DISABLED: "false"` in both compose files.
- **NRC** (`infra/opennotificaties/setup_configuration/data.yaml`): an `objecten` kanaal
alongside `zaken`.
- **`objecten-celery`**: a worker container on the Objecten image (`/celery_worker.sh`),
mirroring `oz-celery`, with `CELERY_BROKER_URL`/`CELERY_RESULT_BACKEND` on
`objecten-redis` db 1 (db 0 is already the cache).
### One NRC, one credential, one kanaal per publisher
Objecten reuses the `big-reference-seed` client OpenZaak publishes with. NRC verifies its
JWT and authorizes it against OpenZaak's Autorisaties API (ADR-0007), which grants that
client `heeft_alle_autorisaties` — so no second credential and no publisher-specific
authorization is needed. A second NRC, or a second credential, would buy isolation this
reference application has no use for.
The kanaal name is **not ours to choose**: the Objects API sends
`NOTIFICATIONS_KANAAL = "objecten"`. NRC rejects a publish to an unregistered kanaal
(`"Kanaal met deze naam bestaat niet"`), which is precisely what the failing check for this
slice reported first. Its filter set (`object_type`) matches the kenmerken the Objects API
sends, so an abonnement can narrow to one objecttype instead of receiving every write.
### Writers address Objecten as `objecten.local` — NRC rejects single-label hosts
NRC types a notification's `hoofdObject` and `resourceUrl` as DRF `URLField`s, so Django's
`URLValidator` runs on them — and it refuses a **single-label** host. Objecten fills both
from the object `url` that DRF built with `request.build_absolute_uri`, i.e. **the Host the
caller used**. Write to `http://objecten:8000` and NRC answers every publish with
```
{"hoofdObject":["Voer een geldige URL in."],"resourceUrl":["Voer een geldige URL in."]}
```
which `objecten-celery` then retries with exponential backoff, forever, in the background —
the write itself having returned 201.
`SITE_DOMAIN` does **not** fix this; it is not what builds those URLs. The fix is on the
caller side: the `objecten` service carries an `objecten.local` network alias, and every
component whose writes must be notified — the ACL (`Acl__Objecten__BaseUrl`), the gateway
integration tests, this slice's verify check — addresses it there. An alias rather than a
plain dotted `SITE_DOMAIN` so the host still **resolves in-network**: a subscriber that
follows `resourceUrl` reaches the record it points at, which S-19b-2 will do. Readers are
unaffected and keep using the plain service name.
This is the same class of constraint as ADR-0028's "the ACL's Objecttypen base URL must
match Objecten's configured `api_root`": these modules put request-derived hosts into data
another module then validates or dereferences.
- ponytail ceiling: nothing *enforces* that a new writer uses the alias — it would get a 201
and silently no notification.
- Upgrade path: if a second writer ever appears, rename the compose service to `objecten.local`
so the plain name stops working, rather than adding a lint.
### A worker, not a synchronous send
`notifications_api_common` only schedules the send on transaction commit. Without a worker
the task sits in redis forever and every register write is silently undelivered — the exact
half-wired state ADR-0028 refused to ship. No `beat` for Objecten: it is a publisher, not a
subscriber, and `nrc-beat` already drains NRC's delivery queue.
## Verification
`make verify-objecten-notifications` (`infra/run-objecten-notifications-check.sh`, in the
CI `verify-stack` job) registers an abonnement on the `objecten` kanaal pointing at a
throwaway webhook sink, writes a `RegisterRecord` exactly as the ACL does on approval, and
asserts the notification reaches the sink. That is the whole chain in one assertion:
Objecten → `objecten-celery` → NRC → `nrc-beat` → the callback. Any missing piece — broker,
worker, kanaal, notifications config — shows up as a non-delivery rather than as a green
config.
## Consequences
**Positive**
- A register write is now observable by anything that subscribes, which is what S-19b-2
(#153) needs to make the projection a cache of Objecten rather than a re-derivation of ZGW.
- ADR-0028's ceiling is lifted: the delivery path is proven end to end, not merely configured.
**Negative / costs**
- One more long-running container (`objecten-celery`) on an already memory-tight CI runner.
- A second publisher on the shared `big-reference-seed` credential — a credential rotation
now touches two modules.
- Objecten now has two in-network names, and which one a caller uses silently decides
whether its writes are notified (ceiling above).
- ponytail ceiling: notification delivery has no dead-letter or alerting — a failed publish
is visible only in the worker log.
- Upgrade path: if undelivered register events start mattering, subscribe an audit sink or
read NRC's own delivery admin rather than building a retry layer here.
## Coupling rules touched (CLAUDE.md §8)
None bent. This is infrastructure between two upstream modules, over their documented
APIs; no service reaches another's database. §8.6 (idempotency at every event boundary)
applies to whatever consumes the new kanaal — S-19b-2's problem, not this slice's.
@@ -1,141 +0,0 @@
# ADR-0030: The read projection is sourced from the register, not from ZGW
- **Status:** Accepted
- **Date:** 2026-08-28
- **Deciders:** Respellion engineering
- **Slice:** S-19b-2 (#153), second of the S-19b (#150) split
- **Builds on:** ADR-0008 (read projection store), ADR-0028 (Objecten holds the register), ADR-0029 (Objecten publishes to NRC)
## Context
ADR-0028 moved the authoritative register record into the Objecten API, and said what should
follow: "the read projection can become a cache of Objecten rather than a re-derivation of
ZGW." Until this slice it was still the latter — the Event Subscriber listened on the `zaken`
kanaal and inferred register state from case events:
- a `zaak`/`create` meant INGEDIEND;
- any `status`/`create` was taken to be the approval, so meant INGESCHREVEN — the subscriber
may not read OpenZaak (§8.1), so it could not tell one statustype from another;
- the citizen-facing reference was not in the notification at all, so every projection had a
second hop: ask the ACL for the zaak's identificatie (#78).
So the register — a fact about a person — was reconstructed by guessing at the lifecycle of the
case that happened to produce it. ADR-0029 made the register itself publish. This ADR switches
the projection over to it.
## Decision
**The Event Subscriber listens on the `objecten` kanaal and projects the `RegisterRecord` the
notification points at. The projection is a cache of the register; ZGW is no longer a source.**
- The subscriber's abonnement moves from `zaken` to `objecten` (`register-abonnement.py`, and
the CI projection check).
- An Objecten notification carries **no record data** — only the object URL and the objecttype
as a kenmerk — so the record is read back through the ACL (`POST /register-records/read`).
§8.1 applies to Objecten exactly as ADR-0028 established: the ACL is the only code that talks
to it.
- The accepted acties are `create`, `update` and `partial_update`. The last one is not
defensive breadth: the ACL upserts with PATCH, and DRF routes a PATCH through the notifying
`update()` while naming the action `partial_update` — which is what Objecten publishes. So
every approval arrives as `partial_update`, and accepting only `create`/`update` drops the
one state change this slice exists to project. `destroy` is deliberately not accepted:
removing a registration from the public register is its own decision.
- The record already carries `id`, `status` and `reference`, so the row is the record. The
zaak-shaped surface goes: `IsZaakCreated`, `IsZaakStatusSet`, `ZaakUrl`, `ZaakId`, and
`ToEntry`'s `Resource == "status"` inference are replaced by `IsRegisterRecordWritten` +
`ObjectUrl`, and the ACL enrichment hop disappears.
### The ACL writes an INGEDIEND record on submit
Before this slice only approval wrote a record, so re-sourcing alone would have silently
dropped every INGEDIEND row from the public register. `OpenZaakAsync` therefore upserts a
record with status INGEDIEND after opening the zaak, keyed on the same zaak id that approval
later upserts to INGESCHREVEN.
This is the same two-writes-converging posture ADR-0028 already accepted for approval, now on
the submit path too: both writes are idempotent, so a retried submit updates the record rather
than adding a second one (§8.6). The reference comes from the registration itself, so unlike
approval this path needs no ZGW read-back.
The alternative — a register holding only INGESCHREVEN — is arguably the more correct reading
of "public register", but it narrows what the openbaar portal shows and reads against PRD §68
("~50 register entries with diverse statuses"). Rejected as a behaviour change this slice was
not asked to make.
### The dedup key is the projected row, not the notification
NRC carries no notification id and may redeliver, so the idempotency key is derived from
content (as before). The obvious candidates both break here:
- **the object URL alone** — the ACL upserts *one object per registration*, so submit and
approval notify about the same URL, and the approval would be swallowed as a duplicate;
- **object URL + actie** — a retried approval is a second `update`, so it would be dropped
while genuinely being the same state (harmless), but a *third* distinct state would collide
with it (not harmless).
The key is therefore the object plus the state that write puts in the projection —
`objecten:object:{url}:{status}:{reference}`. A redelivery collapses; a genuine state change
does not. That is exactly the property §8.6 asks for, and it needs no version field from
Objecten's internals.
### The notification log holds the row, not the event
`processed_notifications` stops describing ZGW events (`actie`, `zaak_id`, `resource`) and
holds the projected row itself (`register_id`, `status`, `reference`). A rebuild becomes a
replay with no mapping rules and no upstream reads at all — §8.4 held before via the ACL hop;
now it holds outright.
The migration **drops** the old columns rather than renaming them. EF scaffolded renames
(`resource``register_id`, `zaak_id``status`) that would have carried ZGW values into
columns meaning something else entirely, and a rebuild would then have projected that garbage.
- ponytail ceiling: the migration empties both tables. A pre-slice row describes a zaak event
the new projector cannot reproject, and the registrations behind those rows have no
RegisterRecord in Objecten (only approvals wrote one), so they are not re-derivable from the
new source either.
- Upgrade path: fine while stacks are ephemeral. If a long-lived environment ever needs to keep
them, backfill by walking Objecten's objects rather than replaying the log.
## Consequences
**Positive**
- The register is read from the register. The projection is a derived cache of a first-class
record, not an inference over someone else's lifecycle.
- The "any status-create is the approval" guess is gone — a real source of wrongness the moment
the zaaktype grows a second statustype.
- One hop fewer per notification: the record carries its own reference, so the ACL enrichment
call disappears.
- A rebuild needs nothing but its own log (§8.4).
**Negative / costs**
- Submission is now two writes across two modules and eventually consistent. A failure between
them leaves a zaak with no register record until the submit is retried; nothing repairs that
automatically yet — the same gap ADR-0028 recorded for approval, now on a second path.
- The projection lags the register by a notification round trip, where it used to lag the zaak
by one. In practice the same order of magnitude.
- Projecting now depends on the ACL being reachable, where the reference enrichment used to be
the only ACL dependency. A failed read means the notification is not logged and not
projected — NRC retries, so it converges, but the failure mode is now on the main path.
- OpenZaak still publishes to `zaken` and nothing in the product listens. Kept because the
`verify-nrc` check asserts that path, and turning off a working publisher to save nothing
would be its own risk.
## Coupling rules touched (CLAUDE.md §8)
None bent. §8.1 holds — the subscriber reaches Objecten only through the ACL. §8.4 is
strengthened: the projection is rebuildable from its own log, with no upstream reads at all.
§8.6 is what the dedup-key discussion above is about.
## Verification
`make verify-projection` (`infra/run-projection-check.sh`, in CI's `verify-stack`) opens a zaak
**through the ACL** and asserts projection-api serves a row for it with status INGEDIEND — the
whole new chain in one assertion: ACL → Objecten → `objecten-celery` → NRC → `nrc-beat`
Event Subscriber → projection → projection-api. A zaak created behind the ACL's back produces
no row, which is the re-source working rather than a gap.
`RegisterProjectieBijwerken.feature` covers the use case in business language, including the
approval case — the same row moving INGEDIEND → INGESCHREVEN, which is now one registration's
record being updated rather than two unrelated ZGW events.
+7 -32
View File
@@ -341,8 +341,7 @@ services:
# Objecten holds the register, OpenZaak holds the process (S-19a, ADR-0028). Both APIs take a # Objecten holds the register, OpenZaak holds the process (S-19a, ADR-0028). Both APIs take a
# static token, not a ZGW JWT. The objecttype URL is assigned at seed time, so the ACL resolves # static token, not a ZGW JWT. The objecttype URL is assigned at seed time, so the ACL resolves
# it by name — lazily, on the first approval, so no depends_on is needed here. # it by name — lazily, on the first approval, so no depends_on is needed here.
# Dotted host on purpose — see the `objecten.local` alias below (ADR-0029). Acl__Objecten__BaseUrl: http://objecten:8000/
Acl__Objecten__BaseUrl: http://objecten.local:8000/
Acl__Objecten__Token: ${OBJECTEN_TOKEN:-1234567890abcdef1234567890abcdef12345678} Acl__Objecten__Token: ${OBJECTEN_TOKEN:-1234567890abcdef1234567890abcdef12345678}
Acl__Objecten__ObjecttypenBaseUrl: http://objecttypen:8000/ Acl__Objecten__ObjecttypenBaseUrl: http://objecttypen:8000/
Acl__Objecten__ObjecttypenToken: ${OBJECTTYPEN_TOKEN:-0123456789abcdef0123456789abcdef01234567} Acl__Objecten__ObjecttypenToken: ${OBJECTTYPEN_TOKEN:-0123456789abcdef0123456789abcdef01234567}
@@ -692,14 +691,12 @@ services:
CACHE_AXES: objecten-redis:6379/0 CACHE_AXES: objecten-redis:6379/0
DISABLE_2FA: "true" DISABLE_2FA: "true"
OTEL_SDK_DISABLED: "true" OTEL_SDK_DISABLED: "true"
CELERY_BROKER_URL: redis://objecten-redis:6379/1 # S-19a: Objecten refuses every write while its Notificaties config is absent
CELERY_RESULT_BACKEND: redis://objecten-redis:6379/1 # (notifications_api_common raises rather than skipping, so POST /objects 500s). Objecten →
# Publish register-record events to NRC on the `objecten` kanaal (S-19b-1, ADR-0029). The NRC # NRC is not wired yet — there is no broker, worker, kanaal or abonnement for it — so turn
# service + notifications_config are provisioned by setup_configuration # notifications off rather than fake a delivery path that silently drops every message.
# (infra/objecten/setup_configuration/data.yaml), and objecten-celery below actually sends # S-19b (#150) sources the projection from Objecten and turns this back on for real.
# them — notifications_api_common only queues the task. See ADR-0028 for why S-19a left this NOTIFICATIONS_DISABLED: "true"
# off until all four pieces existed.
NOTIFICATIONS_DISABLED: "false"
RUN_SETUP_CONFIG: "true" RUN_SETUP_CONFIG: "true"
command: /setup_configuration.sh command: /setup_configuration.sh
volumes: volumes:
@@ -724,28 +721,6 @@ services:
start_period: 30s start_period: 30s
ports: ports:
- "8021:8000" - "8021:8000"
depends_on:
objecten-init:
condition: service_completed_successfully
networks:
cg:
# Objecten reflects the *request* Host into the `url` it returns, and
# notifications_api_common publishes that url as the notification's hoofdObject /
# resourceUrl — which NRC types as a URLField, and Django's URLValidator rejects a
# single-label host ("Voer een geldige URL in."). So every caller whose writes must be
# notified addresses Objecten by this dotted alias instead of `objecten` (ADR-0029).
# Reads are unaffected and still use the plain service name.
aliases:
- objecten.local
# The celery worker that actually delivers Objecten's notifications to NRC (S-19b-1, ADR-0029).
# notifications_api_common only schedules the send on transaction commit; without a worker the
# task sits in redis forever and every register write is silently undelivered. Mirrors oz-celery.
# No beat: Objecten is a publisher, not a subscriber — nrc-beat drains the delivery queue.
objecten-celery:
image: docker.io/maykinmedia/objects-api:${OBJECTS_TAG:-3.4.0}
environment: *objecten-env-local
command: /celery_worker.sh
depends_on: depends_on:
objecten-init: objecten-init:
condition: service_completed_successfully condition: service_completed_successfully
+7 -32
View File
@@ -326,8 +326,7 @@ services:
# Objecten holds the register, OpenZaak holds the process (S-19a, ADR-0028). Both APIs take a # Objecten holds the register, OpenZaak holds the process (S-19a, ADR-0028). Both APIs take a
# static token, not a ZGW JWT. The objecttype URL is assigned at seed time, so the ACL resolves # static token, not a ZGW JWT. The objecttype URL is assigned at seed time, so the ACL resolves
# it by name — lazily, on the first approval, so no depends_on is needed here. # it by name — lazily, on the first approval, so no depends_on is needed here.
# Dotted host on purpose — see the `objecten.local` alias below (ADR-0029). Acl__Objecten__BaseUrl: http://objecten:8000/
Acl__Objecten__BaseUrl: http://objecten.local:8000/
Acl__Objecten__Token: ${OBJECTEN_TOKEN:-1234567890abcdef1234567890abcdef12345678} Acl__Objecten__Token: ${OBJECTEN_TOKEN:-1234567890abcdef1234567890abcdef12345678}
Acl__Objecten__ObjecttypenBaseUrl: http://objecttypen:8000/ Acl__Objecten__ObjecttypenBaseUrl: http://objecttypen:8000/
Acl__Objecten__ObjecttypenToken: ${OBJECTTYPEN_TOKEN:-0123456789abcdef0123456789abcdef01234567} Acl__Objecten__ObjecttypenToken: ${OBJECTTYPEN_TOKEN:-0123456789abcdef0123456789abcdef01234567}
@@ -718,14 +717,12 @@ services:
CACHE_AXES: objecten-redis:6379/0 CACHE_AXES: objecten-redis:6379/0
DISABLE_2FA: "true" DISABLE_2FA: "true"
OTEL_SDK_DISABLED: "true" OTEL_SDK_DISABLED: "true"
CELERY_BROKER_URL: redis://objecten-redis:6379/1 # S-19a: Objecten refuses every write while its Notificaties config is absent
CELERY_RESULT_BACKEND: redis://objecten-redis:6379/1 # (notifications_api_common raises rather than skipping, so POST /objects 500s). Objecten →
# Publish register-record events to NRC on the `objecten` kanaal (S-19b-1, ADR-0029). The NRC # NRC is not wired yet — there is no broker, worker, kanaal or abonnement for it — so turn
# service + notifications_config are provisioned by setup_configuration # notifications off rather than fake a delivery path that silently drops every message.
# (infra/objecten/setup_configuration/data.yaml), and objecten-celery below actually sends # S-19b (#150) sources the projection from Objecten and turns this back on for real.
# them — notifications_api_common only queues the task. See ADR-0028 for why S-19a left this NOTIFICATIONS_DISABLED: "true"
# off until all four pieces existed.
NOTIFICATIONS_DISABLED: "false"
RUN_SETUP_CONFIG: "true" RUN_SETUP_CONFIG: "true"
command: /setup_configuration.sh command: /setup_configuration.sh
# data.yaml is streamed into this external volume by infra/seed-config.sh before start. # data.yaml is streamed into this external volume by infra/seed-config.sh before start.
@@ -753,28 +750,6 @@ services:
start_period: 30s start_period: 30s
ports: ports:
- "8021:8000" - "8021:8000"
depends_on:
objecten-init:
condition: service_completed_successfully
networks:
cg:
# Objecten reflects the *request* Host into the `url` it returns, and
# notifications_api_common publishes that url as the notification's hoofdObject /
# resourceUrl — which NRC types as a URLField, and Django's URLValidator rejects a
# single-label host ("Voer een geldige URL in."). So every caller whose writes must be
# notified addresses Objecten by this dotted alias instead of `objecten` (ADR-0029).
# Reads are unaffected and still use the plain service name.
aliases:
- objecten.local
# The celery worker that actually delivers Objecten's notifications to NRC (S-19b-1, ADR-0029).
# notifications_api_common only schedules the send on transaction commit; without a worker the
# task sits in redis forever and every register write is silently undelivered. Mirrors oz-celery.
# No beat: Objecten is a publisher, not a subscriber — nrc-beat drains the delivery queue.
objecten-celery:
image: docker.io/maykinmedia/objects-api:${OBJECTS_TAG:-3.4.0}
environment: *objecten-env
command: /celery_worker.sh
depends_on: depends_on:
objecten-init: objecten-init:
condition: service_completed_successfully condition: service_completed_successfully
+6 -11
View File
@@ -2,10 +2,10 @@
"""Local-stack bootstrap (S-B04, #110, ADR-0020) — register the NRC abonnement. """Local-stack bootstrap (S-B04, #110, ADR-0020) — register the NRC abonnement.
Runs as the `nrc-subscribe` init container of infra/docker-compose.local.yml. Registers an Runs as the `nrc-subscribe` init container of infra/docker-compose.local.yml. Registers an
abonnement on the `objecten` kanaal pointing at the event-subscriber's /notifications callback, so abonnement on the `zaken` kanaal pointing at the event-subscriber's /notifications callback, so
the register writes the ACL makes (INGEDIEND on submit, INGESCHREVEN on approval) reach the OpenZaak's notifications (zaak create + status set) reach the projection — without this the openbaar
projection — without this the openbaar (public) register stays empty. Since S-19b-2 the projection (public) register stays empty. This is what infra/verify-notification-driver.py does for CI (minus
is sourced from the register in Objecten, not from ZGW zaak events (ADR-0030). the test zaak it also creates).
The callback host is the event-subscriber's resolved **container IP**, not `event-subscriber`, because The callback host is the event-subscriber's resolved **container IP**, not `event-subscriber`, because
NRC validates callbackUrl with Django's URLValidator (a single-label host is rejected — same reason the NRC validates callbackUrl with Django's URLValidator (a single-label host is rejected — same reason the
@@ -22,8 +22,6 @@ SINK_PORT = os.environ.get("SINK_PORT", "8080")
SINK_AUTH = os.environ.get("SINK_AUTH", "Bearer big-reference-notifications") SINK_AUTH = os.environ.get("SINK_AUTH", "Bearer big-reference-notifications")
CID = os.environ.get("OZ_CLIENT_ID", "big-reference-seed") CID = os.environ.get("OZ_CLIENT_ID", "big-reference-seed")
SECRET = os.environ.get("OZ_SECRET", "insecure-dev-secret-change-me") SECRET = os.environ.get("OZ_SECRET", "insecure-dev-secret-change-me")
# The projection is sourced from the register in Objecten, not from ZGW zaak events (S-19b-2).
KANAAL = "objecten"
def token(): def token():
@@ -62,10 +60,7 @@ def main():
status, body = call("GET", f"{NRC}/api/v1/abonnement") status, body = call("GET", f"{NRC}/api/v1/abonnement")
for ab in (body or []) if status == 200 else []: for ab in (body or []) if status == 200 else []:
if str(ab.get("callbackUrl", "")).endswith("/notifications"): if str(ab.get("callbackUrl", "")).endswith("/notifications"):
# The kanaal is part of "current": an abonnement left over from before S-19b-2 points at if ab.get("callbackUrl") == callback:
# the right callback but listens on `zaken`, and would never be replaced on IP alone.
kanalen = [k.get("naam") for k in ab.get("kanalen", [])]
if ab.get("callbackUrl") == callback and kanalen == [KANAAL]:
print(f"abonnement already current: {ab['url']}") print(f"abonnement already current: {ab['url']}")
return return
call("DELETE", ab["url"]) call("DELETE", ab["url"])
@@ -73,7 +68,7 @@ def main():
status, ab = call("POST", f"{NRC}/api/v1/abonnement", { status, ab = call("POST", f"{NRC}/api/v1/abonnement", {
"callbackUrl": callback, "auth": SINK_AUTH, "callbackUrl": callback, "auth": SINK_AUTH,
"kanalen": [{"naam": KANAAL, "filters": {}}]}) "kanalen": [{"naam": "zaken", "filters": {}}]})
if status != 201: if status != 201:
sys.exit(f"create abonnement -> {status}: {json.dumps(ab)}") sys.exit(f"create abonnement -> {status}: {json.dumps(ab)}")
print(f"abonnement registered: {ab['url']} -> {callback}") print(f"abonnement registered: {ab['url']} -> {callback}")
-121
View File
@@ -1,121 +0,0 @@
#!/usr/bin/env python3
"""S-19b-1 (#152): driver for the Objecten → NRC notification check.
Registers an abonnement on the `objecten` kanaal pointing at the webhook sink, then writes a
RegisterRecord object exactly as the ACL's ObjectenGateway does (S-19a). The caller
(run-objecten-notifications-check.sh) watches the sink for the delivery — this only sets it up,
and prints `OBJECT_URL <url>` for the caller to grep on.
Delivery exercises the whole chain: Objecten → its celery worker → NRC → nrc-beat → the callback.
Anything missing (broker, worker, kanaal, notifications config) shows up as a non-delivery.
Stdlib only so it runs in a bare python:3-slim container on the compose network.
"""
import base64
import hashlib
import hmac
import json
import os
import sys
import time
import urllib.error
import urllib.request
OBJECTEN = os.environ["OBJECTEN"] # http://objecten:8000
OBJECTEN_TOKEN = os.environ["OBJECTEN_TOKEN"]
OBJECTTYPEN = os.environ["OBJECTTYPEN"] # http://objecttypen:8000
OBJECTTYPEN_TOKEN = os.environ["OBJECTTYPEN_TOKEN"]
NRC_BASE = os.environ["NRC_BASE"] # http://<nrc-ip>:8000
SINK_CALLBACK = os.environ["SINK_CALLBACK"] # http://<sink-ip>:9000/
SINK_AUTH = os.environ["SINK_AUTH"]
CLIENT_ID = os.environ.get("NRC_CLIENT_ID", "big-reference-seed")
SECRET = os.environ.get("NRC_SECRET", "insecure-dev-secret-change-me")
KANAAL = "objecten"
def mint():
"""The HS256 JWT NRC expects (same shape as infra/local/register-abonnement.py)."""
def seg(d):
return base64.urlsafe_b64encode(json.dumps(d).encode()).rstrip(b"=")
payload = seg({
"iss": CLIENT_ID, "iat": int(time.time()), "client_id": CLIENT_ID,
"user_id": CLIENT_ID, "user_representation": CLIENT_ID,
})
signing_input = seg({"typ": "JWT", "alg": "HS256"}) + b"." + payload
signature = base64.urlsafe_b64encode(
hmac.new(SECRET.encode(), signing_input, hashlib.sha256).digest()).rstrip(b"=")
return (signing_input + b"." + signature).decode()
def nrc(method, url, body=None):
"""Call NRC. `url` may be a path or an absolute URL (the list returns absolute ones)."""
data = json.dumps(body).encode() if body is not None else None
req = urllib.request.Request(
url if url.startswith("http") else f"{NRC_BASE}{url}", data=data, method=method,
headers={"Authorization": f"Bearer {mint()}", "Content-Type": "application/json"})
try:
with urllib.request.urlopen(req, timeout=15) as r:
return json.load(r) if r.length != 0 else {}
except urllib.error.HTTPError as e:
# The body carries the reason (e.g. an unregistered kanaal); the status alone does not.
raise SystemExit(f"FAIL — NRC {method} {url}{e.code}: {e.read().decode(errors='replace')[:400]}")
def token_api(base, token, method, path, body=None, crs=False):
data = json.dumps(body).encode() if body is not None else None
headers = {"Authorization": f"Token {token}"}
if body is not None:
headers["Content-Type"] = "application/json"
if crs:
headers["Accept-Crs"] = "EPSG:4326"
if body is not None:
headers["Content-Crs"] = "EPSG:4326"
req = urllib.request.Request(f"{base}{path}", data=data, method=method, headers=headers)
with urllib.request.urlopen(req, timeout=15) as r:
return json.load(r) if r.length != 0 else {}
def subscribe():
"""Register an abonnement on the objecten kanaal, replacing a stale one for the same callback."""
# NRC returns a bare list here, not a paginated envelope.
for existing in nrc("GET", "/api/v1/abonnement") or []:
if existing.get("callbackUrl") == SINK_CALLBACK:
nrc("DELETE", existing["url"])
nrc("POST", "/api/v1/abonnement", {
"callbackUrl": SINK_CALLBACK,
"auth": SINK_AUTH,
"kanalen": [{"naam": KANAAL, "filters": {}}],
})
print(f">> abonnement on '{KANAAL}' -> {SINK_CALLBACK}")
def objecttype_url():
results = token_api(OBJECTTYPEN, OBJECTTYPEN_TOKEN, "GET", "/api/v2/objecttypes").get("results", [])
match = next((o for o in results if o.get("name") == "RegisterRecord"), None)
if not match:
print("FAIL — no RegisterRecord objecttype in Objecttypen", file=sys.stderr)
raise SystemExit(1)
return match["url"]
def main():
subscribe()
reference = f"NOTIF-{int(time.time())}"
created = token_api(OBJECTEN, OBJECTEN_TOKEN, "POST", "/api/v2/objects", {
"type": objecttype_url(),
"record": {
"typeVersion": 1,
"data": {"id": f"zaak-{reference}", "status": "INGESCHREVEN", "reference": reference},
"startAt": time.strftime("%Y-%m-%d"),
},
}, crs=True)
print(f">> wrote RegisterRecord {created['url']}")
# An NRC notification carries no record data — only hoofdObject/resourceUrl — so the object
# URL, not the reference in its data, is what the caller can correlate the delivery on.
print(f"OBJECT_URL {created['url']}")
return 0
if __name__ == "__main__":
sys.exit(main())
@@ -18,16 +18,6 @@ zgw_consumers:
auth_type: api_key auth_type: api_key
header_key: Authorization header_key: Authorization
header_value: Token 0123456789abcdef0123456789abcdef01234567 header_value: Token 0123456789abcdef0123456789abcdef01234567
# (1b) The NRC Objecten publishes register-record events to (S-19b-1, ADR-0029). Same shape and
# same big-reference-seed credential OpenZaak publishes with — NRC verifies the JWT and
# authorizes it via OpenZaak's AC, which grants that client heeft_alle_autorisaties.
- identifier: nrc
label: Open Notificaties
api_type: nrc
api_root: http://nrc-web:8000/api/v1/
auth_type: zgw
client_id: big-reference-seed
secret: insecure-dev-secret-change-me
# (2) Permit the RegisterRecord objecttype (S-19a). Objecten refuses to store an object whose # (2) Permit the RegisterRecord objecttype (S-19a). Objecten refuses to store an object whose
# objecttype it has not been configured with ("ObjectType with url=… is not configured"), and it # objecttype it has not been configured with ("ObjectType with url=… is not configured"), and it
@@ -50,10 +40,3 @@ tokenauth:
email: admin@localhost email: admin@localhost
organization: Respellion organization: Respellion
is_superuser: true is_superuser: true
# (4) Point Objecten's notifications at that NRC service (S-19b-1, ADR-0029). Requires
# NOTIFICATIONS_DISABLED=false plus a celery broker + worker — without the worker the message is
# queued and never sent, which is exactly the half-wired state S-19a refused to ship (ADR-0028).
notifications_config_enable: true
notifications_config:
notifications_api_service_identifier: nrc
@@ -29,9 +29,7 @@ autorisaties_api_config_enable: true
autorisaties_api: autorisaties_api:
authorizations_api_service_identifier: openzaak-ac authorizations_api_service_identifier: openzaak-ac
# 4. The kanalen publishers announce on: `zaken` (OpenZaak) and `objecten` (Objecten, S-19b-1). # 4. The kanaal OpenZaak publishes zaak events on.
# Both authenticate with the big-reference-seed credential above, which OpenZaak's AC grants
# heeft_alle_autorisaties — so no separate publisher authorization is needed for Objecten.
notifications_kanalen_config_enable: true notifications_kanalen_config_enable: true
notifications_kanalen_config: notifications_kanalen_config:
items: items:
@@ -41,11 +39,3 @@ notifications_kanalen_config:
- bronorganisatie - bronorganisatie
- zaaktype - zaaktype
- vertrouwelijkheidaanduiding - vertrouwelijkheidaanduiding
# 5. The kanaal Objecten publishes register-record events on (S-19b-1, ADR-0029). Its name is
# fixed by the Objects API itself (NOTIFICATIONS_KANAAL = "objecten"), not chosen here. The
# filter set matches what the Objects API sends as kenmerken, so an abonnement can narrow by
# objecttype rather than receiving every object write in the register.
- naam: objecten
documentatie_link: https://objects-and-objecttypes-api.readthedocs.io/
filters:
- object_type
-83
View File
@@ -1,83 +0,0 @@
#!/usr/bin/env bash
#
# S-19b-1 (#152): verify the Objecten → NRC notification path against an ALREADY-RUNNING full
# stack. Registers an abonnement on the `objecten` kanaal pointing at a throwaway webhook sink,
# writes a RegisterRecord object (exactly as the ACL does on approval, S-19a), and asserts the sink
# receives the notification.
#
# This is the whole publish chain in one assertion: Objecten → its celery worker → NRC → nrc-beat →
# the subscriber callback. S-19a deliberately left it disconnected (ADR-0028); this proves it is
# connected for real, rather than merely configured.
#
# All in-network, reaching services by container IP (a single-label host isn't URL-valid for NRC's
# callbackUrl validator; the runner can't reach published ports — gitea-actions-gotchas.md §5/§6).
# EXCEPT Objecttypen, which must be reached by SERVICE NAME: it echoes the request Host into the
# objecttype `url` and Objecten only accepts the one matching its configured api_root (ADR-0028);
# and Objecten, reached by its `objecten.local` alias because it reflects the request Host into the
# notification's hoofdObject/resourceUrl, which NRC validates as a URL (ADR-0029).
#
# Does NOT manage the stack lifecycle, but cleans up the sink/driver it creates.
set -euo pipefail
here="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
SINK_AUTH="Bearer objecten-notification-sink-token"
cleanup() { docker rm -f rr-osink rr-overify >/dev/null 2>&1 || true; }
trap cleanup EXIT
ip() { docker inspect -f '{{range .NetworkSettings.Networks}}{{.IPAddress}}{{end}}' "$1"; }
# Anchored on the compose replica suffix so they don't also match objecten-db / objecten-redis.
obj="$(docker ps -q --filter 'name=objecten[-_][0-9]+$' | head -1)"
nrc="$(docker ps -q --filter 'name=nrc-web' | head -1)"
[ -n "$obj" ] || { echo "ERROR: no running objecten container — bring the stack up first" >&2; exit 1; }
[ -n "$nrc" ] || { echo "ERROR: no running nrc-web container — bring the stack up first" >&2; exit 1; }
net="$(docker inspect -f '{{range $k,$_ := .NetworkSettings.Networks}}{{$k}}{{"\n"}}{{end}}' "$obj" | head -1)"
nrc_ip="$(ip "$nrc")"
echo ">> network=$net nrc=$nrc_ip"
echo ">> starting the webhook sink"
docker rm -f rr-osink >/dev/null 2>&1 || true
sink="$(docker create --network "$net" --name rr-osink -e "EXPECTED_AUTH=$SINK_AUTH" \
python:3-slim python /sink.py)"
docker cp "$here/notification-sink.py" "$sink:/sink.py" >/dev/null
docker start "$sink" >/dev/null
sleep 1
sink_ip="$(ip rr-osink)"
echo ">> sink at $sink_ip:9000"
echo ">> registering the abonnement + writing a RegisterRecord"
docker rm -f rr-overify >/dev/null 2>&1 || true
drv="$(docker create --network "$net" --name rr-overify \
-e "OBJECTEN=http://objecten.local:8000" \
-e "OBJECTEN_TOKEN=${OBJECTEN_TOKEN:-1234567890abcdef1234567890abcdef12345678}" \
-e "OBJECTTYPEN=http://objecttypen:8000" \
-e "OBJECTTYPEN_TOKEN=${OBJECTTYPEN_TOKEN:-0123456789abcdef0123456789abcdef01234567}" \
-e "NRC_BASE=http://$nrc_ip:8000" \
-e "SINK_CALLBACK=http://$sink_ip:9000/" -e "SINK_AUTH=$SINK_AUTH" \
python:3-slim python /driver.py)"
docker cp "$here/objecten-notifications-check.py" "$drv:/driver.py" >/dev/null
docker start -a "$drv"
object_url="$(docker logs rr-overify 2>/dev/null | sed -n 's/^OBJECT_URL //p' | head -1)"
docker rm -f rr-overify >/dev/null
[ -n "$object_url" ] || { echo "FAIL — the driver did not write a RegisterRecord" >&2; exit 1; }
echo ">> wrote $object_url"
# Correlate on the object URL: a notification carries hoofdObject/resourceUrl, never the record
# data, so the reference inside the record is not in the delivered message.
echo ">> waiting for the notification to reach the sink"
for _ in $(seq 1 "${NOTIFICATION_TRIES:-40}"); do
if docker logs rr-osink 2>&1 | grep -qF "$object_url"; then
echo "OK — Objecten published to NRC and the abonnement delivered it:"
docker logs rr-osink 2>&1 | grep -F "$object_url" | tail -1 | cut -c1-500
exit 0
fi
sleep 2
done
echo "FAIL — no 'objecten' notification for $object_url reached the sink." >&2
echo " Objecten accepted the write, so the gap is downstream: the celery broker/worker," >&2
echo " the kanaal registration, or Objecten's notifications_config." >&2
echo "--- sink log ---" >&2; docker logs rr-osink 2>&1 | tail -8 >&2
echo "--- objecten log ---" >&2; docker logs "$obj" 2>&1 | tail -15 >&2
exit 1
+18 -50
View File
@@ -1,26 +1,18 @@
#!/usr/bin/env bash #!/usr/bin/env bash
# #
# Verify the end-to-end read-projection path (S-06, re-sourced by S-19b-2) against an ALREADY-RUNNING # Verify the end-to-end read-projection path (S-06) against an ALREADY-RUNNING full stack:
# full stack: ACL → Objecten → NRC → Event Subscriber → projection → projection-api. Seeds a # OpenZaak → NRC → Event Subscriber → projection → projection-api. Seeds a published BIG
# published BIG zaaktype (idempotent), registers an abonnement on the `objecten` kanaal pointing at # zaaktype (idempotent), registers an abonnement on the `zaken` kanaal pointing at the real
# the real Event Subscriber's /notifications callback (with the bearer it enforces), opens a zaak # Event Subscriber's /notifications callback (with the bearer it enforces), creates a zaak,
# *through the ACL*, and asserts projection-api serves a row for it with status INGEDIEND. # and asserts projection-api serves a row for that zaak with status INGEDIEND.
#
# The zaak is opened through the ACL, not straight against OpenZaak: since ADR-0030 the projection is
# derived from the RegisterRecord in Objecten, and the ACL is what writes that record (INGEDIEND on
# submit). A zaak created behind the ACL's back produces no register write and so no projection row —
# which is the point of the re-source.
# #
# All in-network, reaching services by container IP — single-label hosts aren't URL-valid and # All in-network, reaching services by container IP — single-label hosts aren't URL-valid and
# the runner can't reach published ports (gitea-actions-gotchas.md §5/§6). Does not own the stack # the runner can't reach published ports (gitea-actions-gotchas.md §5/§6). Reuses the
# lifecycle (the caller brings it up and tears it down), but does recreate the `acl` service to # notification driver to register the abonnement + create the zaak. Does NOT manage the stack
# repoint it — see below, and run-domain-check.sh, which does the same. Plain docker primitives only. # lifecycle (the caller owns bring-up + teardown). Plain docker primitives only. See ADR-0007/0008.
# See ADR-0007/0008/0030.
set -euo pipefail set -euo pipefail
here="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" here="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
root="$(cd "$here/.." && pwd)"
compose="$root/infra/docker-compose.yml"
WEBHOOK_AUTH="${NOTIFICATION_WEBHOOK_TOKEN:-Bearer big-reference-notifications}" WEBHOOK_AUTH="${NOTIFICATION_WEBHOOK_TOKEN:-Bearer big-reference-notifications}"
cleanup() { docker rm -f rr-pverify rr-pquery >/dev/null 2>&1 || true; } cleanup() { docker rm -f rr-pverify rr-pquery >/dev/null 2>&1 || true; }
@@ -32,13 +24,11 @@ oz="$(docker ps -q --filter 'name=[-_]openzaak[-_]' | head -1)"
nrc="$(docker ps -q --filter 'name=nrc-web' | head -1)" nrc="$(docker ps -q --filter 'name=nrc-web' | head -1)"
es="$(docker ps -q --filter 'name=event-subscriber' | head -1)" es="$(docker ps -q --filter 'name=event-subscriber' | head -1)"
proj="$(docker ps -q --filter 'name=projection-api' | head -1)" proj="$(docker ps -q --filter 'name=projection-api' | head -1)"
acl="$(docker ps -q --filter 'name=[-_]acl[-_]' | head -1)"
[ -n "$oz" ] && [ -n "$nrc" ] || { echo "ERROR: OpenZaak and/or NRC not running — bring the stack up first" >&2; exit 1; } [ -n "$oz" ] && [ -n "$nrc" ] || { echo "ERROR: OpenZaak and/or NRC not running — bring the stack up first" >&2; exit 1; }
[ -n "$es" ] && [ -n "$proj" ] || { echo "ERROR: event-subscriber and/or projection-api not running — bring the stack up first" >&2; exit 1; } [ -n "$es" ] && [ -n "$proj" ] || { echo "ERROR: event-subscriber and/or projection-api not running — bring the stack up first" >&2; exit 1; }
[ -n "$acl" ] || { echo "ERROR: acl not running — bring the stack up first" >&2; exit 1; }
net="$(docker inspect -f '{{range $k,$_ := .NetworkSettings.Networks}}{{$k}}{{"\n"}}{{end}}' "$oz" | head -1)" net="$(docker inspect -f '{{range $k,$_ := .NetworkSettings.Networks}}{{$k}}{{"\n"}}{{end}}' "$oz" | head -1)"
oz_ip="$(ip "$oz")"; nrc_ip="$(ip "$nrc")"; es_ip="$(ip "$es")"; proj_ip="$(ip "$proj")"; acl_ip="$(ip "$acl")" oz_ip="$(ip "$oz")"; nrc_ip="$(ip "$nrc")"; es_ip="$(ip "$es")"; proj_ip="$(ip "$proj")"
echo ">> network=$net openzaak=$oz_ip nrc=$nrc_ip event-subscriber=$es_ip projection-api=$proj_ip acl=$acl_ip" echo ">> network=$net openzaak=$oz_ip nrc=$nrc_ip event-subscriber=$es_ip projection-api=$proj_ip"
echo ">> seeding a published BIG zaaktype (idempotent)" echo ">> seeding a published BIG zaaktype (idempotent)"
sid="$(docker create --network "$net" -e "OZ_BASE=http://$oz_ip:8000" -e OZ_PUBLISH=1 \ sid="$(docker create --network "$net" -e "OZ_BASE=http://$oz_ip:8000" -e OZ_PUBLISH=1 \
@@ -47,39 +37,19 @@ docker cp "$here/openzaak/seed_catalogus.py" "$sid:/seed.py" >/dev/null
docker start -a "$sid" docker start -a "$sid"
docker rm -f "$sid" >/dev/null docker rm -f "$sid" >/dev/null
echo ">> registering the event-subscriber abonnement on the objecten kanaal" echo ">> registering abonnement at the Event Subscriber + creating a zaak"
docker rm -f rr-pverify >/dev/null 2>&1 || true docker rm -f rr-pverify >/dev/null 2>&1 || true
# The same script the local stack uses (ADR-0020), so both paths register the identical abonnement.
drv="$(docker create --network "$net" --name rr-pverify \ drv="$(docker create --network "$net" --name rr-pverify \
-e "NRC_BASE=http://$nrc_ip:8000" \ -e "OZ_BASE=http://$oz_ip:8000" -e "NRC_BASE=http://$nrc_ip:8000" \
-e "SINK_HOST=$es_ip" -e "SINK_PORT=8080" -e "SINK_AUTH=$WEBHOOK_AUTH" \ -e "SINK_CALLBACK=http://$es_ip:8080/notifications" -e "SINK_AUTH=$WEBHOOK_AUTH" \
python:3-slim python /subscribe.py)" python:3-slim python /driver.py)"
docker cp "$here/local/register-abonnement.py" "$drv:/subscribe.py" >/dev/null docker cp "$here/verify-notification-driver.py" "$drv:/driver.py" >/dev/null
docker start -a "$drv" docker start -a "$drv"
zaak_url="$(docker logs rr-pverify 2>/dev/null | sed -n 's/^ZAAK_CREATED //p' | head -1)"
docker rm -f rr-pverify >/dev/null docker rm -f rr-pverify >/dev/null
[ -n "$zaak_url" ] || { echo "ERROR: driver did not create a zaak" >&2; exit 1; }
# OpenZaak reflects the request Host into the zaaktype `url` it returns, and then rejects that same
# URL on zaak-create when the host is single-label ("Voer een geldige URL in."). The stack's ACL is
# configured with `http://openzaak:8000/`, so it must be repointed at OpenZaak's container IP before
# it can open a zaak — exactly what run-domain-check.sh does, and the same class of constraint as the
# `objecten.local` alias (ADR-0029). The ACL resolves the zaaktype itself (S-27, ADR-0021), so the
# base URL is the only thing to inject.
echo ">> recreating the acl service pointed at OpenZaak's IP"
ACL_OPENZAAK_BASEURL="http://$oz_ip:8000/" docker compose -f "$compose" up -d acl
WAIT_TIMEOUT="${WAIT_TIMEOUT:-120}" bash "$here/wait-healthy.sh" acl
# The container is replaced, so its IP may have changed.
acl="$(docker ps -q --filter 'name=[-_]acl[-_]' | head -1)"
acl_ip="$(ip "$acl")"
echo ">> opening a zaak through the ACL (which writes the INGEDIEND register record)"
reference="PROJ-$(date +%s)"
zaak_url="$(docker run --rm --network "$net" curlimages/curl:latest \
-fsS -X POST "http://$acl_ip:8080/zaken" -H 'Content-Type: application/json' \
-d "{\"bsn\":\"123456782\",\"reference\":\"$reference\"}" \
| sed -n 's/.*"zaakUrl":"\([^"]*\)".*/\1/p')"
[ -n "$zaak_url" ] || { echo "ERROR: the ACL did not open a zaak" >&2; exit 1; }
zaak_uuid="${zaak_url##*/}" zaak_uuid="${zaak_url##*/}"
echo ">> zaak created: $zaak_url (reference $reference)" echo ">> zaak created: $zaak_url"
echo ">> polling projection-api for the projected row (status INGEDIEND)" echo ">> polling projection-api for the projected row (status INGEDIEND)"
for _ in $(seq 1 30); do for _ in $(seq 1 30); do
@@ -93,8 +63,6 @@ for _ in $(seq 1 30); do
sleep 2 sleep 2
done done
echo "FAIL — projection-api never served an INGEDIEND row for zaak $zaak_uuid" >&2 echo "FAIL — projection-api never served an INGEDIEND row for zaak $zaak_uuid" >&2
echo " The chain is ACL → Objecten → NRC → event-subscriber → projection (ADR-0030)." >&2
echo "--- event-subscriber log ---" >&2; docker logs "$es" 2>&1 | tail -10 >&2 echo "--- event-subscriber log ---" >&2; docker logs "$es" 2>&1 | tail -10 >&2
echo "--- projection-api log ---" >&2; docker logs "$proj" 2>&1 | tail -10 >&2 echo "--- projection-api log ---" >&2; docker logs "$proj" 2>&1 | tail -10 >&2
echo "--- acl log ---" >&2; docker logs "$acl" 2>&1 | tail -10 >&2
exit 1 exit 1
+3 -7
View File
@@ -15,13 +15,9 @@ set -euo pipefail
timeout="${WAIT_TIMEOUT:-420}" timeout="${WAIT_TIMEOUT:-420}"
deadline=$(( $(date +%s) + timeout )) deadline=$(( $(date +%s) + timeout ))
# compose service name -> container id. `--filter name=` is a substring match, so it is anchored on # compose service name -> container id. The name filter matches both docker
# the compose replica suffix — otherwise 'objecten' also matches objecten-db / objecten-redis / # compose ("infra-openzaak-1") and podman-compose ("infra_openzaak_1") naming.
# objecten-celery, and 'objecttypen' matches objecttypen-db. Whichever docker listed first won, so a cid_for() { docker ps -aq --filter "name=$1" | head -1; }
# service with a sibling that has no healthcheck timed out with status=none while it was in fact
# healthy. The pattern matches both docker compose ("infra-objecten-1") and podman-compose
# ("infra_objecten_1") naming; the same anchoring the verify check scripts use.
cid_for() { docker ps -aq --filter "name=$1[-_][0-9]+\$" | head -1; }
for svc in "$@"; do for svc in "$@"; do
echo "waiting for '$svc' to be healthy (timeout ${timeout}s)..." echo "waiting for '$svc' to be healthy (timeout ${timeout}s)..."
-13
View File
@@ -90,16 +90,6 @@ app.MapPost("/zaken/reference", async (ZaakReferenceRequest body, AclService acl
return Results.Ok(new { reference }); return Results.Ok(new { reference });
}); });
// Read the register record an object in Objecten holds. The Event Subscriber projects a register
// write from the notification NRC delivers, which carries only the object URL, and may not talk to
// Objecten itself (§8.1, ADR-0028/ADR-0030). 404 when the object holds no record — the subscriber
// treats that as "nothing to project" rather than an error (§8.6).
app.MapPost("/register-records/read", async (RegisterRecordReadRequest body, AclService acl, CancellationToken ct) =>
{
var record = await acl.GetRegisterRecordAsync(new Uri(body.ObjectUrl), ct);
return record is null ? Results.NotFound() : Results.Ok(record);
});
// Store an uploaded diploma against a zaak (S-10b): the domain sends the file as base64; the ACL // Store an uploaded diploma against a zaak (S-10b): the domain sends the file as base64; the ACL
// creates the ZGW enkelvoudiginformatieobject and relates it to the zaak (§8.1). Returns its URL. // creates the ZGW enkelvoudiginformatieobject and relates it to the zaak (§8.1). Returns its URL.
app.MapPost("/documenten", async (StoreDocumentRequest body, AclService acl, CancellationToken ct) => app.MapPost("/documenten", async (StoreDocumentRequest body, AclService acl, CancellationToken ct) =>
@@ -141,9 +131,6 @@ public sealed record CancelZaakRequest(string ZaakUrl);
public sealed record ZaakReferenceRequest(string ZaakUrl); public sealed record ZaakReferenceRequest(string ZaakUrl);
/// <summary>The object whose register record the Event Subscriber wants read back (S-19b-2).</summary>
public sealed record RegisterRecordReadRequest(string ObjectUrl);
public sealed record StoreDocumentRequest(string ZaakUrl, string ContentBase64, string FileName, string ContentType); public sealed record StoreDocumentRequest(string ZaakUrl, string ContentBase64, string FileName, string ContentType);
public partial class Program; public partial class Program;
+1 -22
View File
@@ -24,16 +24,7 @@ public sealed class AclService(
clock.Today, clock.Today,
registration.Reference); registration.Reference);
var zaakUrl = await gateway.OpenZaakAsync(request, ct); return await gateway.OpenZaakAsync(request, ct);
// The register — not ZGW — is what the read projection is sourced from (ADR-0028/ADR-0030),
// so the record exists from submission, not only from approval. Same two-writes-converging
// posture as ApproveZaakAsync: the upsert is keyed on the zaak id, so a retried submit
// updates the record rather than adding a second one (§8.6).
await register.UpsertAsync(
new RegisterRecord(ZaakId(zaakUrl), RegisterRecordStatus.Ingediend, registration.Reference), ct);
return zaakUrl;
} }
/// <summary> /// <summary>
@@ -61,18 +52,6 @@ public sealed class AclService(
ct); ct);
} }
/// <summary>
/// The register record held by an object in Objecten, for the Event Subscriber (S-19b-2). The
/// subscriber gets only an object URL on the notification and may not read Objecten itself
/// (§8.1, ADR-0028), so the ACL reads it back.
/// </summary>
public Task<RegisterRecord?> GetRegisterRecordAsync(Uri objectUrl, CancellationToken ct = default)
{
ArgumentNullException.ThrowIfNull(objectUrl);
return register.GetAsync(objectUrl, ct);
}
/// <summary>The zaak's UUID — the key the register record and the read projection rows share.</summary> /// <summary>The zaak's UUID — the key the register record and the read projection rows share.</summary>
private static string ZaakId(Uri zaakUrl) => zaakUrl.Segments[^1].TrimEnd('/'); private static string ZaakId(Uri zaakUrl) => zaakUrl.Segments[^1].TrimEnd('/');
@@ -13,14 +13,6 @@ public interface IRegisterRecordGateway
/// the existing object instead of creating a second one (§8.6). /// the existing object instead of creating a second one (§8.6).
/// </summary> /// </summary>
Task UpsertAsync(RegisterRecord record, CancellationToken ct = default); Task UpsertAsync(RegisterRecord record, CancellationToken ct = default);
/// <summary>
/// The register record held by the object at <paramref name="objectUrl"/>, or <c>null</c> if that
/// object holds none. The Event Subscriber projects a register write from the notification NRC
/// delivers, which carries only the object URL — so it reads the record back through the ACL
/// rather than talking to Objecten itself (§8.1, S-19b-2).
/// </summary>
Task<RegisterRecord?> GetAsync(Uri objectUrl, CancellationToken ct = default);
} }
/// <summary> /// <summary>
@@ -1,4 +1,3 @@
using System.Net;
using System.Net.Http.Headers; using System.Net.Http.Headers;
using System.Net.Http.Json; using System.Net.Http.Json;
using System.Text.Json.Serialization; using System.Text.Json.Serialization;
@@ -39,30 +38,6 @@ public sealed class ObjectenGateway(HttpClient http, ObjectenOptions options, IC
"Updating the register record", ct); "Updating the register record", ct);
} }
public async Task<RegisterRecord?> GetAsync(Uri objectUrl, CancellationToken ct = default)
{
ArgumentNullException.ThrowIfNull(objectUrl);
// Fetched by the URL the notification carried, so no objecttype resolution and no search —
// unlike a write, which has to find the object for a registration id.
using var message = new HttpRequestMessage(HttpMethod.Get, objectUrl);
message.Headers.Authorization = new AuthenticationHeaderValue("Token", options.Token);
message.Headers.Add("Accept-Crs", "EPSG:4326");
using var response = await http.SendAsync(message, ct);
// The object may be gone by the time a (possibly redelivered) notification is handled —
// there is simply nothing to project, which is not a failure (§8.6).
if (response.StatusCode == HttpStatusCode.NotFound)
return null;
await EnsureSuccessAsync(response, "Reading the register record", ct);
var body = await response.Content.ReadFromJsonAsync<ReadObjectDto>(ct)
?? throw new InvalidOperationException("Objecten returned an empty object response");
var data = body.Record?.Data;
return data is null ? null : new RegisterRecord(data.Id, data.Status, data.Reference);
}
private RecordDto NewRecord(int typeVersion, RecordDataDto data) => private RecordDto NewRecord(int typeVersion, RecordDataDto data) =>
new(typeVersion, data, clock.Today.ToString("yyyy-MM-dd")); new(typeVersion, data, clock.Today.ToString("yyyy-MM-dd"));
@@ -166,12 +141,6 @@ public sealed class ObjectenGateway(HttpClient http, ObjectenOptions options, IC
private sealed record ObjectDto( private sealed record ObjectDto(
[property: JsonPropertyName("url")] string Url); [property: JsonPropertyName("url")] string Url);
private sealed record ReadObjectDto(
[property: JsonPropertyName("record")] ReadRecordDto? Record);
private sealed record ReadRecordDto(
[property: JsonPropertyName("data")] RecordDataDto? Data);
private sealed record CreateObjectDto( private sealed record CreateObjectDto(
[property: JsonPropertyName("type")] string Type, [property: JsonPropertyName("type")] string Type,
[property: JsonPropertyName("record")] RecordDto Record); [property: JsonPropertyName("record")] RecordDto Record);
@@ -20,7 +20,7 @@ public sealed class ObjectenGatewayIntegrationTests
new HttpClient(), new HttpClient(),
new ObjectenOptions new ObjectenOptions
{ {
BaseUrl = new(Env("OBJECTEN_BASE", "http://objecten.local:8000")), BaseUrl = new(Env("OBJECTEN_BASE", "http://objecten:8000")),
Token = Env("OBJECTEN_TOKEN", "1234567890abcdef1234567890abcdef12345678"), Token = Env("OBJECTEN_TOKEN", "1234567890abcdef1234567890abcdef12345678"),
ObjecttypenBaseUrl = new(Env("OBJECTTYPEN_BASE", "http://objecttypen:8000")), ObjecttypenBaseUrl = new(Env("OBJECTTYPEN_BASE", "http://objecttypen:8000")),
ObjecttypenToken = Env("OBJECTTYPEN_TOKEN", "0123456789abcdef0123456789abcdef01234567"), ObjecttypenToken = Env("OBJECTTYPEN_TOKEN", "0123456789abcdef0123456789abcdef01234567"),
@@ -65,7 +65,7 @@ public sealed class ObjectenGatewayIntegrationTests
{ {
using var http = new HttpClient(); using var http = new HttpClient();
var objecttype = await ResolveObjecttypeUrlAsync(http); var objecttype = await ResolveObjecttypeUrlAsync(http);
var query = new Uri(new Uri(Env("OBJECTEN_BASE", "http://objecten.local:8000")), var query = new Uri(new Uri(Env("OBJECTEN_BASE", "http://objecten:8000")),
"/api/v2/objects?type=" + Uri.EscapeDataString(objecttype) + "/api/v2/objects?type=" + Uri.EscapeDataString(objecttype) +
"&data_attrs=id__exact__" + Uri.EscapeDataString(id)); "&data_attrs=id__exact__" + Uri.EscapeDataString(id));
-56
View File
@@ -79,21 +79,11 @@ public class AclServiceTests
{ {
public readonly List<RegisterRecord> Upserted = []; public readonly List<RegisterRecord> Upserted = [];
public RegisterRecord? Stored;
public Uri? ReadFrom;
public Task UpsertAsync(RegisterRecord record, CancellationToken ct = default) public Task UpsertAsync(RegisterRecord record, CancellationToken ct = default)
{ {
Upserted.Add(record); Upserted.Add(record);
return Task.CompletedTask; return Task.CompletedTask;
} }
public Task<RegisterRecord?> GetAsync(Uri objectUrl, CancellationToken ct = default)
{
ReadFrom = objectUrl;
return Task.FromResult(Stored);
}
} }
private static AclDefaults Defaults() => new() private static AclDefaults Defaults() => new()
@@ -140,52 +130,6 @@ public class AclServiceTests
Assert.Equal("reg-77", req.Identificatie); Assert.Equal("reg-77", req.Identificatie);
} }
[Fact]
public async Task Opening_a_zaak_also_writes_an_ingediend_register_record(/* S-19b-2 */)
{
var gateway = new FakeGateway();
var register = new FakeRegisterRecordGateway();
var service = ServiceWith(gateway, register, Defaults(), new DateOnly(2026, 6, 4));
await service.OpenZaakAsync(new DomainRegistration("123456782", "reg-77"));
// The register — not ZGW — is what the read projection is sourced from (ADR-0028), so a
// submitted registration has to exist there the moment the zaak is opened, not only on
// approval. Approval upserts this same record to INGESCHREVEN.
var record = Assert.Single(register.Upserted);
Assert.Equal("abc", record.Id);
Assert.Equal("INGEDIEND", record.Status);
// The reference comes from the registration itself — no ZGW read-back needed on this path.
Assert.Equal("reg-77", record.Reference);
}
[Fact]
public async Task Reading_a_register_record_goes_through_the_objecten_gateway(/* S-19b-2 */)
{
var gateway = new FakeGateway();
var register = new FakeRegisterRecordGateway { Stored = new RegisterRecord("abc", "INGESCHREVEN", "reg-77") };
var service = ServiceWith(gateway, register, Defaults(), new DateOnly(2026, 6, 4));
var objectUrl = new Uri("http://objecten.local:8000/api/v2/objects/9de4a2ca");
var record = await service.GetRegisterRecordAsync(objectUrl);
Assert.Equal(objectUrl, register.ReadFrom);
Assert.Equal("abc", record!.Id);
Assert.Equal("INGESCHREVEN", record.Status);
Assert.Equal("reg-77", record.Reference);
}
[Fact]
public async Task Reading_a_register_record_from_a_null_url_is_rejected(/* S-19b-2 */)
{
var gateway = new FakeGateway();
var register = new FakeRegisterRecordGateway();
var service = ServiceWith(gateway, register, Defaults(), new DateOnly(2026, 6, 4));
await Assert.ThrowsAsync<ArgumentNullException>(() => service.GetRegisterRecordAsync(null!));
Assert.Null(register.ReadFrom);
}
[Fact] [Fact]
public async Task Opening_a_zaak_reflects_a_default_fill_update(/* S-15b */) public async Task Opening_a_zaak_reflects_a_default_fill_update(/* S-15b */)
{ {
@@ -83,43 +83,6 @@ public class ObjectenGatewayTests
private static RegisterRecord Record() => new("zaak-uuid-1", RegisterRecordStatus.Ingeschreven, "REG-2026-0001"); private static RegisterRecord Record() => new("zaak-uuid-1", RegisterRecordStatus.Ingeschreven, "REG-2026-0001");
[Fact]
public async Task Reads_a_register_record_back_from_its_object_url(/* S-19b-2 */)
{
var sent = new List<Sent>();
var objectUrl = new Uri("http://objecten:8000/api/v2/objects/obj-9");
var gateway = Gateway(sent, _ => Json(new
{
url = objectUrl.ToString(),
record = new { data = new { id = "zaak-uuid-1", status = "INGESCHREVEN", reference = "REG-2026-0001" } },
}));
var record = await gateway.GetAsync(objectUrl);
// The object is fetched directly by the URL the notification carried — no objecttype
// resolution and no search, unlike a write.
var read = Assert.Single(sent);
Assert.Equal(HttpMethod.Get, read.Method);
Assert.Equal(objectUrl, read.Uri);
// Objecten is a geo API: the CRS header is required on reads too.
Assert.Equal("EPSG:4326", read.AcceptCrs);
Assert.Equal("Token objecten-token", read.Auth);
Assert.Equal("zaak-uuid-1", record!.Id);
Assert.Equal("INGESCHREVEN", record.Status);
Assert.Equal("REG-2026-0001", record.Reference);
}
[Fact]
public async Task Reading_an_object_that_is_gone_yields_no_record(/* S-19b-2 */)
{
var sent = new List<Sent>();
var gateway = Gateway(sent, _ => new HttpResponseMessage(HttpStatusCode.NotFound));
// A record deleted between the notification and the read is not an error — there is simply
// nothing to project (§8.6: the subscriber tolerates whatever order deliveries arrive in).
Assert.Null(await gateway.GetAsync(new Uri("http://objecten:8000/api/v2/objects/gone")));
}
[Fact] [Fact]
public async Task Creates_the_object_when_none_exists_for_the_registration() public async Task Creates_the_object_when_none_exists_for_the_registration()
{ {
@@ -1,4 +1,3 @@
using System.Net;
using System.Net.Http.Json; using System.Net.Http.Json;
using System.Text.Json.Serialization; using System.Text.Json.Serialization;
using EventSubscriber.Application; using EventSubscriber.Application;
@@ -6,28 +5,26 @@ using EventSubscriber.Application;
namespace EventSubscriber.Api; namespace EventSubscriber.Api;
/// <summary> /// <summary>
/// HTTP client to the ACL service. An Objecten notification carries only the object URL, so the /// HTTP client to the ACL service. The subscriber enriches the projection with the zaak's reference
/// subscriber reads the register record back through the ACL — the only code that may talk to /// (identificatie) by asking the ACL — the only code that may read ZGW (§8.1) — rather than reading
/// Objecten (§8.1, ADR-0028/ADR-0030) — rather than reading Objecten itself. /// OpenZaak itself (adr-proposal #78).
/// </summary> /// </summary>
public sealed class AclHttpClient(HttpClient http) : IAclClient public sealed class AclHttpClient(HttpClient http) : IAclClient
{ {
public async Task<RegisterRecord?> GetRegisterRecordAsync(Uri objectUrl, CancellationToken ct = default) public async Task<string> GetZaakReferenceAsync(Uri zaakUrl, CancellationToken ct = default)
{ {
ArgumentNullException.ThrowIfNull(objectUrl); ArgumentNullException.ThrowIfNull(zaakUrl);
using var response = await http.PostAsJsonAsync( using var response = await http.PostAsJsonAsync(
new Uri(http.BaseAddress!, "register-records/read"), new Uri(http.BaseAddress!, "zaken/reference"), new ReferenceRequest(zaakUrl.ToString()), ct);
new ReadRequest(objectUrl.ToString()), ct);
// The object holds no register record (deleted, or never one) — nothing to project (§8.6).
if (response.StatusCode == HttpStatusCode.NotFound)
return null;
response.EnsureSuccessStatusCode(); response.EnsureSuccessStatusCode();
return await response.Content.ReadFromJsonAsync<RegisterRecord>(ct)
?? throw new InvalidOperationException("The ACL returned an empty register record response."); var body = await response.Content.ReadFromJsonAsync<ReferenceResponse>(ct)
?? throw new InvalidOperationException("The ACL returned an empty reference response.");
return body.Reference;
} }
private sealed record ReadRequest([property: JsonPropertyName("objectUrl")] string ObjectUrl); private sealed record ReferenceRequest([property: JsonPropertyName("zaakUrl")] string ZaakUrl);
private sealed record ReferenceResponse([property: JsonPropertyName("reference")] string Reference);
} }
@@ -84,12 +84,11 @@ app.MapPost("/admin/rebuild", async (NotificationProjector projector, Cancellati
await app.RunAsync(); await app.RunAsync();
/// <summary>The NRC notification body, as Open Notificaties POSTs it. Only the fields the projector /// <summary>The NRC notification body, as Open Notificaties POSTs it. Only the fields the
/// needs are bound; <c>aanmaakdatum</c>, <c>kenmerken</c> and <c>hoofdObject</c> are ignored for a /// projection needs are bound; <c>aanmaakdatum</c>/<c>kenmerken</c> are ignored for the minimal slice.</summary>
/// register write hoofdObject is the same object as resourceUrl (ADR-0030).</summary> public sealed record NotificationDto(string Kanaal, string Resource, string Actie, Uri ResourceUrl, Uri? HoofdObject = null)
public sealed record NotificationDto(string Kanaal, string Resource, string Actie, Uri ResourceUrl)
{ {
public Notification ToNotification() => new(Kanaal, Resource, Actie, ResourceUrl); public Notification ToNotification() => new(Kanaal, Resource, Actie, ResourceUrl, HoofdObject);
} }
public partial class Program public partial class Program
@@ -2,39 +2,40 @@ namespace EventSubscriber.Application;
/// <summary> /// <summary>
/// An inbound NRC (Open Notificaties) notification, as Open Notificaties POSTs it to an /// An inbound NRC (Open Notificaties) notification, as Open Notificaties POSTs it to an
/// abonnement callback. Only the fields the projection needs are modelled. /// abonnement callback. Only the fields the projection needs are modelled; the full ZGW
/// "Notificatie" resource also carries <c>aanmaakdatum</c> and <c>kenmerken</c> which the
/// minimal projection ignores (bsn is deferred — see ADR-0008). For a <c>zaken</c>/<c>zaak</c>/<c>create</c>
/// notification <c>hoofdObject</c> and <c>resourceUrl</c> are both the created zaak's URL.
/// </summary> /// </summary>
/// <remarks>
/// Since S-19b-2 the subscriber listens on the <c>objecten</c> kanaal, not <c>zaken</c>: the
/// register record in Objecten is what the projection is derived from (ADR-0030), so the
/// projection is a cache of the register rather than a re-derivation of the case system. An
/// Objecten notification carries <b>no record data</b> — only the object URL (as both
/// <c>hoofdObject</c> and <c>resourceUrl</c>) and the objecttype as a kenmerk — so the record
/// itself is read back through the ACL.
/// </remarks>
public sealed record Notification( public sealed record Notification(
string Kanaal, string Kanaal,
string Resource, string Resource,
string Actie, string Actie,
Uri ResourceUrl) Uri ResourceUrl,
Uri? HoofdObject = null)
{ {
/// <summary> /// <summary>A zaak being created — projected as INGEDIEND.</summary>
/// A register record written to Objecten — <c>create</c> on submit and <c>partial_update</c> on public bool IsZaakCreated =>
/// approval, since the ACL upserts the same object for a registration (§8.6). Kanaal == "zaken" && Resource == "zaak" && Actie == "create";
/// </summary>
/// <remarks>
/// <c>partial_update</c> is what a PATCH actually reports: DRF routes it through the notifying
/// <c>update()</c> but names the action <c>partial_update</c>, and that is what Objecten puts in
/// the notification. <c>update</c> is accepted too, so a PUT-shaped write would project the same
/// way. <c>destroy</c> is deliberately not: removing a registration from the public register is
/// its own decision, not a side effect of this one.
/// </remarks>
public bool IsRegisterRecordWritten =>
Kanaal == "objecten" && Resource == "object"
&& Actie is "create" or "update" or "partial_update";
/// <summary>The object holding the register record. For a <c>resource: object</c> notification /// <summary>A status being set on a zaak — the approval, projected as INGESCHREVEN (S-09b). In the
/// Objecten sends the object as both <c>hoofdObject</c> and <c>resourceUrl</c> — the object is /// walking skeleton the only status ever set after creation is the approval, and the subscriber may
/// the main resource — so the notification's own <c>hoofdObject</c> is not modelled.</summary> /// not read OpenZaak (§8.1), so any status-create is taken as the approval.</summary>
public Uri ObjectUrl => ResourceUrl; public bool IsZaakStatusSet =>
Kanaal == "zaken" && Resource == "status" && Actie == "create";
/// <summary>The zaak URL this notification concerns — <c>hoofdObject</c> (the zaak) for a status
/// notification, else the resource URL (which, for a zaak-create, is the zaak).</summary>
public Uri ZaakUrl => HoofdObject ?? ResourceUrl;
/// <summary>The zaak UUID used as the projection key — the trailing segment of <see cref="ZaakUrl"/>.</summary>
public string ZaakId => ZaakUrl.Segments[^1].Trim('/');
/// <summary>
/// A deterministic dedup key. Open Notificaties carries no notification id and may
/// redeliver, so the key is derived from the immutable notification content: two
/// deliveries of the same zaak-create collapse to one. (NRC may also deliver
/// out of order; the projector tolerates that — order does not change the outcome.)
/// </summary>
public string IdempotencyKey => $"{Kanaal}:{Resource}:{Actie}:{ResourceUrl}";
} }
@@ -3,27 +3,21 @@ namespace EventSubscriber.Application;
/// <summary> /// <summary>
/// Projects inbound NRC notifications into the read projection. Tolerates duplicate and /// Projects inbound NRC notifications into the read projection. Tolerates duplicate and
/// out-of-order deliveries (CLAUDE.md §8.6): the notification log dedups, and the projection /// out-of-order deliveries (CLAUDE.md §8.6): the notification log dedups, and the projection
/// upsert is idempotent on the register id. Rebuilds the projection by replaying the log. /// upsert is idempotent on the zaak id. Rebuilds the projection by replaying the log.
/// </summary> /// </summary>
public sealed class NotificationProjector(INotificationLog log, IProjectionStore store, IAclClient acl) public sealed class NotificationProjector(INotificationLog log, IProjectionStore store, IAclClient acl)
{ {
/// <summary>Handle one inbound notification. Reacts to a register record being written to /// <summary>Handle one inbound notification. Reacts to a zaak being created (INGEDIEND) and a
/// Objecten (S-19b-2, ADR-0030) and ignores everything else. The notification carries only the /// status being set (INGESCHREVEN); ignores everything else. Enriches the row with the zaak's
/// object URL, so the record is read back through the ACL (§8.1) and becomes the row verbatim.</summary> /// reference via the ACL (§8.1) and records it so a rebuild needs no ZGW access (#78).</summary>
public async Task HandleAsync(Notification notification, CancellationToken ct = default) public async Task HandleAsync(Notification notification, CancellationToken ct = default)
{ {
ArgumentNullException.ThrowIfNull(notification); if (!notification.IsZaakCreated && !notification.IsZaakStatusSet)
if (!notification.IsRegisterRecordWritten)
return;
var record = await acl.GetRegisterRecordAsync(notification.ObjectUrl, ct);
// The object is gone, or holds no register record — nothing to project (§8.6).
if (record is null)
return; return;
var reference = await acl.GetZaakReferenceAsync(notification.ZaakUrl, ct);
var recorded = new RecordedNotification( var recorded = new RecordedNotification(
KeyFor(notification.ObjectUrl, record), record.Id, record.Status, record.Reference); notification.IdempotencyKey, notification.Actie, notification.ZaakId, notification.Resource, reference);
// Atomic record-or-skip: a duplicate (or concurrent) delivery is recognised and dropped // Atomic record-or-skip: a duplicate (or concurrent) delivery is recognised and dropped
// before it touches the projection, so the projection stays a faithful derived artefact. // before it touches the projection, so the projection stays a faithful derived artefact.
@@ -33,20 +27,6 @@ public sealed class NotificationProjector(INotificationLog log, IProjectionStore
await store.UpsertAsync(ToEntry(recorded), ct); await store.UpsertAsync(ToEntry(recorded), ct);
} }
/// <summary>
/// A deterministic dedup key: the object, plus the state that write puts in the projection.
/// </summary>
/// <remarks>
/// Open Notificaties carries no notification id and may redeliver, so the key is derived from
/// content. It cannot be the object URL alone — the ACL upserts one object per registration, so
/// submit and approval both notify about the *same* URL and the approval would be swallowed as a
/// duplicate. Nor can it include the actie: a retried approval would be a second `update`. Keying
/// on the projected row means a redelivery collapses and a genuine state change does not, which
/// is exactly the property §8.6 asks for.
/// </remarks>
private static string KeyFor(Uri objectUrl, RegisterRecord record)
=> $"objecten:object:{objectUrl}:{record.Status}:{record.Reference}";
/// <summary>Rebuild the projection from the durable notification log (PRD §8.4).</summary> /// <summary>Rebuild the projection from the durable notification log (PRD §8.4).</summary>
public async Task RebuildAsync(CancellationToken ct = default) public async Task RebuildAsync(CancellationToken ct = default)
{ {
@@ -55,9 +35,11 @@ public sealed class NotificationProjector(INotificationLog log, IProjectionStore
await store.UpsertAsync(ToEntry(recorded), ct); await store.UpsertAsync(ToEntry(recorded), ct);
} }
/// <summary>The projection row for an accepted notification. The log already holds exactly the /// <summary>The projection row for an accepted notification: a status-set maps to INGESCHREVEN,
/// row's fields, so a rebuild needs no mapping rules and no upstream reads. bsn/naam stay /// a zaak-create to INGEDIEND. bsn/naam are deferred (ADR-0008).</summary>
/// deferred — the register record is public-safe by construction (ADR-0027).</summary>
private static RegisterEntry ToEntry(RecordedNotification recorded) private static RegisterEntry ToEntry(RecordedNotification recorded)
=> new(recorded.RegisterId, recorded.Status, recorded.Reference); => new(
recorded.ZaakId,
recorded.Resource == "status" ? RegistrationStatus.Ingeschreven : RegistrationStatus.Ingediend,
Reference: recorded.Reference);
} }
@@ -4,7 +4,7 @@ namespace EventSubscriber.Application;
/// The durable log of notifications the subscriber has accepted. It is both the idempotency /// The durable log of notifications the subscriber has accepted. It is both the idempotency
/// guard (a replayed notification is recognised and dropped) and the rebuild source: the /// guard (a replayed notification is recognised and dropped) and the rebuild source: the
/// projection is a derived artefact (PRD §8.4) regenerated by replaying this log, so a rebuild /// projection is a derived artefact (PRD §8.4) regenerated by replaying this log, so a rebuild
/// needs no access to Objecten or ZGW (CLAUDE.md §8.1). Implemented in Infrastructure over Postgres. /// needs no access to OpenZaak (CLAUDE.md §8.1). Implemented in Infrastructure over Postgres.
/// </summary> /// </summary>
public interface INotificationLog public interface INotificationLog
{ {
@@ -19,29 +19,22 @@ public interface INotificationLog
Task<IReadOnlyList<RecordedNotification>> AllAsync(CancellationToken ct = default); Task<IReadOnlyList<RecordedNotification>> AllAsync(CancellationToken ct = default);
} }
/// <summary> /// <summary>A notification that has been accepted, retaining what a rebuild needs to recompute its
/// An accepted notification, retaining exactly the projection row it produced — so a rebuild /// projection row — the ZGW <c>resource</c> (zaak-create → INGEDIEND vs status-set → INGESCHREVEN) and
/// reproduces the row by replaying the log, without re-reading Objecten (S-19b-2, ADR-0030). /// the zaak <c>reference</c> (identificatie), so a rebuild reproduces the row without re-reading ZGW (#78).</summary>
/// </summary> public sealed record RecordedNotification(string Key, string Actie, string ZaakId, string Resource, string? Reference);
public sealed record RecordedNotification(string Key, string RegisterId, string Status, string? Reference);
/// <summary> /// <summary>
/// Port to the Anti-Corruption Layer. An Objecten notification carries only the object URL, so the /// Port to the Anti-Corruption Layer. The subscriber enriches the projection with the zaak's
/// subscriber reads the register record back through the ACL — the only code that may talk to /// public-safe reference (its identificatie) by asking the ACL — the only code that may read ZGW
/// Objecten (§8.1, ADR-0028) — rather than reading Objecten itself. /// (§8.1) — rather than reading OpenZaak itself (adr-proposal #78).
/// </summary> /// </summary>
public interface IAclClient public interface IAclClient
{ {
/// <summary>The register record the object at <paramref name="objectUrl"/> holds, or /// <summary>The zaak's reference (identificatie) for the read projection.</summary>
/// <c>null</c> if it holds none — the object may be gone by the time a redelivered Task<string> GetZaakReferenceAsync(Uri zaakUrl, CancellationToken ct = default);
/// notification is handled, which is not an error (§8.6).</summary>
Task<RegisterRecord?> GetRegisterRecordAsync(Uri objectUrl, CancellationToken ct = default);
} }
/// <summary>The public-safe register record as the ACL returns it — the RegisterRecord objecttype's
/// schema (ADR-0027). No bsn, no name: the register is world-readable.</summary>
public sealed record RegisterRecord(string Id, string Status, string? Reference);
/// <summary>The read projection store. Owned by the projection bounded context (ADR-0008); the /// <summary>The read projection store. Owned by the projection bounded context (ADR-0008); the
/// subscriber writes to it and the projection-api reads it.</summary> /// subscriber writes to it and the projection-api reads it.</summary>
public interface IProjectionStore public interface IProjectionStore
@@ -5,42 +5,27 @@ using EventSubscriber.Api;
namespace EventSubscriber.Tests; namespace EventSubscriber.Tests;
/// <summary> /// <summary>
/// Unit tests for the subscriber's ACL client, which reads a register record through the ACL — the /// Unit tests for the subscriber's ACL client, which reads a zaak's reference (identificatie) through
/// only code allowed to talk to Objecten (§8.1, ADR-0028/ADR-0030). Uses a scripted message handler /// the ACL — the only code allowed to talk to ZGW (§8.1, #78). Uses a scripted message handler so no
/// so no real ACL is required. /// real ACL is required.
/// </summary> /// </summary>
public class AclHttpClientTests public class AclHttpClientTests
{ {
private const string ObjectUrl = "http://objecten.local:8000/api/v2/objects/obj-9";
private static AclHttpClient Client(StubHandler handler) => private static AclHttpClient Client(StubHandler handler) =>
new(new HttpClient(handler) { BaseAddress = new Uri("http://acl/") }); new(new HttpClient(handler) { BaseAddress = new Uri("http://acl/") });
[Fact] [Fact]
public async Task Reads_a_register_record_by_posting_the_object_url() public async Task Reads_a_zaak_reference_by_posting_the_zaak_url_and_returns_it()
{ {
var capture = new RequestCapture(); var capture = new RequestCapture();
var client = Client(capture.Responds( var client = Client(capture.Responds(HttpStatusCode.OK, """{"reference":"REG-42"}"""));
HttpStatusCode.OK, """{"id":"zaak-1","status":"INGESCHREVEN","reference":"REG-42"}"""));
var record = await client.GetRegisterRecordAsync(new Uri(ObjectUrl)); var reference = await client.GetZaakReferenceAsync(new Uri("http://openzaak/zaken/api/v1/zaken/abc"));
Assert.Equal("zaak-1", record!.Id); Assert.Equal("REG-42", reference);
Assert.Equal("INGESCHREVEN", record.Status);
Assert.Equal("REG-42", record.Reference);
Assert.Equal(HttpMethod.Post, capture.Seen!.Method); Assert.Equal(HttpMethod.Post, capture.Seen!.Method);
Assert.Equal("http://acl/register-records/read", capture.Seen.RequestUri!.ToString()); Assert.Equal("http://acl/zaken/reference", capture.Seen.RequestUri!.ToString());
Assert.Contains($"\"objectUrl\":\"{ObjectUrl}\"", capture.Body); Assert.Contains("\"zaakUrl\":\"http://openzaak/zaken/api/v1/zaken/abc\"", capture.Body);
}
[Fact]
public async Task Reads_a_missing_record_as_nothing_to_project()
{
var capture = new RequestCapture();
var client = Client(capture.Responds(HttpStatusCode.NotFound));
// The object may be gone by the time a redelivered notification is handled (§8.6).
Assert.Null(await client.GetRegisterRecordAsync(new Uri(ObjectUrl)));
} }
[Fact] [Fact]
@@ -50,7 +35,7 @@ public class AclHttpClientTests
var client = Client(capture.Responds(HttpStatusCode.BadGateway)); var client = Client(capture.Responds(HttpStatusCode.BadGateway));
await Assert.ThrowsAsync<HttpRequestException>( await Assert.ThrowsAsync<HttpRequestException>(
() => client.GetRegisterRecordAsync(new Uri(ObjectUrl))); () => client.GetZaakReferenceAsync(new Uri("http://openzaak/zaken/api/v1/zaken/abc")));
} }
[Fact] [Fact]
@@ -60,17 +45,17 @@ public class AclHttpClientTests
var client = Client(capture.Responds(HttpStatusCode.OK, "null")); var client = Client(capture.Responds(HttpStatusCode.OK, "null"));
var ex = await Assert.ThrowsAsync<InvalidOperationException>( var ex = await Assert.ThrowsAsync<InvalidOperationException>(
() => client.GetRegisterRecordAsync(new Uri(ObjectUrl))); () => client.GetZaakReferenceAsync(new Uri("http://openzaak/zaken/api/v1/zaken/abc")));
Assert.Contains("empty", ex.Message, StringComparison.OrdinalIgnoreCase); Assert.Contains("empty", ex.Message, StringComparison.OrdinalIgnoreCase);
} }
[Fact] [Fact]
public async Task Rejects_a_null_object_url_without_sending_a_request() public async Task Rejects_a_null_zaak_url_without_sending_a_request()
{ {
var capture = new RequestCapture(); var capture = new RequestCapture();
var client = Client(capture.Responds(HttpStatusCode.OK, "{}")); var client = Client(capture.Responds(HttpStatusCode.OK, """{"reference":"REG-1"}"""));
await Assert.ThrowsAsync<ArgumentNullException>(() => client.GetRegisterRecordAsync(null!)); await Assert.ThrowsAsync<ArgumentNullException>(() => client.GetZaakReferenceAsync(null!));
Assert.Null(capture.Seen); Assert.Null(capture.Seen);
} }
} }
@@ -5,18 +5,16 @@ namespace EventSubscriber.Tests;
/// <summary>In-memory stand-ins for the projection store and notification log, so the /// <summary>In-memory stand-ins for the projection store and notification log, so the
/// projector's behaviour is exercised without Postgres (hand-written stubs, the repo's /// projector's behaviour is exercised without Postgres (hand-written stubs, the repo's
/// convention — no mocking library).</summary> /// convention — no mocking library).</summary>
/// <summary>A fake ACL client standing in for the register records Objecten holds: a test seeds a /// <summary>A fake ACL client that returns a fixed reference derived from the zaak, and records
/// record per object URL, and the call count proves a rebuild does not re-read through the ACL.</summary> /// how many times it was called (to prove a rebuild does not re-read via the ACL).</summary>
internal sealed class FakeAclClient : IAclClient internal sealed class FakeAclClient : IAclClient
{ {
public Dictionary<string, RegisterRecord> Records { get; } = [];
public int CallCount { get; private set; } public int CallCount { get; private set; }
public Task<RegisterRecord?> GetRegisterRecordAsync(Uri objectUrl, CancellationToken ct = default) public Task<string> GetZaakReferenceAsync(Uri zaakUrl, CancellationToken ct = default)
{ {
CallCount++; CallCount++;
return Task.FromResult(Records.TryGetValue(objectUrl.ToString(), out var record) ? record : null); return Task.FromResult("REG-" + zaakUrl.Segments[^1].Trim('/'));
} }
} }
@@ -2,14 +2,13 @@ using EventSubscriber.Application;
namespace EventSubscriber.Tests; namespace EventSubscriber.Tests;
/// <summary>Behaviour of the projector that turns NRC notifications into projection rows. Since /// <summary>Behaviour of the projector that turns NRC notifications into projection rows.
/// S-19b-2 the source is the register in Objecten (ADR-0030), not ZGW zaak events: a notification /// The walking skeleton reacts only to a zaak being created (status INGEDIEND) and must
/// carries only the object URL, so the record is read back through the ACL. Duplicate and /// tolerate duplicate and out-of-order deliveries (CLAUDE.md §8.6).</summary>
/// out-of-order deliveries must be tolerated (CLAUDE.md §8.6).</summary>
public sealed class NotificationProjectorTests public sealed class NotificationProjectorTests
{ {
private const string ObjectUrl = "http://objecten.local:8000/api/v2/objects/11111111-1111-1111-1111-111111111111"; private const string ZaakUrl = "http://openzaak:8000/zaken/api/v1/zaken/11111111-1111-1111-1111-111111111111";
private const string ZaakId = "99999999-9999-9999-9999-999999999999"; private const string StatusUrl = "http://openzaak:8000/zaken/api/v1/statussen/22222222-2222-2222-2222-222222222222";
private readonly InMemoryNotificationLog _log = new(); private readonly InMemoryNotificationLog _log = new();
private readonly InMemoryProjectionStore _store = new(); private readonly InMemoryProjectionStore _store = new();
@@ -17,60 +16,46 @@ public sealed class NotificationProjectorTests
private NotificationProjector Projector() => new(_log, _store, _acl); private NotificationProjector Projector() => new(_log, _store, _acl);
/// <summary>A register write as Objecten publishes it: the object is both hoofdObject and private static Notification ZaakCreated(string url = ZaakUrl)
/// resourceUrl, and the record itself is only reachable by reading that object.</summary> => new("zaken", "zaak", "create", new Uri(url));
private Notification RecordWritten(string actie = "create", string url = ObjectUrl, string status = RegistrationStatus.Ingediend, string zaakId = ZaakId)
// A status-set notification: resourceUrl is the status resource, hoofdObject is the zaak it belongs to.
private static Notification StatusSet(string zaakUrl = ZaakUrl, string statusUrl = StatusUrl)
=> new("zaken", "status", "create", new Uri(statusUrl), new Uri(zaakUrl));
[Fact]
public async Task creating_a_zaak_writes_one_row_with_status_ingediend()
{ {
_acl.Records[url] = new RegisterRecord(zaakId, status, "REG-2026-0001"); await Projector().HandleAsync(ZaakCreated());
return new Notification("objecten", "object", actie, new Uri(url));
var entry = Assert.Single(await _store.AllAsync());
Assert.Equal("11111111-1111-1111-1111-111111111111", entry.Id);
Assert.Equal(RegistrationStatus.Ingediend, entry.Status);
// Enriched with the zaak's reference (identificatie), fetched via the ACL (#78).
Assert.Equal("REG-11111111-1111-1111-1111-111111111111", entry.Reference);
} }
[Fact] [Fact]
public async Task a_register_record_write_is_projected_as_a_row_keyed_on_the_registration() public async Task rebuild_reproduces_the_reference_without_re_reading_via_the_acl()
{
await Projector().HandleAsync(RecordWritten());
var entry = Assert.Single(await _store.AllAsync());
// Keyed on the record's own id (the zaak id), not on the Objecten object's uuid — the
// projection row and the register record are the same registration.
Assert.Equal(ZaakId, entry.Id);
Assert.Equal(RegistrationStatus.Ingediend, entry.Status);
Assert.Equal("REG-2026-0001", entry.Reference);
}
// The ACL PATCHes the same object on approval. DRF routes a PATCH through `update()` but reports
// the action as `partial_update`, which is what Objecten puts in the notification — so accepting
// only `create`/`update` silently drops every approval.
[Theory]
[InlineData("partial_update")]
[InlineData("update")]
public async Task approval_updates_the_same_row_from_ingediend_to_ingeschreven(string actie)
{ {
var projector = Projector(); var projector = Projector();
await projector.HandleAsync(RecordWritten()); await projector.HandleAsync(ZaakCreated());
await projector.HandleAsync(RecordWritten(actie, status: RegistrationStatus.Ingeschreven)); var callsAfterProjection = _acl.CallCount;
await projector.RebuildAsync();
var entry = Assert.Single(await _store.AllAsync()); var entry = Assert.Single(await _store.AllAsync());
Assert.Equal(ZaakId, entry.Id); Assert.Equal("REG-11111111-1111-1111-1111-111111111111", entry.Reference);
Assert.Equal(RegistrationStatus.Ingeschreven, entry.Status); // Rebuild replays the log (which stored the reference) — no extra ACL calls (#78, ADR-0008).
} Assert.Equal(callsAfterProjection, _acl.CallCount);
[Fact]
public async Task an_object_whose_record_is_gone_is_not_projected()
{
// Nothing seeded in the fake ACL: the object was deleted before this (redelivered)
// notification was handled. Not an error — there is simply nothing to project (§8.6).
await Projector().HandleAsync(new Notification("objecten", "object", "create", new Uri(ObjectUrl)));
Assert.Empty(await _store.AllAsync());
} }
[Fact] [Fact]
public async Task replaying_the_same_notification_keeps_a_single_row() public async Task replaying_the_same_notification_keeps_a_single_row()
{ {
var projector = Projector(); var projector = Projector();
await projector.HandleAsync(RecordWritten()); await projector.HandleAsync(ZaakCreated());
await projector.HandleAsync(RecordWritten()); await projector.HandleAsync(ZaakCreated());
Assert.Single(await _store.AllAsync()); Assert.Single(await _store.AllAsync());
} }
@@ -79,8 +64,8 @@ public sealed class NotificationProjectorTests
public async Task a_replayed_notification_never_reaches_the_projection_store() public async Task a_replayed_notification_never_reaches_the_projection_store()
{ {
var projector = Projector(); var projector = Projector();
await projector.HandleAsync(RecordWritten()); await projector.HandleAsync(ZaakCreated());
await projector.HandleAsync(RecordWritten()); await projector.HandleAsync(ZaakCreated());
// The duplicate is dropped at the log, before the (idempotent) upsert — so the store // The duplicate is dropped at the log, before the (idempotent) upsert — so the store
// is written exactly once. Row count alone can't see this; the upsert count can. // is written exactly once. Row count alone can't see this; the upsert count can.
@@ -88,59 +73,77 @@ public sealed class NotificationProjectorTests
} }
[Fact] [Fact]
public async Task two_different_registrations_each_get_their_own_row() public async Task two_different_zaken_each_get_their_own_row()
{ {
var projector = Projector(); var projector = Projector();
await projector.HandleAsync(RecordWritten()); await projector.HandleAsync(ZaakCreated());
await projector.HandleAsync(RecordWritten(url: ObjectUrl[..^1] + "2", zaakId: "other-zaak")); await projector.HandleAsync(ZaakCreated(ZaakUrl[..^1] + "2")); // a distinct zaak url
Assert.Equal(2, (await _store.AllAsync()).Count); Assert.Equal(2, (await _store.AllAsync()).Count);
} }
[Theory] [Theory]
[InlineData("zaken", "zaak", "create")] // the ZGW source S-19b-2 replaced [InlineData("documenten", "enkelvoudiginformatieobject", "create")] // wrong kanaal + resource
[InlineData("zaken", "status", "create")] // ditto [InlineData("documenten", "zaak", "create")] // wrong kanaal only
[InlineData("objecten", "object", "destroy")] // a delete we do not project [InlineData("zaken", "zaak", "update")] // wrong actie
[InlineData("documenten", "object", "create")] // wrong kanaal [InlineData("zaken", "zaak", "destroy")] // wrong actie
[InlineData("zaken", "status", "update")] // a status change we ignore
[InlineData("zaken", "resultaat", "create")] // not a status we project
public async Task an_unrelated_notification_is_not_projected(string kanaal, string resource, string actie) public async Task an_unrelated_notification_is_not_projected(string kanaal, string resource, string actie)
{ {
_acl.Records[ObjectUrl] = new RegisterRecord(ZaakId, RegistrationStatus.Ingediend, "REG-2026-0001"); await Projector().HandleAsync(new Notification(kanaal, resource, actie, new Uri(ZaakUrl)));
await Projector().HandleAsync(new Notification(kanaal, resource, actie, new Uri(ObjectUrl)));
Assert.Empty(await _store.AllAsync()); Assert.Empty(await _store.AllAsync());
} }
[Fact] [Fact]
public async Task rebuild_reproduces_the_row_without_re_reading_through_the_acl() public async Task setting_a_status_projects_ingeschreven_keyed_on_the_zaak_not_the_status()
{
await Projector().HandleAsync(StatusSet());
var entry = Assert.Single(await _store.AllAsync());
// Keyed on the zaak (hoofdObject), not the status resource URL.
Assert.Equal("11111111-1111-1111-1111-111111111111", entry.Id);
Assert.Equal(RegistrationStatus.Ingeschreven, entry.Status);
}
[Fact]
public async Task approving_updates_the_existing_zaak_row_from_ingediend_to_ingeschreven()
{ {
var projector = Projector(); var projector = Projector();
await projector.HandleAsync(RecordWritten()); await projector.HandleAsync(ZaakCreated());
await projector.HandleAsync(RecordWritten("partial_update", status: RegistrationStatus.Ingeschreven)); await projector.HandleAsync(StatusSet());
var callsAfterProjection = _acl.CallCount;
var entry = Assert.Single(await _store.AllAsync());
Assert.Equal("11111111-1111-1111-1111-111111111111", entry.Id);
Assert.Equal(RegistrationStatus.Ingeschreven, entry.Status);
}
[Fact]
public async Task rebuild_reproduces_the_approved_status()
{
var projector = Projector();
await projector.HandleAsync(ZaakCreated());
await projector.HandleAsync(StatusSet());
await projector.RebuildAsync(); await projector.RebuildAsync();
var entry = Assert.Single(await _store.AllAsync()); var entry = Assert.Single(await _store.AllAsync());
Assert.Equal(RegistrationStatus.Ingeschreven, entry.Status); Assert.Equal(RegistrationStatus.Ingeschreven, entry.Status);
Assert.Equal("REG-2026-0001", entry.Reference);
// The log holds the projected row itself, so a rebuild needs neither the ACL nor
// Objecten (§8.4, ADR-0030).
Assert.Equal(callsAfterProjection, _acl.CallCount);
} }
[Fact] [Fact]
public async Task rebuild_clears_stale_rows_and_repopulates_from_the_notification_log() public async Task rebuild_clears_stale_rows_and_repopulates_from_the_notification_log()
{ {
var projector = Projector(); var projector = Projector();
await projector.HandleAsync(RecordWritten()); await projector.HandleAsync(ZaakCreated());
// A stale row that is not backed by any logged notification must not survive a rebuild. // A stale row that is not backed by any logged notification must not survive a rebuild.
await _store.UpsertAsync(new RegisterEntry("stale-9999", RegistrationStatus.Ingediend)); await _store.UpsertAsync(new RegisterEntry("stale-9999", RegistrationStatus.Ingediend));
await projector.RebuildAsync(); await projector.RebuildAsync();
var entry = Assert.Single(await _store.AllAsync()); var entry = Assert.Single(await _store.AllAsync());
Assert.Equal(ZaakId, entry.Id); Assert.Equal("11111111-1111-1111-1111-111111111111", entry.Id);
Assert.Equal(RegistrationStatus.Ingediend, entry.Status); Assert.Equal(RegistrationStatus.Ingediend, entry.Status);
} }
} }
@@ -13,8 +13,9 @@ public sealed class EfNotificationLog(ProjectionDbContext db) : INotificationLog
db.ProcessedNotifications.Add(new ProcessedNotificationRow db.ProcessedNotifications.Add(new ProcessedNotificationRow
{ {
Key = notification.Key, Key = notification.Key,
RegisterId = notification.RegisterId, Actie = notification.Actie,
Status = notification.Status, ZaakId = notification.ZaakId,
Resource = notification.Resource,
Reference = notification.Reference, Reference = notification.Reference,
ReceivedAt = DateTimeOffset.UtcNow, ReceivedAt = DateTimeOffset.UtcNow,
}); });
@@ -35,6 +36,6 @@ public sealed class EfNotificationLog(ProjectionDbContext db) : INotificationLog
public async Task<IReadOnlyList<RecordedNotification>> AllAsync(CancellationToken ct = default) public async Task<IReadOnlyList<RecordedNotification>> AllAsync(CancellationToken ct = default)
=> await db.ProcessedNotifications => await db.ProcessedNotifications
.OrderBy(r => r.ReceivedAt) .OrderBy(r => r.ReceivedAt)
.Select(r => new RecordedNotification(r.Key, r.RegisterId, r.Status, r.Reference)) .Select(r => new RecordedNotification(r.Key, r.Actie, r.ZaakId, r.Resource, r.Reference))
.ToListAsync(ct); .ToListAsync(ct);
} }
@@ -1,87 +0,0 @@
// <auto-generated />
using System;
using Microsoft.EntityFrameworkCore;
using Microsoft.EntityFrameworkCore.Infrastructure;
using Microsoft.EntityFrameworkCore.Migrations;
using Microsoft.EntityFrameworkCore.Storage.ValueConversion;
using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata;
using Projection.ReadModel;
#nullable disable
namespace Projection.ReadModel.Migrations
{
[DbContext(typeof(ProjectionDbContext))]
[Migration("20260828103132_ProjectionSourcedFromObjecten")]
partial class ProjectionSourcedFromObjecten
{
/// <inheritdoc />
protected override void BuildTargetModel(ModelBuilder modelBuilder)
{
#pragma warning disable 612, 618
modelBuilder
.HasAnnotation("ProductVersion", "10.0.0")
.HasAnnotation("Relational:MaxIdentifierLength", 63);
NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder);
modelBuilder.Entity("Projection.ReadModel.ProcessedNotificationRow", b =>
{
b.Property<string>("Key")
.HasColumnType("text")
.HasColumnName("key");
b.Property<DateTimeOffset>("ReceivedAt")
.HasColumnType("timestamp with time zone")
.HasColumnName("received_at");
b.Property<string>("Reference")
.HasColumnType("text")
.HasColumnName("reference");
b.Property<string>("RegisterId")
.IsRequired()
.HasColumnType("text")
.HasColumnName("register_id");
b.Property<string>("Status")
.IsRequired()
.HasColumnType("text")
.HasColumnName("status");
b.HasKey("Key");
b.ToTable("processed_notifications", (string)null);
});
modelBuilder.Entity("Projection.ReadModel.RegisterEntryRow", b =>
{
b.Property<string>("Id")
.HasColumnType("text")
.HasColumnName("id");
b.Property<string>("Bsn")
.HasColumnType("text")
.HasColumnName("bsn");
b.Property<string>("NaamPlaceholder")
.HasColumnType("text")
.HasColumnName("naam_placeholder");
b.Property<string>("Reference")
.HasColumnType("text")
.HasColumnName("reference");
b.Property<string>("Status")
.IsRequired()
.HasColumnType("text")
.HasColumnName("status");
b.HasKey("Id");
b.ToTable("register_projection", (string)null);
});
#pragma warning restore 612, 618
}
}
}
@@ -1,69 +0,0 @@
using Microsoft.EntityFrameworkCore.Migrations;
#nullable disable
namespace Projection.ReadModel.Migrations
{
/// <summary>
/// S-19b-2 (ADR-0030): the notification log stops describing ZGW zaak events and starts holding
/// the projected register row itself (register id, status, reference).
/// </summary>
/// <remarks>
/// The old columns are dropped and the new ones added rather than renamed. EF scaffolded renames
/// (<c>resource</c> → <c>register_id</c>, <c>zaak_id</c> → <c>status</c>), which would carry ZGW
/// values into columns that mean something else entirely — "zaak"/"status" as a register id, a
/// zaak uuid as a register status — and a rebuild would then project that garbage.
///
/// Both tables are emptied instead. A pre-existing row describes a zaak event the new projector
/// cannot reproject, and the registrations behind those rows have no RegisterRecord in Objecten
/// (only approvals wrote one before this slice), so they are not re-derivable from the new source
/// either. The projection is a derived artefact (§8.4) and repopulates as register writes arrive.
/// </remarks>
public partial class ProjectionSourcedFromObjecten : Migration
{
/// <inheritdoc />
protected override void Up(MigrationBuilder migrationBuilder)
{
// ponytail: drops the pre-slice register rather than backfilling it. Fine while stacks are
// ephemeral (a fresh `docker compose up` is the norm). If a long-lived environment ever
// needs to keep them, backfill by walking Objecten's objects instead of replaying the log.
migrationBuilder.Sql("DELETE FROM processed_notifications;");
migrationBuilder.Sql("DELETE FROM register_projection;");
migrationBuilder.DropColumn(name: "actie", table: "processed_notifications");
migrationBuilder.DropColumn(name: "zaak_id", table: "processed_notifications");
migrationBuilder.DropColumn(name: "resource", table: "processed_notifications");
migrationBuilder.AddColumn<string>(
name: "register_id",
table: "processed_notifications",
type: "text",
nullable: false,
defaultValue: "");
migrationBuilder.AddColumn<string>(
name: "status",
table: "processed_notifications",
type: "text",
nullable: false,
defaultValue: "");
}
/// <inheritdoc />
protected override void Down(MigrationBuilder migrationBuilder)
{
migrationBuilder.Sql("DELETE FROM processed_notifications;");
migrationBuilder.Sql("DELETE FROM register_projection;");
migrationBuilder.DropColumn(name: "register_id", table: "processed_notifications");
migrationBuilder.DropColumn(name: "status", table: "processed_notifications");
migrationBuilder.AddColumn<string>(
name: "actie", table: "processed_notifications", type: "text", nullable: false, defaultValue: "");
migrationBuilder.AddColumn<string>(
name: "zaak_id", table: "processed_notifications", type: "text", nullable: false, defaultValue: "");
migrationBuilder.AddColumn<string>(
name: "resource", table: "processed_notifications", type: "text", nullable: false, defaultValue: "");
}
}
}
@@ -28,6 +28,11 @@ namespace Projection.ReadModel.Migrations
.HasColumnType("text") .HasColumnType("text")
.HasColumnName("key"); .HasColumnName("key");
b.Property<string>("Actie")
.IsRequired()
.HasColumnType("text")
.HasColumnName("actie");
b.Property<DateTimeOffset>("ReceivedAt") b.Property<DateTimeOffset>("ReceivedAt")
.HasColumnType("timestamp with time zone") .HasColumnType("timestamp with time zone")
.HasColumnName("received_at"); .HasColumnName("received_at");
@@ -36,15 +41,15 @@ namespace Projection.ReadModel.Migrations
.HasColumnType("text") .HasColumnType("text")
.HasColumnName("reference"); .HasColumnName("reference");
b.Property<string>("RegisterId") b.Property<string>("Resource")
.IsRequired() .IsRequired()
.HasColumnType("text") .HasColumnType("text")
.HasColumnName("register_id"); .HasColumnName("resource");
b.Property<string>("Status") b.Property<string>("ZaakId")
.IsRequired() .IsRequired()
.HasColumnType("text") .HasColumnType("text")
.HasColumnName("status"); .HasColumnName("zaak_id");
b.HasKey("Key"); b.HasKey("Key");
@@ -34,8 +34,9 @@ public sealed class ProjectionDbContext(DbContextOptions<ProjectionDbContext> op
e.ToTable("processed_notifications"); e.ToTable("processed_notifications");
e.HasKey(r => r.Key); e.HasKey(r => r.Key);
e.Property(r => r.Key).HasColumnName("key"); e.Property(r => r.Key).HasColumnName("key");
e.Property(r => r.RegisterId).HasColumnName("register_id").IsRequired(); e.Property(r => r.Actie).HasColumnName("actie").IsRequired();
e.Property(r => r.Status).HasColumnName("status").IsRequired(); e.Property(r => r.ZaakId).HasColumnName("zaak_id").IsRequired();
e.Property(r => r.Resource).HasColumnName("resource").IsRequired();
e.Property(r => r.Reference).HasColumnName("reference"); e.Property(r => r.Reference).HasColumnName("reference");
e.Property(r => r.ReceivedAt).HasColumnName("received_at"); e.Property(r => r.ReceivedAt).HasColumnName("received_at");
}); });
@@ -55,20 +56,18 @@ public sealed class RegisterEntryRow
public string? NaamPlaceholder { get; set; } public string? NaamPlaceholder { get; set; }
} }
/// <summary>An accepted notification, retained so the projection can be rebuilt without reading /// <summary>An accepted notification, retained so the projection can be rebuilt without OpenZaak (§8.1).</summary>
/// Objecten or ZGW (§8.1, §8.4). Since S-19b-2 it holds the projected row itself — the register
/// record's id, status and reference — so a rebuild is a replay with no mapping rules (ADR-0030).</summary>
public sealed class ProcessedNotificationRow public sealed class ProcessedNotificationRow
{ {
public required string Key { get; set; } public required string Key { get; set; }
public required string Actie { get; set; }
public required string ZaakId { get; set; }
/// <summary>The registration this record is for (the zaak id) — the projection row's key.</summary> /// <summary>The ZGW resource (e.g. <c>zaak</c> or <c>status</c>) — retained so a rebuild reprojects
public required string RegisterId { get; set; } /// the right status without reading OpenZaak (S-09b).</summary>
public required string Resource { get; set; }
/// <summary>The register status the record carried (INGEDIEND / INGESCHREVEN).</summary> /// <summary>The zaak reference (identificatie), retained so a rebuild reprojects it without the ACL (#78).</summary>
public required string Status { get; set; }
/// <summary>The citizen-facing reference the record carried — matches the submit confirmation (#78).</summary>
public string? Reference { get; set; } public string? Reference { get; set; }
public DateTimeOffset ReceivedAt { get; set; } public DateTimeOffset ReceivedAt { get; set; }
@@ -1,28 +1,19 @@
# language: en # language: en
# Drives S-19b-2 (#153), re-sourcing S-06 (#7). The read projection is derived from the # Drives S-06 (#7). On a zaak-created notification from NRC the Event Subscriber writes a
# RegisterRecord in Objecten (ADR-0030), not from ZGW zaak events: the ACL records a registration # rebuildable read-projection row (PRD §8.4). This scenario exercises the use case against an
# in the register, Objecten notifies, and the Event Subscriber projects the record that # in-memory stand-in for the projection store and notification log; real OpenZaak → NRC →
# notification points at. This scenario exercises the use case against in-memory stand-ins for the # subscriber delivery is verified by the live-stack check (verify-projection, ADR-0007/#58).
# register, the projection store and the notification log; real Objecten → NRC → subscriber Feature: Register-projectie bijwerken op een zaaknotificatie
# delivery is verified by the live-stack check (verify-projection, ADR-0007/0030). Als openbaar register wil ik dat een aangemaakte zaak in de projectie verschijnt
Feature: Register-projectie bijwerken op een registerwijziging zodat het register de ingediende registratie kan tonen.
Als openbaar register wil ik dat een registratie in de projectie verschijnt zodra zij
in het register is vastgelegd, zodat het register haar actuele status kan tonen.
Scenario: Een ingediende registratie levert een rij met status INGEDIEND Scenario: Een zaaknotificatie levert een rij met status INGEDIEND
Given registration "11111111-1111-1111-1111-111111111111" is recorded in the register with status "INGEDIEND" Given a zaak is created in OpenZaak with id "11111111-1111-1111-1111-111111111111"
When the register notification is delivered to the event subscriber When the NRC notification for that zaak is delivered to the event subscriber
Then the register projection contains a row for "11111111-1111-1111-1111-111111111111" with status "INGEDIEND" Then the register projection contains a row for "11111111-1111-1111-1111-111111111111" with status "INGEDIEND"
Scenario: Een goedgekeurde registratie werkt dezelfde rij bij
Given registration "33333333-3333-3333-3333-333333333333" is recorded in the register with status "INGEDIEND"
And the register notification is delivered to the event subscriber
When registration "33333333-3333-3333-3333-333333333333" is recorded in the register with status "INGESCHREVEN"
And the register notification is delivered to the event subscriber
Then the register projection contains a row for "33333333-3333-3333-3333-333333333333" with status "INGESCHREVEN"
Scenario: Dezelfde notificatie tweemaal levert geen duplicaat Scenario: Dezelfde notificatie tweemaal levert geen duplicaat
Given registration "22222222-2222-2222-2222-222222222222" is recorded in the register with status "INGEDIEND" Given a zaak is created in OpenZaak with id "22222222-2222-2222-2222-222222222222"
When the register notification is delivered to the event subscriber When the NRC notification for that zaak is delivered to the event subscriber
And the same register notification is delivered again And the same NRC notification is delivered again
Then the register projection contains exactly one row for "22222222-2222-2222-2222-222222222222" Then the register projection contains exactly one row for "22222222-2222-2222-2222-222222222222"
@@ -5,39 +5,31 @@ using Xunit;
namespace Acceptance.Steps; namespace Acceptance.Steps;
/// <summary>Bindings for <c>RegisterProjectieBijwerken.feature</c> (S-06, re-sourced by S-19b-2). /// <summary>Bindings for <c>RegisterProjectieBijwerken.feature</c> (S-06). Reqnroll creates
/// Reqnroll creates one instance per scenario, so instance fields hold scenario-scoped state.</summary> /// one instance per scenario, so instance fields hold scenario-scoped state.</summary>
[Binding] [Binding]
public sealed class RegisterProjectieBijwerkenSteps public sealed class RegisterProjectieBijwerkenSteps
{ {
private const string ObjectBase = "http://objecten.local:8000/api/v2/objects/"; private const string ZaakBase = "http://openzaak:8000/zaken/api/v1/zaken/";
private readonly InMemoryNotificationLog _log = new(); private readonly InMemoryNotificationLog _log = new();
private readonly InMemoryProjectionStore _store = new(); private readonly InMemoryProjectionStore _store = new();
private readonly InMemoryRegisterRecordClient _register = new();
private readonly NotificationProjector _projector; private readonly NotificationProjector _projector;
private Notification? _notification; private Notification? _notification;
public RegisterProjectieBijwerkenSteps() public RegisterProjectieBijwerkenSteps()
=> _projector = new NotificationProjector(_log, _store, _register); => _projector = new NotificationProjector(_log, _store, new InMemoryAclReferenceClient());
[Given("registration \"(.*)\" is recorded in the register with status \"(.*)\"")] [Given("a zaak is created in OpenZaak with id \"(.*)\"")]
[When("registration \"(.*)\" is recorded in the register with status \"(.*)\"")] public void GivenAZaakIsCreatedInOpenZaakWithId(string id)
public void RegistrationIsRecorded(string id, string status) => _notification = new Notification("zaken", "zaak", "create", new Uri(ZaakBase + id));
{
// The ACL upserts one object per registration, so submit and approval share an object URL.
var objectUrl = ObjectBase + id;
_register.Records[objectUrl] = new RegisterRecord(id, status, "REG-" + id);
_notification = new Notification("objecten", "object", "create", new Uri(objectUrl));
}
[Given("the register notification is delivered to the event subscriber")] [When("the NRC notification for that zaak is delivered to the event subscriber")]
[When("the register notification is delivered to the event subscriber")] public Task WhenTheNotificationIsDelivered()
public Task TheNotificationIsDelivered()
=> _projector.HandleAsync(_notification!); => _projector.HandleAsync(_notification!);
[When("the same register notification is delivered again")] [When("the same NRC notification is delivered again")]
public Task TheSameNotificationIsDeliveredAgain() public Task WhenTheSameNotificationIsDeliveredAgain()
=> _projector.HandleAsync(_notification!); => _projector.HandleAsync(_notification!);
[Then("the register projection contains a row for \"(.*)\" with status \"(.*)\"")] [Then("the register projection contains a row for \"(.*)\" with status \"(.*)\"")]
@@ -39,12 +39,10 @@ public sealed class InMemoryProjectionStore : IProjectionStore
=> [.. _byId.Values.Where(e => e.Id == id)]; => [.. _byId.Values.Where(e => e.Id == id)];
} }
/// <summary>An in-memory stand-in for the register the ACL reads back for the projector, so the /// <summary>A fake ACL client for the projection acceptance scenario: returns a reference derived
/// scenario runs without a running ACL or Objecten (S-19b-2, ADR-0030).</summary> /// from the zaak, so the projector can enrich rows without a running ACL (#78).</summary>
public sealed class InMemoryRegisterRecordClient : IAclClient public sealed class InMemoryAclReferenceClient : IAclClient
{ {
public Dictionary<string, RegisterRecord> Records { get; } = []; public Task<string> GetZaakReferenceAsync(Uri zaakUrl, CancellationToken ct = default)
=> Task.FromResult("REG-" + zaakUrl.Segments[^1].Trim('/'));
public Task<RegisterRecord?> GetRegisterRecordAsync(Uri objectUrl, CancellationToken ct = default)
=> Task.FromResult(Records.TryGetValue(objectUrl.ToString(), out var record) ? record : null);
} }
@@ -65,8 +65,4 @@ public sealed class InMemoryRegisterRecordGateway : IRegisterRecordGateway
Upserted.Add(record); Upserted.Add(record);
return Task.CompletedTask; return Task.CompletedTask;
} }
/// <summary>The most recently written record — scenarios never read one back by object URL.</summary>
public Task<RegisterRecord?> GetAsync(Uri objectUrl, CancellationToken ct = default)
=> Task.FromResult(Upserted.Count == 0 ? null : Upserted[^1]);
} }
+6 -11
View File
@@ -1,16 +1,11 @@
import { expect, request, test } from '@playwright/test'; import { expect, request, test } from '@playwright/test';
// Walking-skeleton happy path (S-08d + S-09 + S-09b + S-12 + S-10a + S-19b-2): a zorgprofessional // Walking-skeleton happy path (S-08d + S-09 + S-09b + S-12 + S-10a): a zorgprofessional logs in via
// logs in via mock DigiD and submits through the self-service portal → BFF → domain; the entry // mock DigiD and submits through the self-service portal → BFF → domain; the entry appears in the
// appears in the openbaar register as INGEDIEND; the citizen supplies the documents the process is // openbaar register as INGEDIEND; the citizen supplies the documents the process is waiting for
// waiting for (S-10a); a behandelaar then logs in to the behandel portal, finds the registration in // (S-10a); a behandelaar then logs in to the behandel portal, finds the registration in the werkbak,
// the werkbak, and approves it (goedkeuren); the decision completes the Flowable Beoordelen task and // and approves it (goedkeuren); the decision completes the Flowable Beoordelen task and flows via the
// flows via the ACL → Objecten → NRC → event-subscriber → projection, and the openbaar register // ACL → NRC → event-subscriber → projection, and the openbaar register shows INGESCHREVEN.
// shows INGESCHREVEN.
//
// Since ADR-0030 both public statuses come from the register in Objecten, not from ZGW zaak events:
// the ACL writes the record on submit (INGEDIEND) and upserts it on approval (INGESCHREVEN), so the
// INGEDIEND assertion below is itself proof of the re-sourced path.
test('DigiD submit → public INGEDIEND → documenten → behandelaar goedkeurt → public INGESCHREVEN', async ({ test('DigiD submit → public INGEDIEND → documenten → behandelaar goedkeurt → public INGESCHREVEN', async ({
page, page,
context, context,