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
id: projection
run: make verify-projection
- name: Objecten → NRC notification delivery
id: objecten_nrc
run: make verify-objecten-notifications
- name: Domain → Flowable → ACL → OpenZaak
id: domain
run: make verify-domain
@@ -248,7 +245,6 @@ jobs:
OBJECTTYPEN: ${{ steps.objecttypen.outcome }}
OBJECTEN: ${{ steps.objecten.outcome }}
REGISTERRECORD: ${{ steps.registerrecord.outcome }}
OBJECTEN_NOTIFICATIONS: ${{ steps.objecten_nrc.outcome }}
ACL: ${{ steps.acl.outcome }}
NRC: ${{ steps.nrc.outcome }}
PROJECTION: ${{ steps.projection.outcome }}
@@ -270,7 +266,6 @@ jobs:
echo "| Objecttypen API + token | $(icon "$OBJECTTYPEN") |"
echo "| Objecten API + token | $(icon "$OBJECTEN") |"
echo "| RegisterRecord objecttype | $(icon "$REGISTERRECORD") |"
echo "| Objecten → NRC | $(icon "$OBJECTEN_NOTIFICATIONS") |"
echo "| ACL ↔ OpenZaak | $(icon "$ACL") |"
echo "| OpenZaak → NRC | $(icon "$NRC") |"
echo "| NRC → Event Subscriber → projection | $(icon "$PROJECTION") |"
@@ -290,7 +285,7 @@ jobs:
# Log dump must precede teardown (which removes the containers).
- name: Dump container logs on 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
if: always()
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):
- **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-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.
- **S-19b** (#150) · Read projection sourced from Objecten instead of NRC zaak events. Depends on S-19a.
---
+1 -7
View File
@@ -43,7 +43,7 @@ export DOCKER_HOST := unix://$(PODMAN_SOCK)
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)
## `verify` is the live-stack stage (full stack up once → ACL + notification checks).
@@ -201,11 +201,6 @@ verify-objecten:
verify-registerrecord:
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,
## tear down (always). For fast single-concern local iteration use `integration`
## (oz-only) or `verify-notifications` (oz+nrc) instead.
@@ -217,7 +212,6 @@ verify:
&& bash infra/run-acl-integration.sh \
&& bash infra/run-notification-check.sh \
&& bash infra/run-projection-check.sh \
&& bash infra/run-objecten-notifications-check.sh \
&& bash infra/run-domain-check.sh \
&& bash infra/run-bff-check.sh \
&& 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
register write yet.
- **Lifted by ADR-0029** (S-19b-1, #152): broker, worker, `objecten` kanaal and
notifications config now exist, and `NOTIFICATIONS_DISABLED` is `false`.
- Upgrade path: S-19b (#150) needs those notifications to source the projection from
Objecten, and turns them on together with the broker, worker, kanaal and abonnement.
## 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.
- The disclosure boundary is enforced by Objecten's schema validation (ADR-0027), not by
discipline in projection code.
- The read projection can become a cache of Objecten rather than a re-derivation of ZGW
done in S-19b-2 (#153), ADR-0030.
- The read projection can become a cache of Objecten rather than a re-derivation of ZGW
(S-19b, #150).
**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.
- Two new hand-kept constants: the pinned objecttype UUID (two files) and the objecttype
name (compose + `register.py`).
- ~~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 projection is now derived from the register, so there is only one source to agree with.
- 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.
## 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
# 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.
# Dotted host on purpose — see the `objecten.local` alias below (ADR-0029).
Acl__Objecten__BaseUrl: http://objecten.local:8000/
Acl__Objecten__BaseUrl: http://objecten:8000/
Acl__Objecten__Token: ${OBJECTEN_TOKEN:-1234567890abcdef1234567890abcdef12345678}
Acl__Objecten__ObjecttypenBaseUrl: http://objecttypen:8000/
Acl__Objecten__ObjecttypenToken: ${OBJECTTYPEN_TOKEN:-0123456789abcdef0123456789abcdef01234567}
@@ -692,14 +691,12 @@ services:
CACHE_AXES: objecten-redis:6379/0
DISABLE_2FA: "true"
OTEL_SDK_DISABLED: "true"
CELERY_BROKER_URL: redis://objecten-redis:6379/1
CELERY_RESULT_BACKEND: redis://objecten-redis:6379/1
# Publish register-record events to NRC on the `objecten` kanaal (S-19b-1, ADR-0029). The NRC
# service + notifications_config are provisioned by setup_configuration
# (infra/objecten/setup_configuration/data.yaml), and objecten-celery below actually sends
# them — notifications_api_common only queues the task. See ADR-0028 for why S-19a left this
# off until all four pieces existed.
NOTIFICATIONS_DISABLED: "false"
# S-19a: Objecten refuses every write while its Notificaties config is absent
# (notifications_api_common raises rather than skipping, so POST /objects 500s). Objecten →
# NRC is not wired yet — there is no broker, worker, kanaal or abonnement for it — so turn
# notifications off rather than fake a delivery path that silently drops every message.
# S-19b (#150) sources the projection from Objecten and turns this back on for real.
NOTIFICATIONS_DISABLED: "true"
RUN_SETUP_CONFIG: "true"
command: /setup_configuration.sh
volumes:
@@ -724,28 +721,6 @@ services:
start_period: 30s
ports:
- "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:
objecten-init:
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
# 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.
# Dotted host on purpose — see the `objecten.local` alias below (ADR-0029).
Acl__Objecten__BaseUrl: http://objecten.local:8000/
Acl__Objecten__BaseUrl: http://objecten:8000/
Acl__Objecten__Token: ${OBJECTEN_TOKEN:-1234567890abcdef1234567890abcdef12345678}
Acl__Objecten__ObjecttypenBaseUrl: http://objecttypen:8000/
Acl__Objecten__ObjecttypenToken: ${OBJECTTYPEN_TOKEN:-0123456789abcdef0123456789abcdef01234567}
@@ -718,14 +717,12 @@ services:
CACHE_AXES: objecten-redis:6379/0
DISABLE_2FA: "true"
OTEL_SDK_DISABLED: "true"
CELERY_BROKER_URL: redis://objecten-redis:6379/1
CELERY_RESULT_BACKEND: redis://objecten-redis:6379/1
# Publish register-record events to NRC on the `objecten` kanaal (S-19b-1, ADR-0029). The NRC
# service + notifications_config are provisioned by setup_configuration
# (infra/objecten/setup_configuration/data.yaml), and objecten-celery below actually sends
# them — notifications_api_common only queues the task. See ADR-0028 for why S-19a left this
# off until all four pieces existed.
NOTIFICATIONS_DISABLED: "false"
# S-19a: Objecten refuses every write while its Notificaties config is absent
# (notifications_api_common raises rather than skipping, so POST /objects 500s). Objecten →
# NRC is not wired yet — there is no broker, worker, kanaal or abonnement for it — so turn
# notifications off rather than fake a delivery path that silently drops every message.
# S-19b (#150) sources the projection from Objecten and turns this back on for real.
NOTIFICATIONS_DISABLED: "true"
RUN_SETUP_CONFIG: "true"
command: /setup_configuration.sh
# data.yaml is streamed into this external volume by infra/seed-config.sh before start.
@@ -753,28 +750,6 @@ services:
start_period: 30s
ports:
- "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:
objecten-init:
condition: service_completed_successfully
+6 -11
View File
@@ -2,10 +2,10 @@
"""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
abonnement on the `objecten` kanaal pointing at the event-subscriber's /notifications callback, so
the register writes the ACL makes (INGEDIEND on submit, INGESCHREVEN on approval) reach the
projection — without this the openbaar (public) register stays empty. Since S-19b-2 the projection
is sourced from the register in Objecten, not from ZGW zaak events (ADR-0030).
abonnement on the `zaken` kanaal pointing at the event-subscriber's /notifications callback, so
OpenZaak's notifications (zaak create + status set) reach the projection — without this the openbaar
(public) register stays empty. This is what infra/verify-notification-driver.py does for CI (minus
the test zaak it also creates).
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
@@ -22,8 +22,6 @@ SINK_PORT = os.environ.get("SINK_PORT", "8080")
SINK_AUTH = os.environ.get("SINK_AUTH", "Bearer big-reference-notifications")
CID = os.environ.get("OZ_CLIENT_ID", "big-reference-seed")
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():
@@ -62,10 +60,7 @@ def main():
status, body = call("GET", f"{NRC}/api/v1/abonnement")
for ab in (body or []) if status == 200 else []:
if str(ab.get("callbackUrl", "")).endswith("/notifications"):
# The kanaal is part of "current": an abonnement left over from before S-19b-2 points at
# 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]:
if ab.get("callbackUrl") == callback:
print(f"abonnement already current: {ab['url']}")
return
call("DELETE", ab["url"])
@@ -73,7 +68,7 @@ def main():
status, ab = call("POST", f"{NRC}/api/v1/abonnement", {
"callbackUrl": callback, "auth": SINK_AUTH,
"kanalen": [{"naam": KANAAL, "filters": {}}]})
"kanalen": [{"naam": "zaken", "filters": {}}]})
if status != 201:
sys.exit(f"create abonnement -> {status}: {json.dumps(ab)}")
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
header_key: Authorization
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
# objecttype it has not been configured with ("ObjectType with url=… is not configured"), and it
@@ -50,10 +40,3 @@ tokenauth:
email: admin@localhost
organization: Respellion
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:
authorizations_api_service_identifier: openzaak-ac
# 4. The kanalen publishers announce on: `zaken` (OpenZaak) and `objecten` (Objecten, S-19b-1).
# 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.
# 4. The kanaal OpenZaak publishes zaak events on.
notifications_kanalen_config_enable: true
notifications_kanalen_config:
items:
@@ -41,11 +39,3 @@ notifications_kanalen_config:
- bronorganisatie
- zaaktype
- 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
#
# Verify the end-to-end read-projection path (S-06, re-sourced by S-19b-2) against an ALREADY-RUNNING
# full stack: ACL → Objecten → NRC → Event Subscriber → projection → projection-api. Seeds a
# published BIG zaaktype (idempotent), registers an abonnement on the `objecten` kanaal pointing at
# the real Event Subscriber's /notifications callback (with the bearer it enforces), opens a zaak
# *through the ACL*, and asserts projection-api serves a row for it 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.
# Verify the end-to-end read-projection path (S-06) against an ALREADY-RUNNING full stack:
# OpenZaak → NRC → Event Subscriber → projection → projection-api. Seeds a published BIG
# zaaktype (idempotent), registers an abonnement on the `zaken` kanaal pointing at the real
# Event Subscriber's /notifications callback (with the bearer it enforces), creates a zaak,
# and asserts projection-api serves a row for that zaak with status INGEDIEND.
#
# 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
# lifecycle (the caller brings it up and tears it down), but does recreate the `acl` service to
# repoint it — see below, and run-domain-check.sh, which does the same. Plain docker primitives only.
# See ADR-0007/0008/0030.
# the runner can't reach published ports (gitea-actions-gotchas.md §5/§6). Reuses the
# notification driver to register the abonnement + create the zaak. Does NOT manage the stack
# lifecycle (the caller owns bring-up + teardown). Plain docker primitives only. See ADR-0007/0008.
set -euo pipefail
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}"
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)"
es="$(docker ps -q --filter 'name=event-subscriber' | 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 "$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)"
oz_ip="$(ip "$oz")"; nrc_ip="$(ip "$nrc")"; es_ip="$(ip "$es")"; proj_ip="$(ip "$proj")"; acl_ip="$(ip "$acl")"
echo ">> network=$net openzaak=$oz_ip nrc=$nrc_ip event-subscriber=$es_ip projection-api=$proj_ip acl=$acl_ip"
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"
echo ">> seeding a published BIG zaaktype (idempotent)"
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 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
# The same script the local stack uses (ADR-0020), so both paths register the identical abonnement.
drv="$(docker create --network "$net" --name rr-pverify \
-e "NRC_BASE=http://$nrc_ip:8000" \
-e "SINK_HOST=$es_ip" -e "SINK_PORT=8080" -e "SINK_AUTH=$WEBHOOK_AUTH" \
python:3-slim python /subscribe.py)"
docker cp "$here/local/register-abonnement.py" "$drv:/subscribe.py" >/dev/null
-e "OZ_BASE=http://$oz_ip:8000" -e "NRC_BASE=http://$nrc_ip:8000" \
-e "SINK_CALLBACK=http://$es_ip:8080/notifications" -e "SINK_AUTH=$WEBHOOK_AUTH" \
python:3-slim python /driver.py)"
docker cp "$here/verify-notification-driver.py" "$drv:/driver.py" >/dev/null
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
# 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; }
[ -n "$zaak_url" ] || { echo "ERROR: driver did not create a zaak" >&2; exit 1; }
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)"
for _ in $(seq 1 30); do
@@ -93,8 +63,6 @@ for _ in $(seq 1 30); do
sleep 2
done
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 "--- 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
+3 -7
View File
@@ -15,13 +15,9 @@ set -euo pipefail
timeout="${WAIT_TIMEOUT:-420}"
deadline=$(( $(date +%s) + timeout ))
# compose service name -> container id. `--filter name=` is a substring match, so it is anchored on
# the compose replica suffix — otherwise 'objecten' also matches objecten-db / objecten-redis /
# objecten-celery, and 'objecttypen' matches objecttypen-db. Whichever docker listed first won, so a
# 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; }
# compose service name -> container id. The name filter matches both docker
# compose ("infra-openzaak-1") and podman-compose ("infra_openzaak_1") naming.
cid_for() { docker ps -aq --filter "name=$1" | head -1; }
for svc in "$@"; do
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 });
});
// 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
// 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) =>
@@ -141,9 +131,6 @@ public sealed record CancelZaakRequest(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 partial class Program;
+1 -22
View File
@@ -24,16 +24,7 @@ public sealed class AclService(
clock.Today,
registration.Reference);
var zaakUrl = 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;
return await gateway.OpenZaakAsync(request, ct);
}
/// <summary>
@@ -61,18 +52,6 @@ public sealed class AclService(
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>
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).
/// </summary>
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>
@@ -1,4 +1,3 @@
using System.Net;
using System.Net.Http.Headers;
using System.Net.Http.Json;
using System.Text.Json.Serialization;
@@ -39,30 +38,6 @@ public sealed class ObjectenGateway(HttpClient http, ObjectenOptions options, IC
"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) =>
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(
[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(
[property: JsonPropertyName("type")] string Type,
[property: JsonPropertyName("record")] RecordDto Record);
@@ -20,7 +20,7 @@ public sealed class ObjectenGatewayIntegrationTests
new HttpClient(),
new ObjectenOptions
{
BaseUrl = new(Env("OBJECTEN_BASE", "http://objecten.local:8000")),
BaseUrl = new(Env("OBJECTEN_BASE", "http://objecten:8000")),
Token = Env("OBJECTEN_TOKEN", "1234567890abcdef1234567890abcdef12345678"),
ObjecttypenBaseUrl = new(Env("OBJECTTYPEN_BASE", "http://objecttypen:8000")),
ObjecttypenToken = Env("OBJECTTYPEN_TOKEN", "0123456789abcdef0123456789abcdef01234567"),
@@ -65,7 +65,7 @@ public sealed class ObjectenGatewayIntegrationTests
{
using var http = new HttpClient();
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) +
"&data_attrs=id__exact__" + Uri.EscapeDataString(id));
-56
View File
@@ -79,21 +79,11 @@ public class AclServiceTests
{
public readonly List<RegisterRecord> Upserted = [];
public RegisterRecord? Stored;
public Uri? ReadFrom;
public Task UpsertAsync(RegisterRecord record, CancellationToken ct = default)
{
Upserted.Add(record);
return Task.CompletedTask;
}
public Task<RegisterRecord?> GetAsync(Uri objectUrl, CancellationToken ct = default)
{
ReadFrom = objectUrl;
return Task.FromResult(Stored);
}
}
private static AclDefaults Defaults() => new()
@@ -140,52 +130,6 @@ public class AclServiceTests
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]
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");
[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]
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.Text.Json.Serialization;
using EventSubscriber.Application;
@@ -6,28 +5,26 @@ using EventSubscriber.Application;
namespace EventSubscriber.Api;
/// <summary>
/// HTTP client to the ACL service. An Objecten notification carries only the object URL, so the
/// subscriber reads the register record back through the ACL — the only code that may talk to
/// Objecten (§8.1, ADR-0028/ADR-0030) — rather than reading Objecten itself.
/// HTTP client to the ACL service. The subscriber enriches the projection with the zaak's reference
/// (identificatie) by asking the ACL — the only code that may read ZGW (§8.1) — rather than reading
/// OpenZaak itself (adr-proposal #78).
/// </summary>
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(
new Uri(http.BaseAddress!, "register-records/read"),
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;
new Uri(http.BaseAddress!, "zaken/reference"), new ReferenceRequest(zaakUrl.ToString()), ct);
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();
/// <summary>The NRC notification body, as Open Notificaties POSTs it. Only the fields the projector
/// needs are bound; <c>aanmaakdatum</c>, <c>kenmerken</c> and <c>hoofdObject</c> are ignored for a
/// register write hoofdObject is the same object as resourceUrl (ADR-0030).</summary>
public sealed record NotificationDto(string Kanaal, string Resource, string Actie, Uri ResourceUrl)
/// <summary>The NRC notification body, as Open Notificaties POSTs it. Only the fields the
/// projection needs are bound; <c>aanmaakdatum</c>/<c>kenmerken</c> are ignored for the minimal slice.</summary>
public sealed record NotificationDto(string Kanaal, string Resource, string Actie, Uri ResourceUrl, Uri? HoofdObject = null)
{
public Notification ToNotification() => new(Kanaal, Resource, Actie, ResourceUrl);
public Notification ToNotification() => new(Kanaal, Resource, Actie, ResourceUrl, HoofdObject);
}
public partial class Program
@@ -2,39 +2,40 @@ namespace EventSubscriber.Application;
/// <summary>
/// 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>
/// <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(
string Kanaal,
string Resource,
string Actie,
Uri ResourceUrl)
Uri ResourceUrl,
Uri? HoofdObject = null)
{
/// <summary>
/// A register record written to Objecten — <c>create</c> on submit and <c>partial_update</c> on
/// approval, since the ACL upserts the same object for a registration (§8.6).
/// </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>A zaak being created — projected as INGEDIEND.</summary>
public bool IsZaakCreated =>
Kanaal == "zaken" && Resource == "zaak" && Actie == "create";
/// <summary>The object holding the register record. For a <c>resource: object</c> notification
/// Objecten sends the object as both <c>hoofdObject</c> and <c>resourceUrl</c> — the object is
/// the main resource — so the notification's own <c>hoofdObject</c> is not modelled.</summary>
public Uri ObjectUrl => ResourceUrl;
/// <summary>A status being set on a zaak — the approval, projected as INGESCHREVEN (S-09b). In the
/// walking skeleton the only status ever set after creation is the approval, and the subscriber may
/// not read OpenZaak (§8.1), so any status-create is taken as the approval.</summary>
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>
/// 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
/// 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>
public sealed class NotificationProjector(INotificationLog log, IProjectionStore store, IAclClient acl)
{
/// <summary>Handle one inbound notification. Reacts to a register record being written to
/// Objecten (S-19b-2, ADR-0030) and ignores everything else. The notification carries only the
/// object URL, so the record is read back through the ACL (§8.1) and becomes the row verbatim.</summary>
/// <summary>Handle one inbound notification. Reacts to a zaak being created (INGEDIEND) and a
/// status being set (INGESCHREVEN); ignores everything else. Enriches the row with the zaak's
/// 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)
{
ArgumentNullException.ThrowIfNull(notification);
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)
if (!notification.IsZaakCreated && !notification.IsZaakStatusSet)
return;
var reference = await acl.GetZaakReferenceAsync(notification.ZaakUrl, ct);
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
// 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);
}
/// <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>
public async Task RebuildAsync(CancellationToken ct = default)
{
@@ -55,9 +35,11 @@ public sealed class NotificationProjector(INotificationLog log, IProjectionStore
await store.UpsertAsync(ToEntry(recorded), ct);
}
/// <summary>The projection row for an accepted notification. The log already holds exactly the
/// row's fields, so a rebuild needs no mapping rules and no upstream reads. bsn/naam stay
/// deferred — the register record is public-safe by construction (ADR-0027).</summary>
/// <summary>The projection row for an accepted notification: a status-set maps to INGESCHREVEN,
/// a zaak-create to INGEDIEND. bsn/naam are deferred (ADR-0008).</summary>
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
/// 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
/// 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>
public interface INotificationLog
{
@@ -19,29 +19,22 @@ public interface INotificationLog
Task<IReadOnlyList<RecordedNotification>> AllAsync(CancellationToken ct = default);
}
/// <summary>
/// An accepted notification, retaining exactly the projection row it produced — so a rebuild
/// reproduces the row by replaying the log, without re-reading Objecten (S-19b-2, ADR-0030).
/// </summary>
public sealed record RecordedNotification(string Key, string RegisterId, string Status, string? Reference);
/// <summary>A notification that has been accepted, retaining what a rebuild needs to recompute its
/// projection row — the ZGW <c>resource</c> (zaak-create → INGEDIEND vs status-set → INGESCHREVEN) and
/// the zaak <c>reference</c> (identificatie), so a rebuild reproduces the row without re-reading ZGW (#78).</summary>
public sealed record RecordedNotification(string Key, string Actie, string ZaakId, string Resource, string? Reference);
/// <summary>
/// Port to the Anti-Corruption Layer. An Objecten notification carries only the object URL, so the
/// subscriber reads the register record back through the ACL — the only code that may talk to
/// Objecten (§8.1, ADR-0028) — rather than reading Objecten itself.
/// Port to the Anti-Corruption Layer. The subscriber enriches the projection with the zaak's
/// public-safe reference (its identificatie) by asking the ACL — the only code that may read ZGW
/// (§8.1) — rather than reading OpenZaak itself (adr-proposal #78).
/// </summary>
public interface IAclClient
{
/// <summary>The register record the object at <paramref name="objectUrl"/> holds, or
/// <c>null</c> if it holds none — the object may be gone by the time a redelivered
/// notification is handled, which is not an error (§8.6).</summary>
Task<RegisterRecord?> GetRegisterRecordAsync(Uri objectUrl, CancellationToken ct = default);
/// <summary>The zaak's reference (identificatie) for the read projection.</summary>
Task<string> GetZaakReferenceAsync(Uri zaakUrl, 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
/// subscriber writes to it and the projection-api reads it.</summary>
public interface IProjectionStore
@@ -5,42 +5,27 @@ using EventSubscriber.Api;
namespace EventSubscriber.Tests;
/// <summary>
/// Unit tests for the subscriber's ACL client, which reads a register record through the ACL — the
/// only code allowed to talk to Objecten (§8.1, ADR-0028/ADR-0030). Uses a scripted message handler
/// so no real ACL is required.
/// Unit tests for the subscriber's ACL client, which reads a zaak's reference (identificatie) through
/// the ACL — the only code allowed to talk to ZGW (§8.1, #78). Uses a scripted message handler so no
/// real ACL is required.
/// </summary>
public class AclHttpClientTests
{
private const string ObjectUrl = "http://objecten.local:8000/api/v2/objects/obj-9";
private static AclHttpClient Client(StubHandler handler) =>
new(new HttpClient(handler) { BaseAddress = new Uri("http://acl/") });
[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 client = Client(capture.Responds(
HttpStatusCode.OK, """{"id":"zaak-1","status":"INGESCHREVEN","reference":"REG-42"}"""));
var client = Client(capture.Responds(HttpStatusCode.OK, """{"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("INGESCHREVEN", record.Status);
Assert.Equal("REG-42", record.Reference);
Assert.Equal("REG-42", reference);
Assert.Equal(HttpMethod.Post, capture.Seen!.Method);
Assert.Equal("http://acl/register-records/read", capture.Seen.RequestUri!.ToString());
Assert.Contains($"\"objectUrl\":\"{ObjectUrl}\"", 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)));
Assert.Equal("http://acl/zaken/reference", capture.Seen.RequestUri!.ToString());
Assert.Contains("\"zaakUrl\":\"http://openzaak/zaken/api/v1/zaken/abc\"", capture.Body);
}
[Fact]
@@ -50,7 +35,7 @@ public class AclHttpClientTests
var client = Client(capture.Responds(HttpStatusCode.BadGateway));
await Assert.ThrowsAsync<HttpRequestException>(
() => client.GetRegisterRecordAsync(new Uri(ObjectUrl)));
() => client.GetZaakReferenceAsync(new Uri("http://openzaak/zaken/api/v1/zaken/abc")));
}
[Fact]
@@ -60,17 +45,17 @@ public class AclHttpClientTests
var client = Client(capture.Responds(HttpStatusCode.OK, "null"));
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);
}
[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 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);
}
}
@@ -5,18 +5,16 @@ namespace EventSubscriber.Tests;
/// <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
/// convention — no mocking library).</summary>
/// <summary>A fake ACL client standing in for the register records Objecten holds: a test seeds a
/// record per object URL, and the call count proves a rebuild does not re-read through the ACL.</summary>
/// <summary>A fake ACL client that returns a fixed reference derived from the zaak, and records
/// how many times it was called (to prove a rebuild does not re-read via the ACL).</summary>
internal sealed class FakeAclClient : IAclClient
{
public Dictionary<string, RegisterRecord> Records { get; } = [];
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++;
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;
/// <summary>Behaviour of the projector that turns NRC notifications into projection rows. Since
/// S-19b-2 the source is the register in Objecten (ADR-0030), not ZGW zaak events: a notification
/// carries only the object URL, so the record is read back through the ACL. Duplicate and
/// out-of-order deliveries must be tolerated (CLAUDE.md §8.6).</summary>
/// <summary>Behaviour of the projector that turns NRC notifications into projection rows.
/// The walking skeleton reacts only to a zaak being created (status INGEDIEND) and must
/// tolerate duplicate and out-of-order deliveries (CLAUDE.md §8.6).</summary>
public sealed class NotificationProjectorTests
{
private const string ObjectUrl = "http://objecten.local:8000/api/v2/objects/11111111-1111-1111-1111-111111111111";
private const string ZaakId = "99999999-9999-9999-9999-999999999999";
private const string ZaakUrl = "http://openzaak:8000/zaken/api/v1/zaken/11111111-1111-1111-1111-111111111111";
private const string StatusUrl = "http://openzaak:8000/zaken/api/v1/statussen/22222222-2222-2222-2222-222222222222";
private readonly InMemoryNotificationLog _log = new();
private readonly InMemoryProjectionStore _store = new();
@@ -17,60 +16,46 @@ public sealed class NotificationProjectorTests
private NotificationProjector Projector() => new(_log, _store, _acl);
/// <summary>A register write as Objecten publishes it: the object is both hoofdObject and
/// resourceUrl, and the record itself is only reachable by reading that object.</summary>
private Notification RecordWritten(string actie = "create", string url = ObjectUrl, string status = RegistrationStatus.Ingediend, string zaakId = ZaakId)
private static Notification ZaakCreated(string url = ZaakUrl)
=> new("zaken", "zaak", "create", new Uri(url));
// 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");
return new Notification("objecten", "object", actie, new Uri(url));
await Projector().HandleAsync(ZaakCreated());
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]
public async Task a_register_record_write_is_projected_as_a_row_keyed_on_the_registration()
{
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)
public async Task rebuild_reproduces_the_reference_without_re_reading_via_the_acl()
{
var projector = Projector();
await projector.HandleAsync(RecordWritten());
await projector.HandleAsync(RecordWritten(actie, status: RegistrationStatus.Ingeschreven));
await projector.HandleAsync(ZaakCreated());
var callsAfterProjection = _acl.CallCount;
await projector.RebuildAsync();
var entry = Assert.Single(await _store.AllAsync());
Assert.Equal(ZaakId, entry.Id);
Assert.Equal(RegistrationStatus.Ingeschreven, entry.Status);
}
[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());
Assert.Equal("REG-11111111-1111-1111-1111-111111111111", entry.Reference);
// Rebuild replays the log (which stored the reference) — no extra ACL calls (#78, ADR-0008).
Assert.Equal(callsAfterProjection, _acl.CallCount);
}
[Fact]
public async Task replaying_the_same_notification_keeps_a_single_row()
{
var projector = Projector();
await projector.HandleAsync(RecordWritten());
await projector.HandleAsync(RecordWritten());
await projector.HandleAsync(ZaakCreated());
await projector.HandleAsync(ZaakCreated());
Assert.Single(await _store.AllAsync());
}
@@ -79,8 +64,8 @@ public sealed class NotificationProjectorTests
public async Task a_replayed_notification_never_reaches_the_projection_store()
{
var projector = Projector();
await projector.HandleAsync(RecordWritten());
await projector.HandleAsync(RecordWritten());
await projector.HandleAsync(ZaakCreated());
await projector.HandleAsync(ZaakCreated());
// 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.
@@ -88,59 +73,77 @@ public sealed class NotificationProjectorTests
}
[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();
await projector.HandleAsync(RecordWritten());
await projector.HandleAsync(RecordWritten(url: ObjectUrl[..^1] + "2", zaakId: "other-zaak"));
await projector.HandleAsync(ZaakCreated());
await projector.HandleAsync(ZaakCreated(ZaakUrl[..^1] + "2")); // a distinct zaak url
Assert.Equal(2, (await _store.AllAsync()).Count);
}
[Theory]
[InlineData("zaken", "zaak", "create")] // the ZGW source S-19b-2 replaced
[InlineData("zaken", "status", "create")] // ditto
[InlineData("objecten", "object", "destroy")] // a delete we do not project
[InlineData("documenten", "object", "create")] // wrong kanaal
[InlineData("documenten", "enkelvoudiginformatieobject", "create")] // wrong kanaal + resource
[InlineData("documenten", "zaak", "create")] // wrong kanaal only
[InlineData("zaken", "zaak", "update")] // wrong actie
[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)
{
_acl.Records[ObjectUrl] = new RegisterRecord(ZaakId, RegistrationStatus.Ingediend, "REG-2026-0001");
await Projector().HandleAsync(new Notification(kanaal, resource, actie, new Uri(ObjectUrl)));
await Projector().HandleAsync(new Notification(kanaal, resource, actie, new Uri(ZaakUrl)));
Assert.Empty(await _store.AllAsync());
}
[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();
await projector.HandleAsync(RecordWritten());
await projector.HandleAsync(RecordWritten("partial_update", status: RegistrationStatus.Ingeschreven));
var callsAfterProjection = _acl.CallCount;
await projector.HandleAsync(ZaakCreated());
await projector.HandleAsync(StatusSet());
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();
var entry = Assert.Single(await _store.AllAsync());
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]
public async Task rebuild_clears_stale_rows_and_repopulates_from_the_notification_log()
{
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.
await _store.UpsertAsync(new RegisterEntry("stale-9999", RegistrationStatus.Ingediend));
await projector.RebuildAsync();
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);
}
}
@@ -13,8 +13,9 @@ public sealed class EfNotificationLog(ProjectionDbContext db) : INotificationLog
db.ProcessedNotifications.Add(new ProcessedNotificationRow
{
Key = notification.Key,
RegisterId = notification.RegisterId,
Status = notification.Status,
Actie = notification.Actie,
ZaakId = notification.ZaakId,
Resource = notification.Resource,
Reference = notification.Reference,
ReceivedAt = DateTimeOffset.UtcNow,
});
@@ -35,6 +36,6 @@ public sealed class EfNotificationLog(ProjectionDbContext db) : INotificationLog
public async Task<IReadOnlyList<RecordedNotification>> AllAsync(CancellationToken ct = default)
=> await db.ProcessedNotifications
.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);
}
@@ -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")
.HasColumnName("key");
b.Property<string>("Actie")
.IsRequired()
.HasColumnType("text")
.HasColumnName("actie");
b.Property<DateTimeOffset>("ReceivedAt")
.HasColumnType("timestamp with time zone")
.HasColumnName("received_at");
@@ -36,15 +41,15 @@ namespace Projection.ReadModel.Migrations
.HasColumnType("text")
.HasColumnName("reference");
b.Property<string>("RegisterId")
b.Property<string>("Resource")
.IsRequired()
.HasColumnType("text")
.HasColumnName("register_id");
.HasColumnName("resource");
b.Property<string>("Status")
b.Property<string>("ZaakId")
.IsRequired()
.HasColumnType("text")
.HasColumnName("status");
.HasColumnName("zaak_id");
b.HasKey("Key");
@@ -34,8 +34,9 @@ public sealed class ProjectionDbContext(DbContextOptions<ProjectionDbContext> op
e.ToTable("processed_notifications");
e.HasKey(r => r.Key);
e.Property(r => r.Key).HasColumnName("key");
e.Property(r => r.RegisterId).HasColumnName("register_id").IsRequired();
e.Property(r => r.Status).HasColumnName("status").IsRequired();
e.Property(r => r.Actie).HasColumnName("actie").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.ReceivedAt).HasColumnName("received_at");
});
@@ -55,20 +56,18 @@ public sealed class RegisterEntryRow
public string? NaamPlaceholder { get; set; }
}
/// <summary>An accepted notification, retained so the projection can be rebuilt without reading
/// 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>
/// <summary>An accepted notification, retained so the projection can be rebuilt without OpenZaak (§8.1).</summary>
public sealed class ProcessedNotificationRow
{
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>
public required string RegisterId { get; set; }
/// <summary>The ZGW resource (e.g. <c>zaak</c> or <c>status</c>) — retained so a rebuild reprojects
/// 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>
public required string Status { get; set; }
/// <summary>The citizen-facing reference the record carried — matches the submit confirmation (#78).</summary>
/// <summary>The zaak reference (identificatie), retained so a rebuild reprojects it without the ACL (#78).</summary>
public string? Reference { get; set; }
public DateTimeOffset ReceivedAt { get; set; }
@@ -1,28 +1,19 @@
# language: en
# Drives S-19b-2 (#153), re-sourcing S-06 (#7). The read projection is derived from the
# RegisterRecord in Objecten (ADR-0030), not from ZGW zaak events: the ACL records a registration
# in the register, Objecten notifies, and the Event Subscriber projects the record that
# notification points at. This scenario exercises the use case against in-memory stand-ins for the
# register, the projection store and the notification log; real Objecten → NRC → subscriber
# delivery is verified by the live-stack check (verify-projection, ADR-0007/0030).
Feature: Register-projectie bijwerken op een registerwijziging
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.
# Drives S-06 (#7). On a zaak-created notification from NRC the Event Subscriber writes a
# rebuildable read-projection row (PRD §8.4). This scenario exercises the use case against an
# in-memory stand-in for the projection store and notification log; real OpenZaak → NRC →
# subscriber delivery is verified by the live-stack check (verify-projection, ADR-0007/#58).
Feature: Register-projectie bijwerken op een zaaknotificatie
Als openbaar register wil ik dat een aangemaakte zaak in de projectie verschijnt
zodat het register de ingediende registratie kan tonen.
Scenario: Een ingediende registratie levert een rij met status INGEDIEND
Given registration "11111111-1111-1111-1111-111111111111" is recorded in the register with status "INGEDIEND"
When the register notification is delivered to the event subscriber
Scenario: Een zaaknotificatie levert een rij met status INGEDIEND
Given a zaak is created in OpenZaak with id "11111111-1111-1111-1111-111111111111"
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"
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
Given registration "22222222-2222-2222-2222-222222222222" is recorded in the register with status "INGEDIEND"
When the register notification is delivered to the event subscriber
And the same register notification is delivered again
Given a zaak is created in OpenZaak with id "22222222-2222-2222-2222-222222222222"
When the NRC notification for that zaak is delivered to the event subscriber
And the same NRC notification is delivered again
Then the register projection contains exactly one row for "22222222-2222-2222-2222-222222222222"
@@ -5,39 +5,31 @@ using Xunit;
namespace Acceptance.Steps;
/// <summary>Bindings for <c>RegisterProjectieBijwerken.feature</c> (S-06, re-sourced by S-19b-2).
/// Reqnroll creates one instance per scenario, so instance fields hold scenario-scoped state.</summary>
/// <summary>Bindings for <c>RegisterProjectieBijwerken.feature</c> (S-06). Reqnroll creates
/// one instance per scenario, so instance fields hold scenario-scoped state.</summary>
[Binding]
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 InMemoryProjectionStore _store = new();
private readonly InMemoryRegisterRecordClient _register = new();
private readonly NotificationProjector _projector;
private Notification? _notification;
public RegisterProjectieBijwerkenSteps()
=> _projector = new NotificationProjector(_log, _store, _register);
=> _projector = new NotificationProjector(_log, _store, new InMemoryAclReferenceClient());
[Given("registration \"(.*)\" is recorded in the register with status \"(.*)\"")]
[When("registration \"(.*)\" is recorded in the register with status \"(.*)\"")]
public void RegistrationIsRecorded(string id, string status)
{
// 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("a zaak is created in OpenZaak with id \"(.*)\"")]
public void GivenAZaakIsCreatedInOpenZaakWithId(string id)
=> _notification = new Notification("zaken", "zaak", "create", new Uri(ZaakBase + id));
[Given("the register notification is delivered to the event subscriber")]
[When("the register notification is delivered to the event subscriber")]
public Task TheNotificationIsDelivered()
[When("the NRC notification for that zaak is delivered to the event subscriber")]
public Task WhenTheNotificationIsDelivered()
=> _projector.HandleAsync(_notification!);
[When("the same register notification is delivered again")]
public Task TheSameNotificationIsDeliveredAgain()
[When("the same NRC notification is delivered again")]
public Task WhenTheSameNotificationIsDeliveredAgain()
=> _projector.HandleAsync(_notification!);
[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)];
}
/// <summary>An in-memory stand-in for the register the ACL reads back for the projector, so the
/// scenario runs without a running ACL or Objecten (S-19b-2, ADR-0030).</summary>
public sealed class InMemoryRegisterRecordClient : IAclClient
/// <summary>A fake ACL client for the projection acceptance scenario: returns a reference derived
/// from the zaak, so the projector can enrich rows without a running ACL (#78).</summary>
public sealed class InMemoryAclReferenceClient : IAclClient
{
public Dictionary<string, RegisterRecord> Records { get; } = [];
public Task<RegisterRecord?> GetRegisterRecordAsync(Uri objectUrl, CancellationToken ct = default)
=> Task.FromResult(Records.TryGetValue(objectUrl.ToString(), out var record) ? record : null);
public Task<string> GetZaakReferenceAsync(Uri zaakUrl, CancellationToken ct = default)
=> Task.FromResult("REG-" + zaakUrl.Segments[^1].Trim('/'));
}
@@ -65,8 +65,4 @@ public sealed class InMemoryRegisterRecordGateway : IRegisterRecordGateway
Upserted.Add(record);
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';
// Walking-skeleton happy path (S-08d + S-09 + S-09b + S-12 + S-10a + S-19b-2): a zorgprofessional
// logs in via mock DigiD and submits through the self-service portal → BFF → domain; the entry
// appears in the openbaar register as INGEDIEND; the citizen supplies the documents the process is
// waiting for (S-10a); a behandelaar then logs in to the behandel portal, finds the registration in
// the werkbak, and approves it (goedkeuren); the decision completes the Flowable Beoordelen task and
// flows via the ACL → Objecten → NRC → event-subscriber → projection, and the openbaar register
// 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.
// Walking-skeleton happy path (S-08d + S-09 + S-09b + S-12 + S-10a): a zorgprofessional logs in via
// mock DigiD and submits through the self-service portal → BFF → domain; the entry appears in the
// openbaar register as INGEDIEND; the citizen supplies the documents the process is waiting for
// (S-10a); a behandelaar then logs in to the behandel portal, finds the registration in the werkbak,
// and approves it (goedkeuren); the decision completes the Flowable Beoordelen task and flows via the
// ACL → NRC → event-subscriber → projection, and the openbaar register shows INGESCHREVEN.
test('DigiD submit → public INGEDIEND → documenten → behandelaar goedkeurt → public INGESCHREVEN', async ({
page,
context,