## What & why S-19b-2, closing out ADR-0028's stated direction: **the read projection is now derived from the `RegisterRecord` in Objecten, not from ZGW zaak events.** Until now the subscriber listened on `zaken` and *inferred* register state from case events — a `zaak/create` meant INGEDIEND, and any `status/create` was assumed to be the approval (it may not read OpenZaak, so it could not tell statustypen apart). The reference wasn't in the notification at all, so every projection made a second hop to the ACL. The register — a fact about a person — was being reconstructed by guessing at the lifecycle of the case that produced it. - The subscriber's abonnement moves to the `objecten` kanaal (S-19b-1 made it publish). - An Objecten notification carries **no record data**, only the object URL, so the record is read back through the ACL (`POST /register-records/read`) — §8.1 applies to Objecten exactly as ADR-0028 established. - The record carries `id`, `status` and `reference`, so the row *is* the record: `IsZaakCreated`, `IsZaakStatusSet`, `ZaakUrl`, `ZaakId` and `ToEntry`'s `Resource == "status"` inference are all gone, and so is the ACL enrichment hop. - **The ACL now writes an INGEDIEND record on submit.** Without it, re-sourcing would silently drop every submitted registration from the public register, since only approval wrote a record. - `processed_notifications` holds the projected row (`register_id`, `status`, `reference`) instead of the ZGW event, so a rebuild is a replay with no mapping rules and no upstream reads at all. **ADR-0030** records it. ADR-0028's open caveat — record written but not yet read, "the two must agree" — is closed: there is one source now. Closes #153 ## Definition of Done - [x] Linked Gitea issue (above). - [x] Failing tests committed before the implementation — two red/green pairs, ACL side (06c0444→566ef7d) and subscriber side (142ed45→8af09b2). - [x] Refactor commit follows (b496ac9). - [x] Conventional Commits referencing the issue (`refs #153`). - [x] CI green — all six jobs onb30fa66, `verify-stack` end to end including the e2e. - [x] `docker compose up` from a fresh clone reaches green health checks within 3 minutes (`verify-stack`'s bring-up step — see the wait-healthy fix below). - [x] Docs updated — ADR-0030 added, ADR-0028's consequence + caveat annotated, BACKLOG.md, e2e header comment. - [x] ADR added in `docs/architecture/`. - [x] Demo note in `docs/demo-script.md` — n/a: no user-visible change. The openbaar register shows the same two statuses for the same registrations; only where they come from changed. ## Notes for reviewers **The decision I'd most like a second opinion on** is the one the issue didn't settle: what happens to INGEDIEND. Objecten held only INGESCHREVEN records, so re-sourcing forced a choice between (a) the ACL also writing on submit, (b) a public register that lists only actual registrations, or (c) a hybrid keeping both kanalen. I took (a): visible behaviour is unchanged and the register holds the whole lifecycle. (b) is arguably the better *semantics* for a public register but narrows what the portal shows and reads against PRD §68 ("~50 register entries with diverse statuses"); (c) leaves the projection half-derived from ZGW, which is the coupling ADR-0028 set out to remove. All three are laid out in ADR-0030. **The dedup key is the projected row**, `objecten:object:{url}:{status}:{reference}` — not the object URL (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) and not URL+actie (a retried approval is a second `update`). Redeliveries collapse, genuine state changes don't. §8.6. **The migration drops columns rather than renaming them.** EF scaffolded renames — `resource` → `register_id`, `zaak_id` → `status` — which would have carried ZGW values into columns meaning something else, and a rebuild would then have projected that garbage. It also empties both tables: a pre-slice row describes a zaak event the new projector can't reproject, and those registrations have no RegisterRecord in Objecten either, so they're not re-derivable from the new source. Stated as a ceiling in the ADR — fine while stacks are ephemeral, backfill from Objecten if a long-lived environment ever needs it. **`run-projection-check.sh` now opens its zaak through the ACL** instead of straight against OpenZaak, because the ACL is what writes the record. A zaak created behind the ACL's back produces no projection row — that's the re-source working, not a gap. ## Three fixes CI found, none of them in the projection logic 1. **`wait-healthy.sh` matched the wrong container** (744f91a). Bring-up timed out with `TIMEOUT: 'objecten' not healthy (status=none)` while the `docker ps` it dumps showed objecten `Up 9 minutes (healthy)`. `--filter name=` is a substring match, so `objecten` also matches `objecten-db`/`objecten-redis`/`objecten-celery`, and `head -1` took whichever docker listed first — the celery worker has no healthcheck, hence `status=none`. Latent since those services landed and decided purely by listing order; `objecttypen` matches `objecttypen-db` the same way. Anchored on the compose replica suffix, which the verify scripts already do. 2. **The ACL had to be repointed at OpenZaak's IP** (7e0897a). Opening the zaak through the ACL put this check in the same bind run-domain-check.sh already handles: `400 {"name":"zaaktype","code":"bad-url","reason":"Voer een geldige URL in."}`. OpenZaak reflects the request Host into the zaaktype URL and then rejects it on zaak-create when single-label — the mechanism compose already documents on `ACL_OPENZAAK_BASEURL`. 3. **Approval arrives as `partial_update`, not `update`** (0dd26a7→b30fa66) — the one real bug in the slice. The ACL upserts with PATCH; DRF routes it through the notifying `update()` but names the action `partial_update`, so the projector dropped every approval. Only the e2e could catch it: `verify-projection` drives a submit, and per ADR-0028 the e2e is the only check that drives a *real* approval. `verify-tracing` also failed once (run 722) on a path this PR doesn't touch, and passed on a plain re-run of the same commit. Tempo logged `pusher failed to consume trace data` / `distributor_pool failing healthcheck` — it dropped spans under runner load rather than the trace chain being broken. Filed as **#156** rather than absorbed here. **Correction to the #152 PR notes:** I wrote there that celery concurrency was "the next knob" if verify-stack got tight. It isn't — `CELERY_WORKER_CONCURRENCY` already defaults to 1 in the Maykin image, so `objecten-celery` is already a single-process worker. Noted in #156. **Possible follow-up, deliberately not done here:** an `openzaak.local` network alias mirroring `objecten.local` would remove the ACL-repoint dance from both run-domain-check.sh and run-projection-check.sh. It changes the host in every zaak URL the system produces, which is too broad a ripple to land inside an unrelated slice — worth its own issue. **Known costs, all in the ADR:** submission is now two writes across two modules and eventually consistent (same posture ADR-0028 accepted for approval); projecting now depends on the ACL being reachable on the main path, not just for enrichment (NRC retries, so it converges); and OpenZaak still publishes to `zaken` with nothing in the product listening — kept because `verify-nrc` asserts that path.Reviewed-on: #155
This commit was merged in pull request #155.
This commit is contained in:
@@ -90,6 +90,16 @@ 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) =>
|
||||
@@ -131,6 +141,9 @@ 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;
|
||||
|
||||
@@ -24,7 +24,16 @@ public sealed class AclService(
|
||||
clock.Today,
|
||||
registration.Reference);
|
||||
|
||||
return await gateway.OpenZaakAsync(request, ct);
|
||||
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;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
@@ -52,6 +61,18 @@ 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,6 +13,14 @@ 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,3 +1,4 @@
|
||||
using System.Net;
|
||||
using System.Net.Http.Headers;
|
||||
using System.Net.Http.Json;
|
||||
using System.Text.Json.Serialization;
|
||||
@@ -38,6 +39,30 @@ 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"));
|
||||
|
||||
@@ -141,6 +166,12 @@ 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);
|
||||
|
||||
@@ -79,11 +79,21 @@ 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()
|
||||
@@ -130,6 +140,52 @@ 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,6 +83,43 @@ 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,3 +1,4 @@
|
||||
using System.Net;
|
||||
using System.Net.Http.Json;
|
||||
using System.Text.Json.Serialization;
|
||||
using EventSubscriber.Application;
|
||||
@@ -5,26 +6,28 @@ using EventSubscriber.Application;
|
||||
namespace EventSubscriber.Api;
|
||||
|
||||
/// <summary>
|
||||
/// 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).
|
||||
/// 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.
|
||||
/// </summary>
|
||||
public sealed class AclHttpClient(HttpClient http) : IAclClient
|
||||
{
|
||||
public async Task<string> GetZaakReferenceAsync(Uri zaakUrl, CancellationToken ct = default)
|
||||
public async Task<RegisterRecord?> GetRegisterRecordAsync(Uri objectUrl, CancellationToken ct = default)
|
||||
{
|
||||
ArgumentNullException.ThrowIfNull(zaakUrl);
|
||||
ArgumentNullException.ThrowIfNull(objectUrl);
|
||||
|
||||
using var response = await http.PostAsJsonAsync(
|
||||
new Uri(http.BaseAddress!, "zaken/reference"), new ReferenceRequest(zaakUrl.ToString()), ct);
|
||||
response.EnsureSuccessStatusCode();
|
||||
new Uri(http.BaseAddress!, "register-records/read"),
|
||||
new ReadRequest(objectUrl.ToString()), ct);
|
||||
|
||||
var body = await response.Content.ReadFromJsonAsync<ReferenceResponse>(ct)
|
||||
?? throw new InvalidOperationException("The ACL returned an empty reference response.");
|
||||
return body.Reference;
|
||||
// The object holds no register record (deleted, or never one) — nothing to project (§8.6).
|
||||
if (response.StatusCode == HttpStatusCode.NotFound)
|
||||
return null;
|
||||
|
||||
response.EnsureSuccessStatusCode();
|
||||
return await response.Content.ReadFromJsonAsync<RegisterRecord>(ct)
|
||||
?? throw new InvalidOperationException("The ACL returned an empty register record response.");
|
||||
}
|
||||
|
||||
private sealed record ReferenceRequest([property: JsonPropertyName("zaakUrl")] string ZaakUrl);
|
||||
|
||||
private sealed record ReferenceResponse([property: JsonPropertyName("reference")] string Reference);
|
||||
private sealed record ReadRequest([property: JsonPropertyName("objectUrl")] string ObjectUrl);
|
||||
}
|
||||
|
||||
@@ -84,11 +84,12 @@ 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
|
||||
/// 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)
|
||||
/// <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)
|
||||
{
|
||||
public Notification ToNotification() => new(Kanaal, Resource, Actie, ResourceUrl, HoofdObject);
|
||||
public Notification ToNotification() => new(Kanaal, Resource, Actie, ResourceUrl);
|
||||
}
|
||||
|
||||
public partial class Program
|
||||
|
||||
@@ -2,40 +2,39 @@ 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; 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.
|
||||
/// abonnement callback. Only the fields the projection needs are modelled.
|
||||
/// </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? HoofdObject = null)
|
||||
Uri ResourceUrl)
|
||||
{
|
||||
/// <summary>A zaak being created — projected as INGEDIEND.</summary>
|
||||
public bool IsZaakCreated =>
|
||||
Kanaal == "zaken" && Resource == "zaak" && Actie == "create";
|
||||
|
||||
/// <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.)
|
||||
/// 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>
|
||||
public string IdempotencyKey => $"{Kanaal}:{Resource}:{Actie}:{ResourceUrl}";
|
||||
/// <remarks>
|
||||
/// <c>partial_update</c> is what a PATCH actually reports: DRF routes it through the notifying
|
||||
/// <c>update()</c> but names the action <c>partial_update</c>, and that is what Objecten puts in
|
||||
/// the notification. <c>update</c> is accepted too, so a PUT-shaped write would project the same
|
||||
/// way. <c>destroy</c> is deliberately not: removing a registration from the public register is
|
||||
/// its own decision, not a side effect of this one.
|
||||
/// </remarks>
|
||||
public bool IsRegisterRecordWritten =>
|
||||
Kanaal == "objecten" && Resource == "object"
|
||||
&& Actie is "create" or "update" or "partial_update";
|
||||
|
||||
/// <summary>The object holding the register record. For a <c>resource: object</c> notification
|
||||
/// 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;
|
||||
}
|
||||
|
||||
@@ -3,21 +3,27 @@ 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 zaak id. Rebuilds the projection by replaying the log.
|
||||
/// upsert is idempotent on the register 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 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>
|
||||
/// <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>
|
||||
public async Task HandleAsync(Notification notification, CancellationToken ct = default)
|
||||
{
|
||||
if (!notification.IsZaakCreated && !notification.IsZaakStatusSet)
|
||||
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)
|
||||
return;
|
||||
|
||||
var reference = await acl.GetZaakReferenceAsync(notification.ZaakUrl, ct);
|
||||
var recorded = new RecordedNotification(
|
||||
notification.IdempotencyKey, notification.Actie, notification.ZaakId, notification.Resource, reference);
|
||||
KeyFor(notification.ObjectUrl, record), record.Id, record.Status, record.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.
|
||||
@@ -27,6 +33,20 @@ 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)
|
||||
{
|
||||
@@ -35,11 +55,9 @@ public sealed class NotificationProjector(INotificationLog log, IProjectionStore
|
||||
await store.UpsertAsync(ToEntry(recorded), ct);
|
||||
}
|
||||
|
||||
/// <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>
|
||||
/// <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>
|
||||
private static RegisterEntry ToEntry(RecordedNotification recorded)
|
||||
=> new(
|
||||
recorded.ZaakId,
|
||||
recorded.Resource == "status" ? RegistrationStatus.Ingeschreven : RegistrationStatus.Ingediend,
|
||||
Reference: recorded.Reference);
|
||||
=> new(recorded.RegisterId, recorded.Status, 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 OpenZaak (CLAUDE.md §8.1). Implemented in Infrastructure over Postgres.
|
||||
/// needs no access to Objecten or ZGW (CLAUDE.md §8.1). Implemented in Infrastructure over Postgres.
|
||||
/// </summary>
|
||||
public interface INotificationLog
|
||||
{
|
||||
@@ -19,22 +19,29 @@ public interface INotificationLog
|
||||
Task<IReadOnlyList<RecordedNotification>> AllAsync(CancellationToken ct = default);
|
||||
}
|
||||
|
||||
/// <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>
|
||||
/// 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>
|
||||
/// 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).
|
||||
/// 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.
|
||||
/// </summary>
|
||||
public interface IAclClient
|
||||
{
|
||||
/// <summary>The zaak's reference (identificatie) for the read projection.</summary>
|
||||
Task<string> GetZaakReferenceAsync(Uri zaakUrl, CancellationToken ct = default);
|
||||
/// <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 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,27 +5,42 @@ using EventSubscriber.Api;
|
||||
namespace EventSubscriber.Tests;
|
||||
|
||||
/// <summary>
|
||||
/// 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.
|
||||
/// 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.
|
||||
/// </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_zaak_reference_by_posting_the_zaak_url_and_returns_it()
|
||||
public async Task Reads_a_register_record_by_posting_the_object_url()
|
||||
{
|
||||
var capture = new RequestCapture();
|
||||
var client = Client(capture.Responds(HttpStatusCode.OK, """{"reference":"REG-42"}"""));
|
||||
var client = Client(capture.Responds(
|
||||
HttpStatusCode.OK, """{"id":"zaak-1","status":"INGESCHREVEN","reference":"REG-42"}"""));
|
||||
|
||||
var reference = await client.GetZaakReferenceAsync(new Uri("http://openzaak/zaken/api/v1/zaken/abc"));
|
||||
var record = await client.GetRegisterRecordAsync(new Uri(ObjectUrl));
|
||||
|
||||
Assert.Equal("REG-42", reference);
|
||||
Assert.Equal("zaak-1", record!.Id);
|
||||
Assert.Equal("INGESCHREVEN", record.Status);
|
||||
Assert.Equal("REG-42", record.Reference);
|
||||
Assert.Equal(HttpMethod.Post, capture.Seen!.Method);
|
||||
Assert.Equal("http://acl/zaken/reference", capture.Seen.RequestUri!.ToString());
|
||||
Assert.Contains("\"zaakUrl\":\"http://openzaak/zaken/api/v1/zaken/abc\"", capture.Body);
|
||||
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)));
|
||||
}
|
||||
|
||||
[Fact]
|
||||
@@ -35,7 +50,7 @@ public class AclHttpClientTests
|
||||
var client = Client(capture.Responds(HttpStatusCode.BadGateway));
|
||||
|
||||
await Assert.ThrowsAsync<HttpRequestException>(
|
||||
() => client.GetZaakReferenceAsync(new Uri("http://openzaak/zaken/api/v1/zaken/abc")));
|
||||
() => client.GetRegisterRecordAsync(new Uri(ObjectUrl)));
|
||||
}
|
||||
|
||||
[Fact]
|
||||
@@ -45,17 +60,17 @@ public class AclHttpClientTests
|
||||
var client = Client(capture.Responds(HttpStatusCode.OK, "null"));
|
||||
|
||||
var ex = await Assert.ThrowsAsync<InvalidOperationException>(
|
||||
() => client.GetZaakReferenceAsync(new Uri("http://openzaak/zaken/api/v1/zaken/abc")));
|
||||
() => client.GetRegisterRecordAsync(new Uri(ObjectUrl)));
|
||||
Assert.Contains("empty", ex.Message, StringComparison.OrdinalIgnoreCase);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task Rejects_a_null_zaak_url_without_sending_a_request()
|
||||
public async Task Rejects_a_null_object_url_without_sending_a_request()
|
||||
{
|
||||
var capture = new RequestCapture();
|
||||
var client = Client(capture.Responds(HttpStatusCode.OK, """{"reference":"REG-1"}"""));
|
||||
var client = Client(capture.Responds(HttpStatusCode.OK, "{}"));
|
||||
|
||||
await Assert.ThrowsAsync<ArgumentNullException>(() => client.GetZaakReferenceAsync(null!));
|
||||
await Assert.ThrowsAsync<ArgumentNullException>(() => client.GetRegisterRecordAsync(null!));
|
||||
Assert.Null(capture.Seen);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -5,16 +5,18 @@ 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 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>
|
||||
/// <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>
|
||||
internal sealed class FakeAclClient : IAclClient
|
||||
{
|
||||
public Dictionary<string, RegisterRecord> Records { get; } = [];
|
||||
|
||||
public int CallCount { get; private set; }
|
||||
|
||||
public Task<string> GetZaakReferenceAsync(Uri zaakUrl, CancellationToken ct = default)
|
||||
public Task<RegisterRecord?> GetRegisterRecordAsync(Uri objectUrl, CancellationToken ct = default)
|
||||
{
|
||||
CallCount++;
|
||||
return Task.FromResult("REG-" + zaakUrl.Segments[^1].Trim('/'));
|
||||
return Task.FromResult(Records.TryGetValue(objectUrl.ToString(), out var record) ? record : null);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -2,13 +2,14 @@ using EventSubscriber.Application;
|
||||
|
||||
namespace EventSubscriber.Tests;
|
||||
|
||||
/// <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>
|
||||
/// <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>
|
||||
public sealed class NotificationProjectorTests
|
||||
{
|
||||
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 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 readonly InMemoryNotificationLog _log = new();
|
||||
private readonly InMemoryProjectionStore _store = new();
|
||||
@@ -16,46 +17,60 @@ public sealed class NotificationProjectorTests
|
||||
|
||||
private NotificationProjector Projector() => new(_log, _store, _acl);
|
||||
|
||||
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()
|
||||
/// <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)
|
||||
{
|
||||
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);
|
||||
_acl.Records[url] = new RegisterRecord(zaakId, status, "REG-2026-0001");
|
||||
return new Notification("objecten", "object", actie, new Uri(url));
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task rebuild_reproduces_the_reference_without_re_reading_via_the_acl()
|
||||
public async Task a_register_record_write_is_projected_as_a_row_keyed_on_the_registration()
|
||||
{
|
||||
var projector = Projector();
|
||||
await projector.HandleAsync(ZaakCreated());
|
||||
var callsAfterProjection = _acl.CallCount;
|
||||
|
||||
await projector.RebuildAsync();
|
||||
await Projector().HandleAsync(RecordWritten());
|
||||
|
||||
var entry = Assert.Single(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);
|
||||
// Keyed on the record's own id (the zaak id), not on the Objecten object's uuid — the
|
||||
// projection row and the register record are the same registration.
|
||||
Assert.Equal(ZaakId, entry.Id);
|
||||
Assert.Equal(RegistrationStatus.Ingediend, entry.Status);
|
||||
Assert.Equal("REG-2026-0001", entry.Reference);
|
||||
}
|
||||
|
||||
// The ACL PATCHes the same object on approval. DRF routes a PATCH through `update()` but reports
|
||||
// the action as `partial_update`, which is what Objecten puts in the notification — so accepting
|
||||
// only `create`/`update` silently drops every approval.
|
||||
[Theory]
|
||||
[InlineData("partial_update")]
|
||||
[InlineData("update")]
|
||||
public async Task approval_updates_the_same_row_from_ingediend_to_ingeschreven(string actie)
|
||||
{
|
||||
var projector = Projector();
|
||||
await projector.HandleAsync(RecordWritten());
|
||||
await projector.HandleAsync(RecordWritten(actie, status: RegistrationStatus.Ingeschreven));
|
||||
|
||||
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());
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task replaying_the_same_notification_keeps_a_single_row()
|
||||
{
|
||||
var projector = Projector();
|
||||
await projector.HandleAsync(ZaakCreated());
|
||||
await projector.HandleAsync(ZaakCreated());
|
||||
await projector.HandleAsync(RecordWritten());
|
||||
await projector.HandleAsync(RecordWritten());
|
||||
|
||||
Assert.Single(await _store.AllAsync());
|
||||
}
|
||||
@@ -64,8 +79,8 @@ public sealed class NotificationProjectorTests
|
||||
public async Task a_replayed_notification_never_reaches_the_projection_store()
|
||||
{
|
||||
var projector = Projector();
|
||||
await projector.HandleAsync(ZaakCreated());
|
||||
await projector.HandleAsync(ZaakCreated());
|
||||
await projector.HandleAsync(RecordWritten());
|
||||
await projector.HandleAsync(RecordWritten());
|
||||
|
||||
// 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.
|
||||
@@ -73,77 +88,59 @@ public sealed class NotificationProjectorTests
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task two_different_zaken_each_get_their_own_row()
|
||||
public async Task two_different_registrations_each_get_their_own_row()
|
||||
{
|
||||
var projector = Projector();
|
||||
await projector.HandleAsync(ZaakCreated());
|
||||
await projector.HandleAsync(ZaakCreated(ZaakUrl[..^1] + "2")); // a distinct zaak url
|
||||
await projector.HandleAsync(RecordWritten());
|
||||
await projector.HandleAsync(RecordWritten(url: ObjectUrl[..^1] + "2", zaakId: "other-zaak"));
|
||||
|
||||
Assert.Equal(2, (await _store.AllAsync()).Count);
|
||||
}
|
||||
|
||||
[Theory]
|
||||
[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
|
||||
[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
|
||||
public async Task an_unrelated_notification_is_not_projected(string kanaal, string resource, string actie)
|
||||
{
|
||||
await Projector().HandleAsync(new Notification(kanaal, resource, actie, new Uri(ZaakUrl)));
|
||||
_acl.Records[ObjectUrl] = new RegisterRecord(ZaakId, RegistrationStatus.Ingediend, "REG-2026-0001");
|
||||
|
||||
await Projector().HandleAsync(new Notification(kanaal, resource, actie, new Uri(ObjectUrl)));
|
||||
|
||||
Assert.Empty(await _store.AllAsync());
|
||||
}
|
||||
|
||||
[Fact]
|
||||
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()
|
||||
public async Task rebuild_reproduces_the_row_without_re_reading_through_the_acl()
|
||||
{
|
||||
var projector = Projector();
|
||||
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.HandleAsync(RecordWritten());
|
||||
await projector.HandleAsync(RecordWritten("partial_update", status: RegistrationStatus.Ingeschreven));
|
||||
var callsAfterProjection = _acl.CallCount;
|
||||
|
||||
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(ZaakCreated());
|
||||
await projector.HandleAsync(RecordWritten());
|
||||
// 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("11111111-1111-1111-1111-111111111111", entry.Id);
|
||||
Assert.Equal(ZaakId, entry.Id);
|
||||
Assert.Equal(RegistrationStatus.Ingediend, entry.Status);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -13,9 +13,8 @@ public sealed class EfNotificationLog(ProjectionDbContext db) : INotificationLog
|
||||
db.ProcessedNotifications.Add(new ProcessedNotificationRow
|
||||
{
|
||||
Key = notification.Key,
|
||||
Actie = notification.Actie,
|
||||
ZaakId = notification.ZaakId,
|
||||
Resource = notification.Resource,
|
||||
RegisterId = notification.RegisterId,
|
||||
Status = notification.Status,
|
||||
Reference = notification.Reference,
|
||||
ReceivedAt = DateTimeOffset.UtcNow,
|
||||
});
|
||||
@@ -36,6 +35,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.Actie, r.ZaakId, r.Resource, r.Reference))
|
||||
.Select(r => new RecordedNotification(r.Key, r.RegisterId, r.Status, r.Reference))
|
||||
.ToListAsync(ct);
|
||||
}
|
||||
|
||||
+87
@@ -0,0 +1,87 @@
|
||||
// <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
|
||||
}
|
||||
}
|
||||
}
|
||||
+69
@@ -0,0 +1,69 @@
|
||||
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: "");
|
||||
}
|
||||
}
|
||||
}
|
||||
+4
-9
@@ -28,11 +28,6 @@ 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");
|
||||
@@ -41,15 +36,15 @@ namespace Projection.ReadModel.Migrations
|
||||
.HasColumnType("text")
|
||||
.HasColumnName("reference");
|
||||
|
||||
b.Property<string>("Resource")
|
||||
b.Property<string>("RegisterId")
|
||||
.IsRequired()
|
||||
.HasColumnType("text")
|
||||
.HasColumnName("resource");
|
||||
.HasColumnName("register_id");
|
||||
|
||||
b.Property<string>("ZaakId")
|
||||
b.Property<string>("Status")
|
||||
.IsRequired()
|
||||
.HasColumnType("text")
|
||||
.HasColumnName("zaak_id");
|
||||
.HasColumnName("status");
|
||||
|
||||
b.HasKey("Key");
|
||||
|
||||
|
||||
@@ -34,9 +34,8 @@ 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.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.RegisterId).HasColumnName("register_id").IsRequired();
|
||||
e.Property(r => r.Status).HasColumnName("status").IsRequired();
|
||||
e.Property(r => r.Reference).HasColumnName("reference");
|
||||
e.Property(r => r.ReceivedAt).HasColumnName("received_at");
|
||||
});
|
||||
@@ -56,18 +55,20 @@ public sealed class RegisterEntryRow
|
||||
public string? NaamPlaceholder { get; set; }
|
||||
}
|
||||
|
||||
/// <summary>An accepted notification, retained so the projection can be rebuilt without OpenZaak (§8.1).</summary>
|
||||
/// <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>
|
||||
public sealed class ProcessedNotificationRow
|
||||
{
|
||||
public required string Key { get; set; }
|
||||
public required string Actie { get; set; }
|
||||
public required string ZaakId { 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 registration this record is for (the zaak id) — the projection row's key.</summary>
|
||||
public required string RegisterId { get; set; }
|
||||
|
||||
/// <summary>The zaak reference (identificatie), retained so a rebuild reprojects it without the ACL (#78).</summary>
|
||||
/// <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>
|
||||
public string? Reference { get; set; }
|
||||
|
||||
public DateTimeOffset ReceivedAt { get; set; }
|
||||
|
||||
Reference in New Issue
Block a user