using EventSubscriber.Application; namespace EventSubscriber.Tests; /// 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 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(); private readonly FakeAclClient _acl = new(); private NotificationProjector Projector() => new(_log, _store, _acl); /// 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) { _acl.Records[url] = new RegisterRecord(zaakId, status, "REG-2026-0001"); return new Notification("objecten", "object", actie, new Uri(url)); } [Fact] public async Task a_register_record_write_is_projected_as_a_row_keyed_on_the_registration() { await Projector().HandleAsync(RecordWritten()); var entry = Assert.Single(await _store.AllAsync()); // Keyed on the record's own id (the zaak id), not on the Objecten object's uuid — the // projection row and the register record are the same registration. Assert.Equal(ZaakId, entry.Id); Assert.Equal(RegistrationStatus.Ingediend, entry.Status); Assert.Equal("REG-2026-0001", entry.Reference); } [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(RecordWritten()); await projector.HandleAsync(RecordWritten()); Assert.Single(await _store.AllAsync()); } [Fact] public async Task a_replayed_notification_never_reaches_the_projection_store() { var projector = Projector(); 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. Assert.Equal(1, _store.UpsertCount); } [Fact] public async Task two_different_registrations_each_get_their_own_row() { var projector = Projector(); await projector.HandleAsync(RecordWritten()); await projector.HandleAsync(RecordWritten(url: ObjectUrl[..^1] + "2", zaakId: "other-zaak")); Assert.Equal(2, (await _store.AllAsync()).Count); } [Theory] [InlineData("zaken", "zaak", "create")] // the ZGW source S-19b-2 replaced [InlineData("zaken", "status", "create")] // ditto [InlineData("objecten", "object", "destroy")] // a delete we do not project [InlineData("documenten", "object", "create")] // wrong kanaal public async Task an_unrelated_notification_is_not_projected(string kanaal, string resource, string actie) { _acl.Records[ObjectUrl] = new RegisterRecord(ZaakId, RegistrationStatus.Ingediend, "REG-2026-0001"); await Projector().HandleAsync(new Notification(kanaal, resource, actie, new Uri(ObjectUrl))); Assert.Empty(await _store.AllAsync()); } [Fact] public async Task rebuild_reproduces_the_row_without_re_reading_through_the_acl() { var projector = Projector(); 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(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(ZaakId, entry.Id); Assert.Equal(RegistrationStatus.Ingediend, entry.Status); } }