diff --git a/services/event-subscriber/EventSubscriber.Api/AclHttpClient.cs b/services/event-subscriber/EventSubscriber.Api/AclHttpClient.cs index c24bad2..d280ed3 100644 --- a/services/event-subscriber/EventSubscriber.Api/AclHttpClient.cs +++ b/services/event-subscriber/EventSubscriber.Api/AclHttpClient.cs @@ -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; /// -/// 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. /// public sealed class AclHttpClient(HttpClient http) : IAclClient { - public async Task GetZaakReferenceAsync(Uri zaakUrl, CancellationToken ct = default) + public async Task 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(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(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); } diff --git a/services/event-subscriber/EventSubscriber.Application/Notification.cs b/services/event-subscriber/EventSubscriber.Application/Notification.cs index ac1d8ff..76a013a 100644 --- a/services/event-subscriber/EventSubscriber.Application/Notification.cs +++ b/services/event-subscriber/EventSubscriber.Application/Notification.cs @@ -2,11 +2,16 @@ namespace EventSubscriber.Application; /// /// 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 aanmaakdatum and kenmerken which the -/// minimal projection ignores (bsn is deferred — see ADR-0008). For a zaken/zaak/create -/// notification hoofdObject and resourceUrl are both the created zaak's URL. +/// abonnement callback. Only the fields the projection needs are modelled. /// +/// +/// Since S-19b-2 the subscriber listens on the objecten kanaal, not zaken: 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 no record data — only the object URL (as both +/// hoofdObject and resourceUrl) and the objecttype as a kenmerk — so the record +/// itself is read back through the ACL. +/// public sealed record Notification( string Kanaal, string Resource, @@ -14,28 +19,12 @@ public sealed record Notification( Uri ResourceUrl, Uri? HoofdObject = null) { - /// A zaak being created — projected as INGEDIEND. - public bool IsZaakCreated => - Kanaal == "zaken" && Resource == "zaak" && Actie == "create"; + /// A register record written to Objecten — create on submit, update on + /// approval, since the ACL upserts the same object for a registration (§8.6). + public bool IsRegisterRecordWritten => + Kanaal == "objecten" && Resource == "object" && Actie is "create" or "update"; - /// 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. - public bool IsZaakStatusSet => - Kanaal == "zaken" && Resource == "status" && Actie == "create"; - - /// The zaak URL this notification concerns — hoofdObject (the zaak) for a status - /// notification, else the resource URL (which, for a zaak-create, is the zaak). - public Uri ZaakUrl => HoofdObject ?? ResourceUrl; - - /// The zaak UUID used as the projection key — the trailing segment of . - public string ZaakId => ZaakUrl.Segments[^1].Trim('/'); - - /// - /// 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.) - /// - public string IdempotencyKey => $"{Kanaal}:{Resource}:{Actie}:{ResourceUrl}"; + /// The object holding the register record. Objecten sets both fields to the object; + /// hoofdObject is the main resource by definition, so prefer it. + public Uri ObjectUrl => HoofdObject ?? ResourceUrl; } diff --git a/services/event-subscriber/EventSubscriber.Application/NotificationProjector.cs b/services/event-subscriber/EventSubscriber.Application/NotificationProjector.cs index 2d08908..90098c1 100644 --- a/services/event-subscriber/EventSubscriber.Application/NotificationProjector.cs +++ b/services/event-subscriber/EventSubscriber.Application/NotificationProjector.cs @@ -3,28 +3,22 @@ namespace EventSubscriber.Application; /// /// 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. /// public sealed class NotificationProjector(INotificationLog log, IProjectionStore store, IAclClient acl) { - /// 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). + /// Handle one inbound notification. Reacts to a register record being written to + /// Objecten (S-19b-2, ADR-0030) and ignores everything else. public async Task HandleAsync(Notification notification, CancellationToken ct = default) { - if (!notification.IsZaakCreated && !notification.IsZaakStatusSet) + ArgumentNullException.ThrowIfNull(notification); + + if (!notification.IsRegisterRecordWritten) return; - var reference = await acl.GetZaakReferenceAsync(notification.ZaakUrl, ct); - var recorded = new RecordedNotification( - 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. - if (!await log.TryRecordAsync(recorded, ct)) - return; - - await store.UpsertAsync(ToEntry(recorded), ct); + // S-19b-2: reading the record back through the ACL and projecting it lands with the + // implementation; today nothing reaches the store. + await Task.CompletedTask; } /// Rebuild the projection from the durable notification log (PRD §8.4). @@ -35,11 +29,9 @@ public sealed class NotificationProjector(INotificationLog log, IProjectionStore await store.UpsertAsync(ToEntry(recorded), ct); } - /// The projection row for an accepted notification: a status-set maps to INGESCHREVEN, - /// a zaak-create to INGEDIEND. bsn/naam are deferred (ADR-0008). + /// 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). 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); } diff --git a/services/event-subscriber/EventSubscriber.Application/Ports.cs b/services/event-subscriber/EventSubscriber.Application/Ports.cs index 7e73309..7bea81b 100644 --- a/services/event-subscriber/EventSubscriber.Application/Ports.cs +++ b/services/event-subscriber/EventSubscriber.Application/Ports.cs @@ -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. /// public interface INotificationLog { @@ -19,22 +19,29 @@ public interface INotificationLog Task> AllAsync(CancellationToken ct = default); } -/// A notification that has been accepted, retaining what a rebuild needs to recompute its -/// projection row — the ZGW resource (zaak-create → INGEDIEND vs status-set → INGESCHREVEN) and -/// the zaak reference (identificatie), so a rebuild reproduces the row without re-reading ZGW (#78). -public sealed record RecordedNotification(string Key, string Actie, string ZaakId, string Resource, string? Reference); +/// +/// 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). +/// +public sealed record RecordedNotification(string Key, string RegisterId, string Status, string? Reference); /// -/// 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. /// public interface IAclClient { - /// The zaak's reference (identificatie) for the read projection. - Task GetZaakReferenceAsync(Uri zaakUrl, CancellationToken ct = default); + /// The register record the object at holds, or + /// null if it holds none — the object may be gone by the time a redelivered + /// notification is handled, which is not an error (§8.6). + Task GetRegisterRecordAsync(Uri objectUrl, CancellationToken ct = default); } +/// 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. +public sealed record RegisterRecord(string Id, string Status, string? Reference); + /// The read projection store. Owned by the projection bounded context (ADR-0008); the /// subscriber writes to it and the projection-api reads it. public interface IProjectionStore diff --git a/services/event-subscriber/EventSubscriber.Tests/AclHttpClientTests.cs b/services/event-subscriber/EventSubscriber.Tests/AclHttpClientTests.cs index ef170df..830546b 100644 --- a/services/event-subscriber/EventSubscriber.Tests/AclHttpClientTests.cs +++ b/services/event-subscriber/EventSubscriber.Tests/AclHttpClientTests.cs @@ -5,27 +5,42 @@ using EventSubscriber.Api; namespace EventSubscriber.Tests; /// -/// 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. /// 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( - () => 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( - () => 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(() => client.GetZaakReferenceAsync(null!)); + await Assert.ThrowsAsync(() => client.GetRegisterRecordAsync(null!)); Assert.Null(capture.Seen); } } diff --git a/services/event-subscriber/EventSubscriber.Tests/InMemoryStores.cs b/services/event-subscriber/EventSubscriber.Tests/InMemoryStores.cs index 5cbad6e..16b7af5 100644 --- a/services/event-subscriber/EventSubscriber.Tests/InMemoryStores.cs +++ b/services/event-subscriber/EventSubscriber.Tests/InMemoryStores.cs @@ -5,16 +5,18 @@ namespace EventSubscriber.Tests; /// 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). -/// 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). +/// 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. internal sealed class FakeAclClient : IAclClient { + public Dictionary Records { get; } = []; + public int CallCount { get; private set; } - public Task GetZaakReferenceAsync(Uri zaakUrl, CancellationToken ct = default) + public Task 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); } } diff --git a/services/event-subscriber/EventSubscriber.Tests/NotificationProjectorTests.cs b/services/event-subscriber/EventSubscriber.Tests/NotificationProjectorTests.cs index ed2a677..25aab00 100644 --- a/services/event-subscriber/EventSubscriber.Tests/NotificationProjectorTests.cs +++ b/services/event-subscriber/EventSubscriber.Tests/NotificationProjectorTests.cs @@ -2,13 +2,14 @@ using EventSubscriber.Application; namespace EventSubscriber.Tests; -/// 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). +/// 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). 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,56 @@ 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() + /// 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. + 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), 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); + } + + [Fact] + public async Task approval_updates_the_same_row_from_ingediend_to_ingeschreven() + { + var projector = Projector(); + await projector.HandleAsync(RecordWritten()); + // The ACL PATCHes the same object on approval, so Objecten publishes an `update`. + await projector.HandleAsync(RecordWritten("update", 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 +75,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 +84,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("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); } } diff --git a/services/projection-api/Projection.ReadModel/EfNotificationLog.cs b/services/projection-api/Projection.ReadModel/EfNotificationLog.cs index b46f510..2b63acb 100644 --- a/services/projection-api/Projection.ReadModel/EfNotificationLog.cs +++ b/services/projection-api/Projection.ReadModel/EfNotificationLog.cs @@ -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> 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); } diff --git a/services/projection-api/Projection.ReadModel/Migrations/20260828103132_ProjectionSourcedFromObjecten.Designer.cs b/services/projection-api/Projection.ReadModel/Migrations/20260828103132_ProjectionSourcedFromObjecten.Designer.cs new file mode 100644 index 0000000..2c26f55 --- /dev/null +++ b/services/projection-api/Projection.ReadModel/Migrations/20260828103132_ProjectionSourcedFromObjecten.Designer.cs @@ -0,0 +1,87 @@ +// +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 + { + /// + 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("Key") + .HasColumnType("text") + .HasColumnName("key"); + + b.Property("ReceivedAt") + .HasColumnType("timestamp with time zone") + .HasColumnName("received_at"); + + b.Property("Reference") + .HasColumnType("text") + .HasColumnName("reference"); + + b.Property("RegisterId") + .IsRequired() + .HasColumnType("text") + .HasColumnName("register_id"); + + b.Property("Status") + .IsRequired() + .HasColumnType("text") + .HasColumnName("status"); + + b.HasKey("Key"); + + b.ToTable("processed_notifications", (string)null); + }); + + modelBuilder.Entity("Projection.ReadModel.RegisterEntryRow", b => + { + b.Property("Id") + .HasColumnType("text") + .HasColumnName("id"); + + b.Property("Bsn") + .HasColumnType("text") + .HasColumnName("bsn"); + + b.Property("NaamPlaceholder") + .HasColumnType("text") + .HasColumnName("naam_placeholder"); + + b.Property("Reference") + .HasColumnType("text") + .HasColumnName("reference"); + + b.Property("Status") + .IsRequired() + .HasColumnType("text") + .HasColumnName("status"); + + b.HasKey("Id"); + + b.ToTable("register_projection", (string)null); + }); +#pragma warning restore 612, 618 + } + } +} diff --git a/services/projection-api/Projection.ReadModel/Migrations/20260828103132_ProjectionSourcedFromObjecten.cs b/services/projection-api/Projection.ReadModel/Migrations/20260828103132_ProjectionSourcedFromObjecten.cs new file mode 100644 index 0000000..81f3217 --- /dev/null +++ b/services/projection-api/Projection.ReadModel/Migrations/20260828103132_ProjectionSourcedFromObjecten.cs @@ -0,0 +1,69 @@ +using Microsoft.EntityFrameworkCore.Migrations; + +#nullable disable + +namespace Projection.ReadModel.Migrations +{ + /// + /// 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). + /// + /// + /// The old columns are dropped and the new ones added rather than renamed. EF scaffolded renames + /// (resourceregister_id, zaak_idstatus), 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. + /// + public partial class ProjectionSourcedFromObjecten : Migration + { + /// + 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( + name: "register_id", + table: "processed_notifications", + type: "text", + nullable: false, + defaultValue: ""); + + migrationBuilder.AddColumn( + name: "status", + table: "processed_notifications", + type: "text", + nullable: false, + defaultValue: ""); + } + + /// + 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( + name: "actie", table: "processed_notifications", type: "text", nullable: false, defaultValue: ""); + migrationBuilder.AddColumn( + name: "zaak_id", table: "processed_notifications", type: "text", nullable: false, defaultValue: ""); + migrationBuilder.AddColumn( + name: "resource", table: "processed_notifications", type: "text", nullable: false, defaultValue: ""); + } + } +} diff --git a/services/projection-api/Projection.ReadModel/Migrations/ProjectionDbContextModelSnapshot.cs b/services/projection-api/Projection.ReadModel/Migrations/ProjectionDbContextModelSnapshot.cs index 182add1..dc6b315 100644 --- a/services/projection-api/Projection.ReadModel/Migrations/ProjectionDbContextModelSnapshot.cs +++ b/services/projection-api/Projection.ReadModel/Migrations/ProjectionDbContextModelSnapshot.cs @@ -28,11 +28,6 @@ namespace Projection.ReadModel.Migrations .HasColumnType("text") .HasColumnName("key"); - b.Property("Actie") - .IsRequired() - .HasColumnType("text") - .HasColumnName("actie"); - b.Property("ReceivedAt") .HasColumnType("timestamp with time zone") .HasColumnName("received_at"); @@ -41,15 +36,15 @@ namespace Projection.ReadModel.Migrations .HasColumnType("text") .HasColumnName("reference"); - b.Property("Resource") + b.Property("RegisterId") .IsRequired() .HasColumnType("text") - .HasColumnName("resource"); + .HasColumnName("register_id"); - b.Property("ZaakId") + b.Property("Status") .IsRequired() .HasColumnType("text") - .HasColumnName("zaak_id"); + .HasColumnName("status"); b.HasKey("Key"); diff --git a/services/projection-api/Projection.ReadModel/ProjectionDbContext.cs b/services/projection-api/Projection.ReadModel/ProjectionDbContext.cs index d9df709..d363139 100644 --- a/services/projection-api/Projection.ReadModel/ProjectionDbContext.cs +++ b/services/projection-api/Projection.ReadModel/ProjectionDbContext.cs @@ -34,9 +34,8 @@ public sealed class ProjectionDbContext(DbContextOptions 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; } } -/// An accepted notification, retained so the projection can be rebuilt without OpenZaak (§8.1). +/// 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). public sealed class ProcessedNotificationRow { public required string Key { get; set; } - public required string Actie { get; set; } - public required string ZaakId { get; set; } - /// The ZGW resource (e.g. zaak or status) — retained so a rebuild reprojects - /// the right status without reading OpenZaak (S-09b). - public required string Resource { get; set; } + /// The registration this record is for (the zaak id) — the projection row's key. + public required string RegisterId { get; set; } - /// The zaak reference (identificatie), retained so a rebuild reprojects it without the ACL (#78). + /// The register status the record carried (INGEDIEND / INGESCHREVEN). + public required string Status { get; set; } + + /// The citizen-facing reference the record carried — matches the submit confirmation (#78). public string? Reference { get; set; } public DateTimeOffset ReceivedAt { get; set; }