From af3cb83c7eeb0759eb58530ed6d30f5e4e58a359 Mon Sep 17 00:00:00 2001 From: Renatillas Date: Mon, 17 Aug 2026 11:39:58 +0200 Subject: [PATCH] v2 maybe --- README.md | 5 + backends/factos_pog/README.md | 94 +- backends/factos_pog/dev/factos_pog_dev.gleam | 54 +- backends/factos_pog/docs/durable-effects.md | 69 +- .../factos_pog/docs/durable-subscriptions.md | 136 +- backends/factos_pog/docs/how-it-works.md | 87 +- backends/factos_pog/gleam.toml | 3 + backends/factos_pog/manifest.toml | 4 + .../dbmate/20260816000100_factos_pog_v2.sql | 6 +- backends/factos_pog/priv/migrations.sql | 6 +- .../factos_pog/src/factos/factos_pog.gleam | 1020 ++++++++-- .../factos_pog/test/factos_pog_test.gleam | 1767 +++++++++++++++-- 12 files changed, 2875 insertions(+), 376 deletions(-) diff --git a/README.md b/README.md index 1431d9d..d192c61 100644 --- a/README.md +++ b/README.md @@ -218,6 +218,11 @@ itself. Factos does not maintain a projection table automatically. Views can always be recomputed if the original events are still decodable. That is why event codec compatibility matters. +The pure `factos` package keeps this storage-independent contract. The +PostgreSQL backend `factos_pog` can additionally maintain an application-owned +projection and its named cursor in one transaction, then make a dispatch wait +for that cursor before the caller reads the projection. + ## How are effects handled? Reactors turn committed event records into effect values: diff --git a/backends/factos_pog/README.md b/backends/factos_pog/README.md index b78354c..f585569 100644 --- a/backends/factos_pog/README.md +++ b/backends/factos_pog/README.md @@ -60,8 +60,8 @@ cursor-based reads recover anything not delivered by the notification session. The dbmate history retains the `1.0.0` event-store baseline and one `2.0.0` upgrade. `20260816000100_factos_pog_v2.sql` converts stored JSON to JSONB, -enforces UUIDv4 event identity, and installs durable subscriptions and event -notifications. +enforces UUIDv4 event identity, and installs generation-checked durable +subscriptions with complete recorded event notification envelopes. ## Define a codec @@ -163,6 +163,30 @@ let assert Ok(dispatch) = |> factos_pog.dispatch(BuyTicket(buyer: "renata"), event_id: uuid.v4_string) ``` +`with_query` controls the facts that must remain stable while the command +transaction commits. It does not wait for a read model. When the caller will +immediately query named PostgreSQL projections, use the post-commit barrier: + +```gleam +let assert Ok(dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "ticket-sale-renata", + decider: ticket_decider(), + codec: ticket_codec(), + ) + |> factos_pog.dispatch_and_wait( + BuyTicket(buyer: "renata"), + event_id: uuid.v4_string, + subscriptions: ["ticket-projection"], + timeout: duration.seconds(5), + ) +``` + +Name only projections the caller intends to query. A post-commit timeout or +missing name returns `DispatchCommittedButNotObserved` with the committed +`Dispatch`; it does not roll the event back. + ## Consume committed events durably A managed subscription owns its PostgreSQL `LISTEN` session and a private @@ -175,16 +199,41 @@ let assert Ok(subscription) = config:, name: "ticket-projection", query: factos.AllEvents, + start_from: factos_pog.Origin, codec: ticket_codec(), - handle: fn(decoded) { - ticket_projection.apply(decoded) + handle: fn(recorded) { + ticket_projection.apply(recorded) }, ) ``` -The backend stores the cursor in `factos_subscriptions`. It processes events -sequentially and checkpoints each event after the handler returns `Ok(Nil)`. -Application code does not load or update the cursor. +The handler receives a complete `factos.Recorded` envelope, including event id, +stream, revision, global position, descriptor, and domain event. The backend +stores the cursor in `factos_subscriptions`, processes events sequentially, and +checkpoints each event after the handler returns `Ok(Nil)`. + +For an application-owned PostgreSQL projection, commit the projection row and +cursor atomically: + +```gleam +let assert Ok(subscription) = + factos_pog.new_projection_subscription( + connection:, + config:, + name: "ticket-projection", + query: factos.AllEvents, + start_from: factos_pog.Origin, + codec: ticket_codec(), + project: fn(transaction_connection, recorded) { + ticket_projection.apply(transaction_connection, recorded) + }, + ) +``` + +`Origin`, `Current`, and `After(position:)` initialize only a new stable name. +Use `subscription_status` to inspect its cursor, generation, global log head, +and exact global-log lag. Use `reset_subscription` to reposition an existing +name and increment its generation; an online worker reloads immediately. For development or embedding outside an application supervisor: @@ -217,20 +266,41 @@ path. A processing failure restarts the subscriber from its last successful event. A notification-session failure restarts the session and subscriber in dependency order. -Delivery is at least once. A crash after an external effect succeeds but before -the automatic cursor update can replay the event. External handlers therefore -need a deterministic idempotency key derived from stable domain data or event -metadata. Retry schedules, dead letters, and poison-event policy remain -application concerns rather than a second backend-owned effect log. +Ordinary handlers are at least once. A crash after an external effect succeeds +but before the separate cursor update can replay the event. Use `Recorded.id` +when one source event maps to one operation, or a stable domain-operation key +when one operation spans events. Atomic PostgreSQL projection callbacks instead +commit their changes and cursor together. Retry schedules, dead letters, and +poison-event policy remain application concerns. Only one running worker should own a subscription name. Use leader election or a session advisory lock when multiple application instances may start the same consumer. +A handler that dispatches a follow-up command must not wait for its own +subscription name: that cursor cannot advance until the handler returns, so the +follow-up commits and then times out. + `listen`, `unlisten`, and `read_after` remain available as low-level primitives for applications that need a custom scheduler, batch model, or checkpoint protocol. +## Run tests + +The integration suite shares one externally managed PostgreSQL service and +creates an isolated database for each test. Start and stop the checked-in +Compose service around the suite: + +```sh +docker compose up --wait -d +gleam test +docker compose down -v +``` + +The defaults match `compose.yml`. Override a remote or CI service with +`FACTOS_POG_TEST_HOST`, `FACTOS_POG_TEST_PORT`, `FACTOS_POG_TEST_DATABASE`, +`FACTOS_POG_TEST_USERNAME`, and `FACTOS_POG_TEST_PASSWORD`. + ## Tradeoff: serializable contention PostgreSQL `SERIALIZABLE` isolation protects arbitrary event-type/tag predicates diff --git a/backends/factos_pog/dev/factos_pog_dev.gleam b/backends/factos_pog/dev/factos_pog_dev.gleam index 480c6f8..49d05cd 100644 --- a/backends/factos_pog/dev/factos_pog_dev.gleam +++ b/backends/factos_pog/dev/factos_pog_dev.gleam @@ -11,6 +11,7 @@ import gleam/option.{Some} import gleam/otp/actor import gleam/result import gleam/string +import gleam/time/duration import gleamy/bench import pog import simplifile @@ -121,6 +122,9 @@ fn start_connection( } fn reset_schema(connection: pog.Connection) -> Nil { + let assert Ok(_) = + pog.query("drop table if exists factos_dev_projection") + |> pog.execute(on: connection) let assert Ok(_) = pog.query("drop table if exists factos_outbox") |> pog.execute(on: connection) @@ -299,39 +303,71 @@ fn dispatch_once( fn smoke_subscription(config: pog.Config, connection: pog.Connection) -> Nil { reset_schema(connection) - let processed = process.new_subject() + let assert Ok(_) = + pog.query( + "create table factos_dev_projection ( + singleton boolean primary key default true check (singleton), + value integer not null + )", + ) + |> pog.execute(on: connection) let assert Ok(subscription) = - factos_pog.new_subscription( + factos_pog.new_projection_subscription( connection:, config:, name: "subscription-smoke", query: factos.AllEvents, + start_from: factos_pog.Origin, codec: codec(), - handle: fn(event) { - process.send(processed, event) - Ok(Nil) + project: fn(transaction_connection, recorded) { + let factos.Recorded(event: Incremented(value:), ..) = recorded + pog.query( + "insert into factos_dev_projection (singleton, value) + values (true, $1) + on conflict (singleton) do update set value = excluded.value", + ) + |> pog.parameter(pog.int(value)) + |> pog.execute(on: transaction_connection) + |> result.map(fn(_) { Nil }) + |> result.map_error(string.inspect) }, ) let assert Ok(actor.Started(data: handle, ..)) = factos_pog.start(subscription) - let assert Ok(_) = + let assert Ok(dispatch) = factos_pog.new_dispatch( connection:, stream: "subscription-smoke", decider: decider(), codec: codec(), ) - |> factos_pog.dispatch(Increment, event_id: uuid.v4_string) - let assert Ok(factos.Decoded(event: Incremented(value:), ..)) = - process.receive(processed, within: 10_000) + |> factos_pog.dispatch_and_wait( + Increment, + event_id: uuid.v4_string, + subscriptions: ["subscription-smoke"], + timeout: duration.seconds(5), + ) + let assert Ok(returned) = + pog.query("select value from factos_dev_projection where singleton = true") + |> pog.returning(int_column_decoder()) + |> pog.execute(on: connection) + let assert [value] = returned.rows assert value == 1 + let assert Ok(factos_pog.SubscriptionStatus(cursor:, events_behind: 0, ..)) = + factos_pog.subscription_status(connection, name: "subscription-smoke") + assert cursor == dispatch.append.position let assert Ok(Nil) = factos_pog.stop(handle) Nil } +fn int_column_decoder() -> decode.Decoder(Int) { + use value <- decode.field(0, decode.int) + decode.success(value) +} + fn decider() -> factos.Decider(Command, State, Event, Nil) { factos.decider(initial: Counter(0), decide:, evolve:) } diff --git a/backends/factos_pog/docs/durable-effects.md b/backends/factos_pog/docs/durable-effects.md index d28251e..e33ea6a 100644 --- a/backends/factos_pog/docs/durable-effects.md +++ b/backends/factos_pog/docs/durable-effects.md @@ -1,8 +1,8 @@ # Durable Effects -`factos_pog` does not maintain an effect outbox. Managed subscribers consume the -authoritative event log from a backend-owned cursor, derive effects from decoded -events, and execute them after the event transaction commits. +`factos_pog` does not maintain an effect outbox. Ordinary managed subscribers +consume the authoritative event log from a backend-owned cursor and execute +effects after the event transaction commits. Read [How `factos_pog` works](how-it-works.html) for the storage model and [Durable Subscriptions](durable-subscriptions.html) for the catch-up loop. @@ -26,9 +26,9 @@ fn ledger_effects(event: Event) -> List(Effect) { } ``` -The managed handler receives the domain event and its descriptor. It can derive -the effect values first, then perform the required IO. Put any stable operation -identity needed for delivery in domain data or event metadata. +The managed handler receives a complete `factos.Recorded` envelope. It can +derive effect values first, then perform the required IO using stable identity +from the record or domain data. ## Delivery is at least once @@ -51,6 +51,11 @@ Send that key to a destination that supports idempotency, or record it in a uniquely constrained application table before applying a local database effect. Do not allocate a fresh identifier on each retry. +When one external operation corresponds to one source event, `Recorded.id` is a +stable retry key. Keep a domain-operation key when one logical operation spans +multiple events; using each source id would incorrectly execute it more than +once. + ## Complete one event before returning The backend advances the subscription cursor only after the handler returns @@ -87,20 +92,54 @@ subscription name, stable domain operation identity, effect identity, failure, attempt history, and replay audit. That table represents application policy rather than a second generic event store. -## Database projections +## Atomic PostgreSQL projections + +Use `new_projection_subscription` when the derived state and subscription +checkpoint live in the same PostgreSQL database: + +```gleam +let assert Ok(subscription) = + factos_pog.new_projection_subscription( + connection:, + config:, + name: "user-projection", + query: factos.AllEvents, + start_from: factos_pog.Origin, + codec: event_codec(), + project: fn(transaction_connection, recorded) { + let factos.Recorded( + id: event_id, + event: UserRegistered(username:), + .., + ) = recorded + pog.query( + "insert into user_projection (event_id, username) values ($1, $2)", + ) + |> pog.parameter(pog.text(event_id)) + |> pog.parameter(pog.text(username)) + |> pog.execute(on: transaction_connection) + |> result.map(fn(_) { Nil }) + |> result.map_error(string.inspect) + }, + ) +``` -The managed cursor update occurs after the handler, not inside the handler's -projection transaction. Projection changes must therefore be replay-safe, for -example through deterministic upserts or a unique logical operation key. +The backend opens one transaction per matching event, passes its +transaction-scoped connection to `project`, and advances the generation-checked +cursor on that same connection only after `project` returns `Ok(Nil)`. Callback +errors, query errors, crashes, and concurrent resets roll back both changes. +Do not start nested Pog transactions or call external services from `project`. -If a projection requires its changes and checkpoint to commit in one PostgreSQL -transaction, use `read_after` with an application-owned checkpoint -instead of the managed subscription. +Use `Recorded.id` as a source-event uniqueness key when each event contributes +one row. Keep domain-operation keys for projections that combine several events. +Ordinary `new_subscription` handlers remain at least once because their work +finishes before a separate cursor update. ## Rebuilds and operational replay -Projection rebuilds can clear their derived state and consume the event log from -`factos.NoPosition` through a custom runtime or a separate subscription name. +Projection rebuilds can clear derived state and reset the existing subscription +to `Origin`, or consume under a separate stable name when both versions must run +at once. External effects need an explicit replay decision. Replaying historical events may resend emails, charge accounts, or publish messages. Require an operator to diff --git a/backends/factos_pog/docs/durable-subscriptions.md b/backends/factos_pog/docs/durable-subscriptions.md index 1545f09..03e8f62 100644 --- a/backends/factos_pog/docs/durable-subscriptions.md +++ b/backends/factos_pog/docs/durable-subscriptions.md @@ -19,9 +19,10 @@ let assert Ok(subscription) = config:, name: "welcome-email", query: factos.AllEvents, + start_from: factos_pog.Origin, codec: event_codec(), - handle: fn(decoded) { - let factos.Decoded(event:, descriptor:) = decoded + handle: fn(recorded) { + let factos.Recorded(event:, descriptor:, ..) = recorded case event { UserRegistered(user_id:, email_address:) -> welcome_email.send( @@ -43,10 +44,92 @@ stable across process and application restarts. The backend stores its cursor in allocates a private notification-session name, so the application does not need to create a second named configuration. -The handler receives a decoded domain event and its `factos.EventDescriptor`. +The handler receives a complete `factos.Recorded` envelope: stable event id, +stream, revision, global position, domain event, and `factos.EventDescriptor`. It returns `Ok(Nil)` only after all required work for that event has completed. Events are handled sequentially in global commit order. +## Choose the initial cursor + +`start_from` applies only when the stable name has no durable row: + +- `Origin` starts before the first event; +- `Current` stores the event-log head observed after both listeners are active; +- `After(position:)` starts after `NoPosition` or a non-negative position that + is not ahead of the log. + +An existing row always wins over the constructor value. Use +`reset_subscription` to reposition it. Each reset changes the cursor and +increments its generation in one transaction, then wakes a live worker. A +stale worker generation cannot overwrite the reset. + +`subscription_status` returns the cursor, generation, current event-log head, +and `events_behind`. That count is global-log lag, not matching-query lag, +because the cursor advances across non-matching rows. + +## Choose the processing transaction + +`new_subscription` is for ordinary at-least-once handlers, including external +services. The handler completes before a separate cursor write. + +`new_projection_subscription` atomically commits an application-owned +PostgreSQL projection and its cursor: + +```gleam +let assert Ok(subscription) = + factos_pog.new_projection_subscription( + connection:, + config:, + name: "user-projection", + query: factos.AllEvents, + start_from: factos_pog.Origin, + codec: event_codec(), + project: fn(transaction_connection, recorded) { + user_projection.apply(transaction_connection, recorded) + }, + ) +``` + +The projection callback receives a transaction-scoped Pog connection. Its +changes and the generation-checked cursor roll back together on callback error, +query error, crash, or reset race. Process one event per transaction; do not +nestedly call `pog.transaction` or run external services in the callback. + +## Wait for selected projections after dispatch + +Factos command-context consistency protects the facts read by a command before +its event commits. `dispatch_and_wait` is a separate post-commit guarantee: it +waits until each selected durable cursor reaches the dispatch position. + +```gleam +let assert Ok(dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "user-renata", + decider: user_decider(), + codec: event_codec(), + ) + |> factos_pog.dispatch_and_wait( + RegisterUser(user_id: "renata"), + event_id: uuid.v4_string, + subscriptions: ["user-projection"], + timeout: duration.seconds(5), + ) +``` + +Successful return means the selected atomic PostgreSQL projections and their +cursors are query-visible. Name only projections the caller will immediately +read. `wait_for_subscriptions` provides the same barrier for an already known +position. + +A commit failure is `DispatchNotCommitted`. Every wait failure is +`DispatchCommittedButNotObserved`, which carries the committed dispatch so the +caller can recover without retrying the command as if it had failed. + +Do not wait for a subscription from a command dispatched inside that same +subscription's callback. The cursor cannot advance until the callback returns, +so the follow-up command commits and its self-wait times out. + ## Start or supervise it For development or embedding outside an application supervisor: @@ -77,22 +160,23 @@ failure restarts both children in dependency order. ## Delivery path -Initialization establishes `LISTEN` before loading the cursor. This closes the -startup race: an event committed while catch-up runs queues a notification for -the already-established listener. +Initialization establishes both event and reset `LISTEN` registrations before +loading the checkpoint. This closes the startup race: a commit or reset queued +while catch-up runs is already observable. For an envelope smaller than PostgreSQL's notification limit, the trigger sends: -- the new event cursor; -- the immediately preceding committed cursor; +- the new event cursor and immediately preceding committed cursor; +- event id, stream, and stream revision; - event type, version, tags, and metadata; - event JSON data. When `previous_cursor` equals the subscriber's current cursor, the notification -is contiguous. The worker applies the query to its type and tags, decodes a -matching event, calls the handler directly, and stores the new cursor after -`Ok(Nil)`. A non-matching contiguous event advances the cursor without invoking -the decoder or handler. +is contiguous. The worker applies the query to its type and tags and decodes a +matching event. Ordinary handlers then checkpoint separately; PostgreSQL +projection callbacks and their generation-checked cursor commit in one +transaction. A non-matching contiguous event advances the cursor without +invoking the decoder or callback. A notification containing a large event is reduced to its cursor pair so the append transaction cannot exceed PostgreSQL's payload limit. The worker then @@ -118,25 +202,27 @@ Duplicate, delayed, and spurious notifications are harmless: cursors at or behind the stored position are ignored. Notifications are visible only after the appending transaction commits. -## Delivery guarantee +## Delivery guarantees -The contract is at least once. The backend checkpoints each event automatically, -but PostgreSQL cannot atomically commit that cursor with an external HTTP -request, message publish, or email send: +Ordinary handlers are at least once. PostgreSQL cannot atomically commit their +cursor with an external HTTP request, message publish, or email send: 1. the handler completes the external effect; 2. the process fails before the cursor update commits; -3. the same event is delivered again after restart. +3. the same recorded event is delivered again after restart. + +Use `Recorded.id` as an idempotency key when one source event maps to one +operation. Use stable domain-operation identity when one operation spans +multiple events. A fresh key on each retry defeats idempotency. -Use a destination-supported idempotency key derived from stable domain data or -event metadata. The managed handler intentionally does not expose the storage -position or event id as application identity. +PostgreSQL projection handlers use the atomic constructor described above, so +their database writes and cursor commit or roll back together. -A handler error stops the subscriber actor abnormally. Supervision restarts it -from the last successful event. There is no backend-wide retry schedule or -dead-letter policy; those decisions depend on the integration. A permanently -failing event blocks later matching events until the application handles or -operationally resolves it. +A callback error stops the subscriber actor abnormally. Supervision restarts it +from durable state. There is no backend-wide retry schedule or dead-letter +policy; those decisions depend on the integration. A permanently failing event +blocks later matching events until the application handles or operationally +resolves it. ## Multiple application instances diff --git a/backends/factos_pog/docs/how-it-works.md b/backends/factos_pog/docs/how-it-works.md index 7bdb22c..f69f92c 100644 --- a/backends/factos_pog/docs/how-it-works.md +++ b/backends/factos_pog/docs/how-it-works.md @@ -1,10 +1,10 @@ # How `factos_pog` Works `factos_pog` is the PostgreSQL event-store backend for Factos. It owns event -persistence, command-context consistency, managed subscription cursors, ordered -recovery reads, notification sessions, and local subscription supervision. -Applications own codecs, effect execution, idempotency, retry policy, and -application-level supervision. +persistence, command-context consistency, managed subscription cursors, atomic +PostgreSQL projection checkpoints, ordered recovery reads, notification +sessions, and local subscription supervision. Applications own codecs, +external effects, idempotency, retry policy, and application-level supervision. ## Dispatch flow @@ -55,11 +55,25 @@ transaction aborts instead of allowing both commands to accept stale context. For stream dispatch, the current stream revision is the append boundary. +Command-context consistency ends when the event transaction commits. It does +not imply that a derived projection has caught up. + +`dispatch_and_wait` adds that separate read-after-write guarantee. After a +successful dispatch it polls all selected stable subscription names in one +cursor query and returns only when every cursor reaches the dispatch position. +A commit failure is distinct from a post-commit wait failure; the latter carries +the committed `Dispatch`. + +Empty name lists and no-event dispatches return without reading subscription +state. Polling removes duplicate names while preserving first-occurrence order, +uses a bounded 20-millisecond interval, and reports only lagging +`SubscriptionStatus` values on timeout. + ## Managed durable subscriptions -`new_subscription` binds a stable subscription name to an event query, codec, -and handler. The backend creates the name's cursor row in -`factos_subscriptions`; the application never has to load or store that cursor. +`new_subscription` and `new_projection_subscription` bind a stable name to an +event query, codec, initial cursor, and callback. The backend owns the name's +cursor and generation in `factos_subscriptions`. `start` starts a local supervision tree and returns a stoppable handle. `stop` waits for that complete tree to terminate. `supervised` returns the same tree as @@ -70,14 +84,15 @@ children. The managed lifecycle is: -1. establish `LISTEN` on a private Pog notification connection; -2. load the backend-owned cursor; -3. read and process any existing events in global position order; -4. handle contiguous committed events directly from notification payloads; -5. checkpoint each event after its handler returns `Ok(Nil)`; +1. establish `LISTEN` for events and subscription resets; +2. load the durable cursor and generation, or create them from `Origin`, + `Current`, or `After(position:)` on the first start; +3. read and process existing events in global position order; +4. handle contiguous committed events directly from full recorded envelopes; +5. generation-check each cursor advance; 6. use ordered reads after a notification gap or cursor-only notification; -7. every 30 seconds, reconcile work missed during a reconnect; -8. after a failure, reload the last durable cursor before retrying. +7. every 30 seconds, reload durable state and reconcile missed work; +8. after a failure or reset, reload the authoritative checkpoint before retrying. Listening before loading the cursor closes the startup race: a commit after `LISTEN` queues a notification while catch-up runs. @@ -104,8 +119,15 @@ subscription runtimes. ## Checkpoints and concurrency -The managed worker processes matching events sequentially and stores its cursor -after every successful handler call. It never advances past failed work. +The managed worker processes matching events sequentially and never advances +past failed work. Ordinary handlers complete before a separate cursor update. +Projection handlers update application rows and the cursor in one PostgreSQL +transaction. + +Every cursor write includes the generation loaded by the worker. A reset changes +the cursor and increments generation transactionally, so an in-flight stale +worker cannot overwrite it. Reset notifications are only a wake-up path; +periodic durable reload recovers a missed notification. One running worker should own each subscription name. The cursor row is durable state, not a leader-election lock; multiple application instances must use @@ -117,21 +139,26 @@ gap-tracking model and therefore belongs in a custom runtime. ## Effects and delivery guarantees -Derive effects deterministically from the decoded event and descriptor, then -execute them after the event transaction has committed. +Every callback receives the full `factos.Recorded` envelope, including id, +stream, revision, global position, descriptor, and domain event. -The managed cursor update occurs after the handler returns, so all managed -delivery is at least once, including database projections. Projection writes -must be replay-safe, for example through deterministic upserts or unique -operation keys. Applications that require a projection update and cursor in one -PostgreSQL transaction can build that transaction with `read_after` and -an application-owned checkpoint. +`new_subscription` is at least once. For an external service, no atomic +transaction spans PostgreSQL and that service. A crash after external success +but before cursor advancement causes redelivery. Use `Recorded.id` when one +source event maps to one operation, and domain-operation identity when one +operation spans several events. -For an external service, no atomic transaction spans PostgreSQL and that -service. A crash after external success but before cursor advancement causes -redelivery, so the handler must use a stable idempotency key. +`new_projection_subscription` opens one Pog transaction per matching event, +passes its scoped connection to the callback, and writes the generation-checked +cursor only after the callback succeeds. Callback errors, query errors, crashes, +and reset races roll back both the application projection and checkpoint. Retries, backoff, dead letters, poison-event handling, and operational replay are -application policy. Keeping those decisions outside `factos_pog` avoids a second -durable effect log and allows each integration to choose the policy its external -system actually supports. +application policy. Keeping those decisions outside `factos_pog` avoids a +second durable effect log and allows each integration to choose the policy its +external system actually supports. + +Callers should wait only for projections they will immediately query. A +subscription callback must not dispatch a command that waits for the same +subscription name: its cursor cannot advance until the callback returns, so the +nested wait can only time out after committing. diff --git a/backends/factos_pog/gleam.toml b/backends/factos_pog/gleam.toml index d6de49c..6de4a29 100644 --- a/backends/factos_pog/gleam.toml +++ b/backends/factos_pog/gleam.toml @@ -34,11 +34,13 @@ source = "./docs/durable-effects.md" factos = { path = "../.." } gleam_stdlib = ">= 1.0.0 and < 2.0.0" pog = { git = "https://github.com/foxfriends/pog.git", ref = "919fd6ac96095ea11fa7c940b17eaece49cc5993" } +gleam_time = ">= 1.8.0 and < 2.0.0" gleam_erlang = ">= 1.0.0 and < 2.0.0" gleam_json = ">= 3.1.0 and < 4.0.0" gleam_otp = ">= 1.2.0 and < 2.0.0" [dev_dependencies] +envoy = ">= 1.0.0 and < 2.0.0" unitest = ">= 1.0.0 and < 2.0.0" gleeunit = ">= 1.0.0 and < 2.0.0" testcontainer = ">= 1.0.2 and < 2.0.0" @@ -46,3 +48,4 @@ testcontainer_formulas = ">= 1.0.0 and < 2.0.0" simplifile = ">= 2.5.0 and < 3.0.0" youid = ">= 1.5.4 and < 2.0.0" gleamy_bench = ">= 0.6.0 and < 1.0.0" +global_value = ">= 1.0.0 and < 2.0.0" diff --git a/backends/factos_pog/manifest.toml b/backends/factos_pog/manifest.toml index fdb9342..b2700f1 100644 --- a/backends/factos_pog/manifest.toml +++ b/backends/factos_pog/manifest.toml @@ -31,6 +31,7 @@ packages = [ { name = "glearray", version = "2.1.2", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "glearray", source = "hex", outer_checksum = "1554E48DD40114D7602F5BFF4D7278B6B3B735F137C7FDEEADFB2FE7951C94BE" }, { name = "gleeunit", version = "1.11.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleeunit", source = "hex", outer_checksum = "EC31ABA74256AEA531EDF8169931D775BBB384FED0A8A1BDC4DD9354E3E21826" }, { name = "glexer", version = "2.5.0", build_tools = ["gleam"], requirements = ["gleam_stdlib", "splitter"], otp_app = "glexer", source = "hex", outer_checksum = "D423CF63B8D4F9654EF53D0AF4904B249EB0679BD41C63440170AC508C9CC65E" }, + { name = "global_value", version = "1.0.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "global_value", source = "hex", outer_checksum = "23F74C91A7B819C43ABCCBF49DAD5BB8799D81F2A3736BA9A534BD47F309FF4F" }, { name = "opentelemetry_api", version = "1.5.0", build_tools = ["rebar3", "mix"], requirements = [], otp_app = "opentelemetry_api", source = "hex", outer_checksum = "F53EC8A1337AE4A487D43AC89DA4BD3A3C99DDF576655D071DEED8B56A2D5DDA" }, { name = "pg_types", version = "0.6.0", build_tools = ["rebar3"], requirements = [], otp_app = "pg_types", source = "hex", outer_checksum = "9949A4849DD13408FA249AB7B745E0D2DFDB9532AEE2B9722326E33CD082A778" }, { name = "pgo", version = "0.20.0", build_tools = ["rebar3"], requirements = ["backoff", "opentelemetry_api", "pg_types"], otp_app = "pgo", source = "hex", outer_checksum = "2F11E6649CEB38E569EF56B16BE1D04874AE5B11A02867080A2817CE423C683B" }, @@ -49,13 +50,16 @@ packages = [ ] [requirements] +envoy = { version = ">= 1.0.0 and < 2.0.0" } factos = { path = "../.." } gleam_erlang = { version = ">= 1.0.0 and < 2.0.0" } gleam_json = { version = ">= 3.1.0 and < 4.0.0" } gleam_otp = { version = ">= 1.2.0 and < 2.0.0" } gleam_stdlib = { version = ">= 1.0.0 and < 2.0.0" } +gleam_time = { version = ">= 1.8.0 and < 2.0.0" } gleamy_bench = { version = ">= 0.6.0 and < 1.0.0" } gleeunit = { version = ">= 1.0.0 and < 2.0.0" } +global_value = { version = ">= 1.0.0 and < 2.0.0" } pog = { git = "https://github.com/foxfriends/pog.git", ref = "919fd6ac96095ea11fa7c940b17eaece49cc5993" } simplifile = { version = ">= 2.5.0 and < 3.0.0" } testcontainer = { version = ">= 1.0.2 and < 2.0.0" } diff --git a/backends/factos_pog/priv/dbmate/20260816000100_factos_pog_v2.sql b/backends/factos_pog/priv/dbmate/20260816000100_factos_pog_v2.sql index 4efe7ff..4f956f6 100644 --- a/backends/factos_pog/priv/dbmate/20260816000100_factos_pog_v2.sql +++ b/backends/factos_pog/priv/dbmate/20260816000100_factos_pog_v2.sql @@ -9,7 +9,8 @@ alter table factos_events create table factos_subscriptions ( name text primary key check (btrim(name) <> ''), - cursor bigint not null default -1 check (cursor >= -1) + cursor bigint not null default -1 check (cursor >= -1), + generation bigint not null default 0 check (generation >= 0) ); -- Sequence values are allocated before commit. Taking this transaction-scoped @@ -46,6 +47,9 @@ begin 'cursor', new.position, 'previous_cursor', previous_cursor, 'event', jsonb_build_object( + 'id', new.id, + 'stream', new.stream, + 'revision', new.revision, 'type', new.type, 'version', new.version, 'tags', case diff --git a/backends/factos_pog/priv/migrations.sql b/backends/factos_pog/priv/migrations.sql index 4d7ca1f..a337c53 100644 --- a/backends/factos_pog/priv/migrations.sql +++ b/backends/factos_pog/priv/migrations.sql @@ -40,7 +40,8 @@ on conflict do nothing; create table if not exists factos_subscriptions ( name text primary key check (btrim(name) <> ''), - cursor bigint not null default -1 check (cursor >= -1) + cursor bigint not null default -1 check (cursor >= -1), + generation bigint not null default 0 check (generation >= 0) ); create or replace function factos_pog_lock_event_append() @@ -76,6 +77,9 @@ begin 'cursor', new.position, 'previous_cursor', previous_cursor, 'event', jsonb_build_object( + 'id', new.id, + 'stream', new.stream, + 'revision', new.revision, 'type', new.type, 'version', new.version, 'tags', case diff --git a/backends/factos_pog/src/factos/factos_pog.gleam b/backends/factos_pog/src/factos/factos_pog.gleam index 918b87b..15135ed 100644 --- a/backends/factos_pog/src/factos/factos_pog.gleam +++ b/backends/factos_pog/src/factos/factos_pog.gleam @@ -29,6 +29,8 @@ import gleam/otp/static_supervisor import gleam/otp/supervision import gleam/result import gleam/string +import gleam/time/duration +import gleam/time/timestamp import pog /// A domain event prepared for PostgreSQL persistence. @@ -120,6 +122,15 @@ pub type Dispatch(event) { Dispatch(append: Append, events: List(factos.Recorded(event))) } +/// A dispatch failure before commit or a post-commit observation failure. +pub type DispatchWaitError(event, domain_error) { + DispatchNotCommitted(error: Error(domain_error)) + DispatchCommittedButNotObserved( + dispatch: Dispatch(event), + reason: SubscriptionError, + ) +} + pub opaque type DispatchBuilder(command, state, event, domain_error) { DispatchBuilder( connection: pog.Connection, @@ -131,18 +142,64 @@ pub opaque type DispatchBuilder(command, state, event, domain_error) { ) } +/// Initial durable cursor used only when a subscription name is first created. +pub type SubscriptionStart { + Origin + Current + After(position: factos.SequencePosition) +} + +/// Durable progress and global event-log lag for a managed subscription. +pub type SubscriptionStatus { + SubscriptionStatus( + name: String, + cursor: factos.SequencePosition, + generation: Int, + event_log_position: factos.SequencePosition, + events_behind: Int, + ) +} + +/// Failure while configuring or coordinating a managed subscription. +pub type SubscriptionError { + InvalidSubscriptionName(name: String) + InvalidSubscriptionPosition(position: factos.SequencePosition) + SubscriptionNotFound(name: String) + SubscriptionStartAfterEventLog( + name: String, + requested: factos.SequencePosition, + event_log_position: factos.SequencePosition, + ) + SubscriptionStoreError(error: pog.QueryError) + SubscriptionWaitTimedOut( + through: factos.SequencePosition, + pending: List(SubscriptionStatus), + ) +} + +type SubscriptionHandler(event, processing_error) { + OrdinarySubscriptionHandler( + handle: fn(factos.Recorded(event)) -> Result(Nil, processing_error), + ) + PostgresProjectionHandler( + project: fn(pog.Connection, factos.Recorded(event)) -> + Result(Nil, processing_error), + ) +} + /// A managed durable event subscription. /// /// The subscription owns notification lifecycle and its durable cursor. -/// Application code only handles decoded events. +/// Application code handles complete recorded event envelopes. pub opaque type Subscription(event, processing_error) { Builder( name: String, queries_connection: pog.Connection, listen_config: pog.Config, query: factos.Query, + start_from: SubscriptionStart, codec: EventCodec(event), - handle: fn(factos.Decoded(event)) -> Result(Nil, processing_error), + handler: SubscriptionHandler(event, processing_error), ) } @@ -156,13 +213,28 @@ pub type SubscriptionStopError { SubscriptionStopTimedOut } +type SubscriptionCheckpoint { + SubscriptionCheckpoint(cursor: factos.SequencePosition, generation: Int) +} + +type SubscriptionCheckpointWriteError { + SubscriptionGenerationChanged + SubscriptionCheckpointError(error: SubscriptionError) +} + +type ProjectionTransactionError(processing_error) { + ProjectionCallbackFailed(error: processing_error) + ProjectionCheckpointFailed(error: SubscriptionCheckpointWriteError) +} + type SubscriptionState(event, processing_error) { SubscriptionState( subscription: Subscription(event, processing_error), subject: process.Subject(SubscriptionMessage), - cursor: factos.SequencePosition, + checkpoint: SubscriptionCheckpoint, notifications_pid: process.Pid, - reference: reference.Reference, + event_reference: reference.Reference, + reset_reference: reference.Reference, monitor: process.Monitor, timer: option.Option(process.Timer), ) @@ -175,6 +247,15 @@ type SubscriptionMessage { SubscriptionListenerDown(down: process.Down) } +type RoutedSubscriptionNotification { + EventNotificationPayload(payload: String) + ResetNotificationPayload(payload: String) +} + +type ResetNotification { + ResetNotification(name: String, generation: Int) +} + type NotificationRouting { NotificationRouting( cursor: factos.SequencePosition, @@ -188,6 +269,9 @@ type NotificationDescriptor { NotificationDescriptor( cursor: factos.SequencePosition, previous_cursor: factos.SequencePosition, + id: String, + stream: String, + revision: Int, descriptor: factos.EventDescriptor, ) } @@ -196,7 +280,7 @@ type InlineNotification(event) { InlineNotification( cursor: factos.SequencePosition, previous_cursor: factos.SequencePosition, - event: factos.Decoded(event), + event: factos.Recorded(event), ) } @@ -204,6 +288,8 @@ const subscription_batch_size = 100 const subscription_reconciliation_interval_milliseconds = 30_000 +const subscription_reset_channel = "factos_pog_subscription_resets" + type DispatchQuery { StreamQuery ContextQuery(factos.Query) @@ -312,6 +398,65 @@ pub fn dispatch( } } +/// Dispatch a command, then wait for selected durable subscriptions. +/// +/// A post-commit wait failure carries the committed dispatch so callers never +/// mistake an observation timeout for a failed command. +pub fn dispatch_and_wait( + builder: DispatchBuilder(command, state, event, domain_error), + command: command, + event_id event_id: fn() -> String, + subscriptions subscriptions: List(String), + timeout timeout: duration.Duration, +) -> Result(Dispatch(event), DispatchWaitError(event, domain_error)) { + let DispatchBuilder(connection:, ..) = builder + case dispatch(builder, command, event_id:) { + Error(error) -> Error(DispatchNotCommitted(error:)) + Ok(dispatch) -> + case + wait_for_subscriptions( + connection, + subscriptions:, + through: dispatch.append.position, + timeout:, + ) + { + Ok(Nil) -> Ok(dispatch) + Error(reason) -> + Error(DispatchCommittedButNotObserved(dispatch:, reason:)) + } + } +} + +/// Wait until every named subscription cursor has observed `through`. +pub fn wait_for_subscriptions( + connection: pog.Connection, + subscriptions subscriptions: List(String), + through through: factos.SequencePosition, + timeout timeout: duration.Duration, +) -> Result(Nil, SubscriptionError) { + case subscriptions, through { + [], _ -> Ok(Nil) + _, factos.NoPosition -> Ok(Nil) + [_, ..], factos.SequencePosition(position) -> + case position < 0 { + True -> Error(InvalidSubscriptionPosition(position: through)) + False -> { + let subscriptions = list.unique(subscriptions) + use _ <- result.try(validate_subscription_names(subscriptions)) + let timeout = non_negative_wait_duration(timeout) + let deadline = timestamp.add(timestamp.system_time(), timeout) + wait_for_subscription_cursors( + connection, + subscriptions, + through, + deadline, + ) + } + } + } +} + pub type Error(domain_error) { /// The decider rejected the command with a domain error. DomainError(domain_error) @@ -354,6 +499,7 @@ fn stop_notifications(pid: process.Pid) -> Nil /// PostgreSQL notifications are advisory and may be coalesced or missed while /// disconnected. Establish the listener before reading durable events after the /// subscription checkpoint, and repeat that catch-up after every reconnect. +@internal pub fn listen( connection: pog.NotificationsConnection, ) -> Result(reference.Reference, Nil) { @@ -375,7 +521,7 @@ pub fn unlisten( /// event rows. Recovery reads run after startup, a notification gap, or periodic /// reconnect reconciliation. /// -/// `handle` receives only the decoded event and descriptor. After it returns +/// `handle` receives the complete recorded event envelope. After it returns /// `Ok(Nil)`, the worker stores the private cursor. A crash between an external /// side effect and that cursor update can replay the event, so external effects /// must be idempotent. @@ -384,24 +530,372 @@ pub fn new_subscription( config config: pog.Config, name name: String, query query: factos.Query, + start_from start_from: SubscriptionStart, + codec codec: EventCodec(event), + handle handle: fn(factos.Recorded(event)) -> Result(Nil, processing_error), +) -> Result(Subscription(event, processing_error), SubscriptionError) { + use _ <- result.try(validate_subscription_name(name)) + Ok(Builder( + name:, + listen_config: config, + queries_connection: connection, + query:, + start_from:, + codec:, + handler: OrdinarySubscriptionHandler(handle:), + )) +} + +/// Configure a subscription whose projection and cursor commit atomically. +/// +/// `project` runs once per matching event inside a Pog transaction. Use only +/// the supplied transaction-scoped connection and do not start a nested +/// transaction or call external services from the callback. +pub fn new_projection_subscription( + connection connection: pog.Connection, + config config: pog.Config, + name name: String, + query query: factos.Query, + start_from start_from: SubscriptionStart, codec codec: EventCodec(event), - handle handle: fn(factos.Decoded(event)) -> Result(Nil, processing_error), -) -> Result(Subscription(event, processing_error), Nil) { + project project: fn(pog.Connection, factos.Recorded(event)) -> + Result(Nil, processing_error), +) -> Result(Subscription(event, processing_error), SubscriptionError) { + use _ <- result.try(validate_subscription_name(name)) + Ok(Builder( + name:, + listen_config: config, + queries_connection: connection, + query:, + start_from:, + codec:, + handler: PostgresProjectionHandler(project:), + )) +} + +/// Read a subscription's durable checkpoint and exact global event-log lag. +/// +/// `events_behind` counts all event-log rows after the cursor, including rows +/// that do not match the subscription query. +pub fn subscription_status( + connection: pog.Connection, + name name: String, +) -> Result(SubscriptionStatus, SubscriptionError) { + use _ <- result.try(validate_subscription_name(name)) + use returned <- result.try( + pog.query( + "select + subscription.cursor, + subscription.generation, + coalesce((select max(position) from factos_events), -1), + ( + select count(*) + from factos_events + where position > subscription.cursor + ) + from factos_subscriptions as subscription + where subscription.name = $1", + ) + |> pog.parameter(pog.text(name)) + |> pog.returning(subscription_status_decoder(name)) + |> pog.execute(on: connection) + |> result.map_error(SubscriptionStoreError), + ) + case returned.rows { + [] -> Error(SubscriptionNotFound(name:)) + [status, ..] -> Ok(status) + } +} + +/// Atomically reposition a durable subscription and notify its live worker. +pub fn reset_subscription( + connection: pog.Connection, + name name: String, + start_from start_from: SubscriptionStart, +) -> Result(Nil, SubscriptionError) { + use _ <- result.try(validate_subscription_name(name)) + run_subscription_transaction(connection, fn(transaction_connection) { + use _ <- result.try(load_subscription_checkpoint( + transaction_connection, + name, + )) + use event_log_position <- result.try(read_event_log_cursor( + transaction_connection, + )) + use cursor <- result.try(resolve_subscription_start( + name, + start_from, + event_log_position, + )) + use updated <- result.try( + pog.query( + "update factos_subscriptions + set cursor = $2, generation = generation + 1 + where name = $1 + returning cursor, generation", + ) + |> pog.parameter(pog.text(name)) + |> pog.parameter(pog.int(cursor_to_int(cursor))) + |> pog.returning(subscription_checkpoint_decoder()) + |> pog.execute(on: transaction_connection) + |> result.map_error(SubscriptionStoreError), + ) + case updated.rows { + [] -> Error(SubscriptionNotFound(name:)) + [SubscriptionCheckpoint(generation:, ..), ..] -> + notify_subscription_reset(transaction_connection, name, generation) + } + }) +} + +fn validate_subscription_name(name: String) -> Result(Nil, SubscriptionError) { case string.trim(name) { - "" -> Error(Nil) - _ -> { - Ok(Builder( - name:, - listen_config: config, - queries_connection: connection, - query:, - codec:, - handle:, - )) + "" -> Error(InvalidSubscriptionName(name:)) + _ -> Ok(Nil) + } +} + +fn validate_subscription_names( + names: List(String), +) -> Result(Nil, SubscriptionError) { + case names { + [] -> Ok(Nil) + [name, ..remaining] -> { + use _ <- result.try(validate_subscription_name(name)) + validate_subscription_names(remaining) + } + } +} + +fn non_negative_wait_duration(timeout: duration.Duration) -> duration.Duration { + case duration.to_milliseconds(timeout) <= 0 { + True -> duration.milliseconds(0) + False -> timeout + } +} + +fn wait_for_subscription_cursors( + connection: pog.Connection, + subscriptions: List(String), + through: factos.SequencePosition, + deadline: timestamp.Timestamp, +) -> Result(Nil, SubscriptionError) { + use cursors <- result.try(read_subscription_cursors(connection, subscriptions)) + use pending_names <- result.try(pending_subscription_names( + subscriptions, + cursors, + through, + )) + case pending_names { + [] -> Ok(Nil) + [_, ..] -> { + let remaining = timestamp.difference(timestamp.system_time(), deadline) + let remaining_milliseconds = duration.to_milliseconds(remaining) + case remaining_milliseconds <= 0 { + True -> { + use statuses <- result.try(load_subscription_statuses( + connection, + pending_names, + )) + let pending = + statuses + |> list.filter(fn(status) { + let SubscriptionStatus(cursor:, ..) = status + cursor_is_after(through, cursor) + }) + case pending { + [] -> Ok(Nil) + [_, ..] -> Error(SubscriptionWaitTimedOut(through:, pending:)) + } + } + False -> { + process.sleep(int.min(remaining_milliseconds, 20)) + wait_for_subscription_cursors( + connection, + subscriptions, + through, + deadline, + ) + } + } } } } +fn read_subscription_cursors( + connection: pog.Connection, + subscriptions: List(String), +) -> Result(dict.Dict(String, factos.SequencePosition), SubscriptionError) { + use returned <- result.try( + pog.query( + "select name, cursor + from factos_subscriptions + where name = any($1::text[])", + ) + |> pog.parameter(pog.array(pog.text, subscriptions)) + |> pog.returning(subscription_cursor_decoder()) + |> pog.execute(on: connection) + |> result.map_error(SubscriptionStoreError), + ) + returned.rows + |> dict.from_list + |> Ok +} + +fn pending_subscription_names( + subscriptions: List(String), + cursors: dict.Dict(String, factos.SequencePosition), + through: factos.SequencePosition, +) -> Result(List(String), SubscriptionError) { + collect_pending_subscription_names( + subscriptions, + cursors, + through, + pending: [], + ) +} + +fn collect_pending_subscription_names( + subscriptions: List(String), + cursors: dict.Dict(String, factos.SequencePosition), + through: factos.SequencePosition, + pending pending: List(String), +) -> Result(List(String), SubscriptionError) { + case subscriptions { + [] -> Ok(list.reverse(pending)) + [name, ..remaining] -> + case dict.get(cursors, name) { + Error(Nil) -> Error(SubscriptionNotFound(name:)) + Ok(cursor) -> { + let pending = case cursor_is_after(through, cursor) { + True -> [name, ..pending] + False -> pending + } + collect_pending_subscription_names( + remaining, + cursors, + through, + pending:, + ) + } + } + } +} + +fn load_subscription_statuses( + connection: pog.Connection, + subscriptions: List(String), +) -> Result(List(SubscriptionStatus), SubscriptionError) { + use returned <- result.try( + pog.query( + "select + subscription.name, + subscription.cursor, + subscription.generation, + coalesce((select max(position) from factos_events), -1), + ( + select count(*) + from factos_events + where position > subscription.cursor + ) + from factos_subscriptions as subscription + where subscription.name = any($1::text[])", + ) + |> pog.parameter(pog.array(pog.text, subscriptions)) + |> pog.returning(named_subscription_status_decoder()) + |> pog.execute(on: connection) + |> result.map_error(SubscriptionStoreError), + ) + let statuses = + returned.rows + |> list.map(fn(status) { + let SubscriptionStatus(name:, ..) = status + #(name, status) + }) + |> dict.from_list + order_subscription_statuses(subscriptions, statuses, ordered: []) +} + +fn order_subscription_statuses( + subscriptions: List(String), + statuses: dict.Dict(String, SubscriptionStatus), + ordered ordered: List(SubscriptionStatus), +) -> Result(List(SubscriptionStatus), SubscriptionError) { + case subscriptions { + [] -> Ok(list.reverse(ordered)) + [name, ..remaining] -> + case dict.get(statuses, name) { + Error(Nil) -> Error(SubscriptionNotFound(name:)) + Ok(status) -> + order_subscription_statuses(remaining, statuses, ordered: [ + status, + ..ordered + ]) + } + } +} + +fn subscription_cursor_decoder() -> decode.Decoder( + #(String, factos.SequencePosition), +) { + use name <- decode.field(0, decode.string) + use cursor <- decode.field(1, decode.int) + decode.success(#(name, cursor_from_int(cursor))) +} + +fn named_subscription_status_decoder() -> decode.Decoder(SubscriptionStatus) { + use name <- decode.field(0, decode.string) + use cursor <- decode.field(1, decode.int) + use generation <- decode.field(2, decode.int) + use event_log_position <- decode.field(3, decode.int) + use events_behind <- decode.field(4, decode.int) + decode.success(SubscriptionStatus( + name:, + cursor: cursor_from_int(cursor), + generation:, + event_log_position: cursor_from_int(event_log_position), + events_behind:, + )) +} + +fn run_subscription_transaction( + connection: pog.Connection, + work: fn(pog.Connection) -> Result(value, SubscriptionError), +) -> Result(value, SubscriptionError) { + case pog.transaction(connection, work) { + Ok(value) -> Ok(value) + Error(pog.TransactionQueryError(error)) -> + Error(SubscriptionStoreError(error:)) + Error(pog.TransactionRolledBack(error)) -> Error(error) + } +} + +fn notify_subscription_reset( + connection: pog.Connection, + name: String, + generation: Int, +) -> Result(Nil, SubscriptionError) { + pog.query( + "with notified as materialized ( + select pg_catalog.pg_notify( + 'factos_pog_subscription_resets', + jsonb_build_object( + 'name', $1::text, + 'generation', $2::bigint + )::text + ) + ) + select 1 + from notified", + ) + |> pog.parameter(pog.text(name)) + |> pog.parameter(pog.int(generation)) + |> pog.returning(int_field_decoder()) + |> pog.execute(on: connection) + |> result.map(fn(_) { Nil }) + |> result.map_error(SubscriptionStoreError) +} + /// Start a subscription's local supervision tree. /// /// The tree starts the dedicated Pog notification session before the @@ -520,11 +1014,7 @@ pub fn read( )) } -/// Read a bounded page of recorded events after a global sequence position. -/// -/// This low-level primitive supports application-owned schedulers, projection -/// checkpoints, and subscription protocols. Events are filtered before decoding -/// and returned in global position order. +@internal pub fn read_after( connection: pog.Connection, query query: factos.Query, @@ -582,40 +1072,49 @@ fn initialise_subscription( ) process.trap_exits(True) let monitor = process.monitor(notifications_pid) - use reference <- result.try( + use event_reference <- result.try( listen(notifications) |> result.map_error(fn(_) { process.demonitor_process(monitor) "subscription notification listener could not LISTEN" }), ) - use cursor <- result.try( - load_subscription_cursor(subscription.queries_connection, subscription.name) - |> result.map_error(fn(reason) { - unlisten(notifications, reference) + use reset_reference <- result.try( + listen_on_channel(notifications, subscription_reset_channel) + |> result.map_error(fn(_) { + unlisten(notifications, event_reference) + process.demonitor_process(monitor) + "subscription reset listener could not LISTEN" + }), + ) + use checkpoint <- result.try( + initialise_subscription_checkpoint( + subscription.queries_connection, + subscription.name, + subscription.start_from, + ) + |> result.map_error(fn(error) { + unlisten(notifications, event_reference) + unlisten(notifications, reset_reference) process.demonitor_process(monitor) - reason + string.inspect(error) }), ) let selector = process.new_selector() |> process.select(subject) - |> pog.select_notifications(fn(notification) { - SubscriptionNotification(notification:) - }) - |> process.select_specific_monitor(monitor, fn(down) { - SubscriptionListenerDown(down:) - }) - |> process.select_trapped_exits(fn(message) { - SubscriptionParentExit(message:) - }) + |> pog.select_notifications(SubscriptionNotification) + |> process.select_specific_monitor(monitor, SubscriptionListenerDown) + |> process.select_trapped_exits(SubscriptionParentExit) + process.send(subject, CheckSubscription) SubscriptionState( subscription:, subject:, - cursor:, + checkpoint:, notifications_pid:, - reference:, + event_reference:, + reset_reference:, monitor:, timer: option.None, ) @@ -631,8 +1130,11 @@ fn handle_subscription_message( case message { CheckSubscription -> run_subscription_cycle(state) SubscriptionNotification(notification:) -> - case subscription_notification_payload(state, notification) { - Ok(payload) -> handle_subscription_notification(state, payload) + case route_subscription_notification(state, notification) { + Ok(EventNotificationPayload(payload:)) -> + handle_subscription_notification(state, payload) + Ok(ResetNotificationPayload(payload:)) -> + handle_subscription_reset_notification(state, payload) Error(Nil) -> actor.continue(state) } SubscriptionParentExit(message: _) -> { @@ -660,19 +1162,28 @@ fn run_subscription_cycle( state: SubscriptionState(event, processing_error), ) -> actor.Next(SubscriptionState(event, processing_error), SubscriptionMessage) { cancel_subscription_timer(state.timer) - case catch_up_subscription(state.subscription, state.cursor) { - Error(reason) -> stop_subscription_abnormally(reason) - Ok(cursor) -> { - let timer = - process.send_after( - state.subject, - subscription_reconciliation_interval_milliseconds, - CheckSubscription, - ) - actor.continue( - SubscriptionState(..state, cursor:, timer: option.Some(timer)), - ) - } + case + load_subscription_checkpoint( + state.subscription.queries_connection, + state.subscription.name, + ) + { + Error(error) -> stop_subscription_abnormally(string.inspect(error)) + Ok(checkpoint) -> + case catch_up_subscription(state.subscription, checkpoint) { + Error(reason) -> stop_subscription_abnormally(reason) + Ok(checkpoint) -> { + let timer = + process.send_after( + state.subject, + subscription_reconciliation_interval_milliseconds, + CheckSubscription, + ) + actor.continue( + SubscriptionState(..state, checkpoint:, timer: option.Some(timer)), + ) + } + } } } @@ -688,16 +1199,12 @@ fn cancel_subscription_timer(timer: option.Option(process.Timer)) -> Nil { fn drain_subscription( configuration: Subscription(event, processing_error), - cursor: factos.SequencePosition, + checkpoint: SubscriptionCheckpoint, through: factos.SequencePosition, -) -> Result(factos.SequencePosition, String) { +) -> Result(SubscriptionCheckpoint, String) { + let SubscriptionCheckpoint(cursor:, ..) = checkpoint case cursor_is_after(through, cursor) { - False -> - store_subscription_cursor( - configuration.queries_connection, - configuration.name, - cursor, - ) + False -> Ok(checkpoint) True -> case read_events_through( @@ -717,35 +1224,64 @@ fn drain_subscription( }), ) Ok([]) -> - store_subscription_cursor( + store_subscription_checkpoint( configuration.queries_connection, configuration.name, + checkpoint, through, ) + |> result.map_error(subscription_checkpoint_write_error_to_string) Ok(events) -> { - use next_cursor <- result.try(handle_recorded_events( + use next_checkpoint <- result.try(handle_recorded_events( configuration, events, - cursor, + checkpoint, )) - drain_subscription(configuration, next_cursor, through) + drain_subscription(configuration, next_checkpoint, through) } } } } -fn subscription_notification_payload( +fn route_subscription_notification( state: SubscriptionState(event, processing_error), notification: pog.Notification, -) -> Result(String, Nil) { +) -> Result(RoutedSubscriptionNotification, Nil) { let pog.Notify(pid:, reference:, channel:, payload:) = notification - case - pid == state.notifications_pid - && reference == state.reference - && channel == event_notification_channel - { - True -> Ok(payload) + case pid == state.notifications_pid { False -> Error(Nil) + True -> + case + reference == state.event_reference + && channel == event_notification_channel + { + True -> Ok(EventNotificationPayload(payload:)) + False -> + case + reference == state.reset_reference + && channel == subscription_reset_channel + { + True -> Ok(ResetNotificationPayload(payload:)) + False -> Error(Nil) + } + } + } +} + +fn handle_subscription_reset_notification( + state: SubscriptionState(event, processing_error), + payload: String, +) -> actor.Next(SubscriptionState(event, processing_error), SubscriptionMessage) { + case json.parse(payload, using: reset_notification_decoder()) { + Error(_) -> actor.continue(state) + Ok(ResetNotification(name:, generation:)) -> { + let SubscriptionCheckpoint(generation: held_generation, ..) = + state.checkpoint + case name == state.subscription.name && generation > held_generation { + True -> run_subscription_cycle(state) + False -> actor.continue(state) + } + } } } @@ -753,12 +1289,13 @@ fn handle_subscription_notification( state: SubscriptionState(event, processing_error), payload: String, ) -> actor.Next(SubscriptionState(event, processing_error), SubscriptionMessage) { + let SubscriptionCheckpoint(cursor: held_cursor, ..) = state.checkpoint case json.parse(payload, using: notification_routing_decoder()) { Error(_) -> run_subscription_cycle(state) Ok(NotificationRouting(cursor:, previous_cursor:, type_:, tags:)) -> case - cursor_is_after(cursor, state.cursor), - previous_cursor == state.cursor + cursor_is_after(cursor, held_cursor), + previous_cursor == held_cursor { False, _ -> actor.continue(state) True, False -> run_subscription_cycle(state) @@ -798,15 +1335,14 @@ fn handle_subscription_notification( fn handle_inline_event( state: SubscriptionState(event, processing_error), - decoded: factos.Decoded(event), - cursor: factos.SequencePosition, + recorded: factos.Recorded(event), + _cursor: factos.SequencePosition, ) -> actor.Next(SubscriptionState(event, processing_error), SubscriptionMessage) { - case state.subscription.handle(decoded) { - Error(error) -> - stop_subscription_abnormally( - "subscription event handler failed: " <> string.inspect(error), - ) - Ok(Nil) -> checkpoint_inline_cursor(state, cursor) + case + process_subscription_event(state.subscription, state.checkpoint, recorded) + { + Error(reason) -> stop_subscription_abnormally(reason) + Ok(checkpoint) -> actor.continue(SubscriptionState(..state, checkpoint:)) } } @@ -815,130 +1351,276 @@ fn checkpoint_inline_cursor( cursor: factos.SequencePosition, ) -> actor.Next(SubscriptionState(event, processing_error), SubscriptionMessage) { case - store_subscription_cursor( + store_subscription_checkpoint( state.subscription.queries_connection, state.subscription.name, + state.checkpoint, cursor, ) { - Error(reason) -> stop_subscription_abnormally(reason) - Ok(cursor) -> actor.continue(SubscriptionState(..state, cursor:)) + Error(error) -> + stop_subscription_abnormally( + subscription_checkpoint_write_error_to_string(error), + ) + Ok(checkpoint) -> actor.continue(SubscriptionState(..state, checkpoint:)) } } fn catch_up_subscription( configuration: Subscription(event, processing_error), - cursor: factos.SequencePosition, -) -> Result(factos.SequencePosition, String) { - use through <- result.try(read_event_log_cursor( - configuration.queries_connection, - )) - drain_subscription(configuration, cursor, through) + checkpoint: SubscriptionCheckpoint, +) -> Result(SubscriptionCheckpoint, String) { + use through <- result.try( + read_event_log_cursor(configuration.queries_connection) + |> result.map_error(string.inspect), + ) + drain_subscription(configuration, checkpoint, through) } fn handle_recorded_events( configuration: Subscription(event, processing_error), events: List(factos.Recorded(event)), - cursor: factos.SequencePosition, -) -> Result(factos.SequencePosition, String) { + checkpoint: SubscriptionCheckpoint, +) -> Result(SubscriptionCheckpoint, String) { case events { - [] -> Ok(cursor) - [factos.Recorded(position:, event:, descriptor:, ..), ..rest] -> { + [] -> Ok(checkpoint) + [recorded, ..rest] -> { + use checkpoint <- result.try(process_subscription_event( + configuration, + checkpoint, + recorded, + )) + handle_recorded_events(configuration, rest, checkpoint) + } + } +} + +fn process_subscription_event( + configuration: Subscription(event, processing_error), + checkpoint: SubscriptionCheckpoint, + recorded: factos.Recorded(event), +) -> Result(SubscriptionCheckpoint, String) { + let factos.Recorded(position:, ..) = recorded + case configuration.handler { + OrdinarySubscriptionHandler(handle:) -> { use _ <- result.try( - configuration.handle(factos.Decoded(event:, descriptor:)) + handle(recorded) |> result.map_error(fn(error) { "subscription event handler failed: " <> string.inspect(error) }), ) - use cursor <- result.try(store_subscription_cursor( + store_subscription_checkpoint( configuration.queries_connection, configuration.name, + checkpoint, position, + ) + |> result.map_error(subscription_checkpoint_write_error_to_string) + } + PostgresProjectionHandler(project:) -> + process_projection_event(configuration, checkpoint, recorded, project) + } +} + +fn process_projection_event( + configuration: Subscription(event, processing_error), + checkpoint: SubscriptionCheckpoint, + recorded: factos.Recorded(event), + project: fn(pog.Connection, factos.Recorded(event)) -> + Result(Nil, processing_error), +) -> Result(SubscriptionCheckpoint, String) { + let factos.Recorded(position:, ..) = recorded + case + pog.transaction( + configuration.queries_connection, + fn(transaction_connection) { + use _ <- result.try( + project(transaction_connection, recorded) + |> result.map_error(fn(error) { ProjectionCallbackFailed(error:) }), + ) + store_subscription_checkpoint( + transaction_connection, + configuration.name, + checkpoint, + position, + ) + |> result.map_error(fn(error) { ProjectionCheckpointFailed(error:) }) + }, + ) + { + Ok(checkpoint) -> Ok(checkpoint) + Error(pog.TransactionQueryError(error)) -> + Error( + "subscription projection transaction failed: " + <> query_error_to_string(error), + ) + Error(pog.TransactionRolledBack(ProjectionCallbackFailed(error:))) -> + Error("subscription projection failed: " <> string.inspect(error)) + Error(pog.TransactionRolledBack(ProjectionCheckpointFailed(error:))) -> + Error(subscription_checkpoint_write_error_to_string(error)) + } +} + +fn initialise_subscription_checkpoint( + connection: pog.Connection, + name: String, + start_from: SubscriptionStart, +) -> Result(SubscriptionCheckpoint, SubscriptionError) { + use existing <- result.try(load_optional_subscription_checkpoint( + connection, + name, + )) + case existing { + option.Some(checkpoint) -> Ok(checkpoint) + option.None -> { + use event_log_position <- result.try(read_event_log_cursor(connection)) + use cursor <- result.try(resolve_subscription_start( + name, + start_from, + event_log_position, )) - handle_recorded_events(configuration, rest, cursor) + use inserted <- result.try( + pog.query( + "insert into factos_subscriptions (name, cursor, generation) + values ($1, $2, 0) + on conflict (name) do nothing + returning cursor, generation", + ) + |> pog.parameter(pog.text(name)) + |> pog.parameter(pog.int(cursor_to_int(cursor))) + |> pog.returning(subscription_checkpoint_decoder()) + |> pog.execute(on: connection) + |> result.map_error(SubscriptionStoreError), + ) + case inserted.rows { + [checkpoint, ..] -> Ok(checkpoint) + [] -> load_subscription_checkpoint(connection, name) + } } } } -fn load_subscription_cursor( +fn resolve_subscription_start( + name: String, + start_from: SubscriptionStart, + event_log_position: factos.SequencePosition, +) -> Result(factos.SequencePosition, SubscriptionError) { + case start_from { + Origin -> Ok(factos.NoPosition) + Current -> Ok(event_log_position) + After(position:) -> + case position { + factos.NoPosition -> Ok(factos.NoPosition) + factos.SequencePosition(value) -> + case value < 0, cursor_is_after(position, event_log_position) { + True, _ -> Error(InvalidSubscriptionPosition(position:)) + False, True -> + Error(SubscriptionStartAfterEventLog( + name:, + requested: position, + event_log_position:, + )) + False, False -> Ok(position) + } + } + } +} + +fn load_subscription_checkpoint( connection: pog.Connection, name: String, -) -> Result(factos.SequencePosition, String) { - use _ <- result.try( - pog.query( - "insert into factos_subscriptions (name, cursor) - values ($1, -1) - on conflict (name) do nothing", - ) - |> pog.parameter(pog.text(name)) - |> pog.execute(on: connection) - |> result.map_error(fn(error) { - "subscription cursor initialization failed: " - <> query_error_to_string(error) - }), - ) +) -> Result(SubscriptionCheckpoint, SubscriptionError) { + use checkpoint <- result.try(load_optional_subscription_checkpoint( + connection, + name, + )) + case checkpoint { + option.Some(checkpoint) -> Ok(checkpoint) + option.None -> Error(SubscriptionNotFound(name:)) + } +} + +fn load_optional_subscription_checkpoint( + connection: pog.Connection, + name: String, +) -> Result(option.Option(SubscriptionCheckpoint), SubscriptionError) { use returned <- result.try( pog.query( - "select cursor + "select cursor, generation from factos_subscriptions where name = $1", ) |> pog.parameter(pog.text(name)) - |> pog.returning(int_field_decoder()) + |> pog.returning(subscription_checkpoint_decoder()) |> pog.execute(on: connection) - |> result.map_error(fn(error) { - "subscription cursor load failed: " <> query_error_to_string(error) - }), + |> result.map_error(SubscriptionStoreError), ) case returned.rows { - [cursor] -> Ok(cursor_from_int(cursor)) - [] -> Error("subscription cursor is missing") - [_, _, ..] -> Error("subscription cursor is duplicated") + [] -> Ok(option.None) + [checkpoint, ..] -> Ok(option.Some(checkpoint)) } } -fn store_subscription_cursor( +fn store_subscription_checkpoint( connection: pog.Connection, name: String, + expected: SubscriptionCheckpoint, cursor: factos.SequencePosition, -) -> Result(factos.SequencePosition, String) { +) -> Result(SubscriptionCheckpoint, SubscriptionCheckpointWriteError) { + let SubscriptionCheckpoint(generation:, ..) = expected use returned <- result.try( pog.query( "update factos_subscriptions - set cursor = greatest(cursor, $2) - where name = $1 - returning cursor", + set cursor = greatest(cursor, $3) + where name = $1 and generation = $2 + returning cursor, generation", ) |> pog.parameter(pog.text(name)) + |> pog.parameter(pog.int(generation)) |> pog.parameter(pog.int(cursor_to_int(cursor))) - |> pog.returning(int_field_decoder()) + |> pog.returning(subscription_checkpoint_decoder()) |> pog.execute(on: connection) |> result.map_error(fn(error) { - "subscription cursor store failed: " <> query_error_to_string(error) + SubscriptionCheckpointError(error: SubscriptionStoreError(error:)) }), ) case returned.rows { - [cursor] -> Ok(cursor_from_int(cursor)) - [] -> Error("subscription cursor is missing") - [_, _, ..] -> Error("subscription cursor is duplicated") + [checkpoint, ..] -> Ok(checkpoint) + [] -> { + use durable <- result.try( + load_optional_subscription_checkpoint(connection, name) + |> result.map_error(fn(error) { SubscriptionCheckpointError(error:) }), + ) + case durable { + option.None -> + Error(SubscriptionCheckpointError(error: SubscriptionNotFound(name:))) + option.Some(_) -> Error(SubscriptionGenerationChanged) + } + } + } +} + +fn subscription_checkpoint_write_error_to_string( + error: SubscriptionCheckpointWriteError, +) -> String { + case error { + SubscriptionGenerationChanged -> "subscription generation changed" + SubscriptionCheckpointError(error:) -> string.inspect(error) } } fn read_event_log_cursor( connection: pog.Connection, -) -> Result(factos.SequencePosition, String) { +) -> Result(factos.SequencePosition, SubscriptionError) { use returned <- result.try( pog.query("select coalesce(max(position), -1) from factos_events") |> pog.returning(int_field_decoder()) |> pog.execute(on: connection) - |> result.map_error(fn(error) { - "subscription event cursor read failed: " <> query_error_to_string(error) - }), + |> result.map_error(SubscriptionStoreError), ) case returned.rows { - [cursor] -> Ok(cursor_from_int(cursor)) - [] -> Error("subscription event cursor is missing") - [_, _, ..] -> Error("subscription event cursor is duplicated") + [] -> Ok(factos.NoPosition) + [cursor, ..] -> Ok(cursor_from_int(cursor)) } } @@ -983,6 +1665,12 @@ fn read_events_through( } } +fn reset_notification_decoder() -> decode.Decoder(ResetNotification) { + use name <- decode.field("name", decode.string) + use generation <- decode.field("generation", decode.int) + decode.success(ResetNotification(name:, generation:)) +} + fn notification_routing_decoder() -> decode.Decoder(NotificationRouting) { use cursor <- decode.field("cursor", decode.int) use previous_cursor <- decode.field("previous_cursor", decode.int) @@ -1005,6 +1693,9 @@ fn notification_routing_decoder() -> decode.Decoder(NotificationRouting) { fn notification_descriptor_decoder() -> decode.Decoder(NotificationDescriptor) { use cursor <- decode.field("cursor", decode.int) use previous_cursor <- decode.field("previous_cursor", decode.int) + use id <- decode.subfield(["event", "id"], decode.string) + use stream <- decode.subfield(["event", "stream"], decode.string) + use revision <- decode.subfield(["event", "revision"], decode.int) use type_ <- decode.subfield( ["event", "type"], decode.string |> decode.map(factos.event_type), @@ -1021,6 +1712,9 @@ fn notification_descriptor_decoder() -> decode.Decoder(NotificationDescriptor) { decode.success(NotificationDescriptor( cursor: cursor_from_int(cursor), previous_cursor: cursor_from_int(previous_cursor), + id:, + stream:, + revision:, descriptor: factos.EventDescriptor( type_:, version:, @@ -1038,8 +1732,14 @@ fn decode_inline_notification( json.parse(payload, using: notification_descriptor_decoder()) |> result.replace_error(InvalidData), ) - let NotificationDescriptor(cursor:, previous_cursor:, descriptor:) = - notification + let NotificationDescriptor( + cursor:, + previous_cursor:, + id:, + stream:, + revision:, + descriptor:, + ) = notification let EventCodec(decode: select_decoder, ..) = codec use event_decoder <- result.try(select_decoder(descriptor)) use event <- result.try( @@ -1049,7 +1749,14 @@ fn decode_inline_notification( Ok(InlineNotification( cursor:, previous_cursor:, - event: factos.Decoded(event:, descriptor:), + event: factos.Recorded( + id:, + stream:, + revision:, + position: cursor, + event:, + descriptor:, + ), )) } @@ -1587,6 +2294,31 @@ fn current_revision( } } +fn subscription_status_decoder( + name: String, +) -> decode.Decoder(SubscriptionStatus) { + use cursor <- decode.field(0, decode.int) + use generation <- decode.field(1, decode.int) + use event_log_position <- decode.field(2, decode.int) + use events_behind <- decode.field(3, decode.int) + decode.success(SubscriptionStatus( + name:, + cursor: cursor_from_int(cursor), + generation:, + event_log_position: cursor_from_int(event_log_position), + events_behind:, + )) +} + +fn subscription_checkpoint_decoder() -> decode.Decoder(SubscriptionCheckpoint) { + use cursor <- decode.field(0, decode.int) + use generation <- decode.field(1, decode.int) + decode.success(SubscriptionCheckpoint( + cursor: cursor_from_int(cursor), + generation:, + )) +} + fn int_field_decoder() -> decode.Decoder(Int) { use value <- decode.field(0, decode.int) decode.success(value) diff --git a/backends/factos_pog/test/factos_pog_test.gleam b/backends/factos_pog/test/factos_pog_test.gleam index 8e731bc..9c9a79f 100644 --- a/backends/factos_pog/test/factos_pog_test.gleam +++ b/backends/factos_pog/test/factos_pog_test.gleam @@ -1,21 +1,23 @@ +import envoy import factos import factos/factos_pog import gleam/dynamic/decode import gleam/erlang/application import gleam/erlang/atom import gleam/erlang/process +import gleam/function import gleam/int import gleam/json import gleam/list import gleam/option.{Some} import gleam/otp/actor import gleam/otp/static_supervisor +import gleam/result import gleam/string +import gleam/time/duration +import global_value import pog import simplifile -import testcontainer -import testcontainer/error as testcontainer_error -import testcontainer_formulas/postgres import unitest import youid/uuid @@ -23,6 +25,17 @@ pub fn main() -> Nil { unitest.main() } +type SharedPostgres { + SharedPostgres( + admin_pool_pid: process.Pid, + admin_connection: pog.Connection, + host: String, + port: Int, + username: String, + password: String, + ) +} + type Command { RegisterUser(username: String) } @@ -46,6 +59,7 @@ type DomainError { type CounterCommand { Increment + DoNothing } type CounterEvent { @@ -57,7 +71,47 @@ type CounterState { } type ManagedSubscriptionMessage { - ManagedSubscriptionEvent(event: factos.Decoded(Event)) + ManagedSubscriptionEvent(event: factos.Recorded(Event)) +} + +type BlockingSubscriptionMessage { + BlockingHandlerStarted(event: factos.Recorded(Event)) +} + +type ProjectionSubscriptionMessage { + ProjectionAttemptStarted(attempt: Int, event: factos.Recorded(Event)) +} + +type ProjectionBarrierMessage { + ProjectionBarrierReady( + name: String, + event: factos.Recorded(Event), + release: process.Subject(Nil), + ) +} + +type DispatchWaitMessage { + DispatchWaitCompleted( + result: Result( + factos_pog.Dispatch(Event), + factos_pog.DispatchWaitError(Event, DomainError), + ), + ) +} + +type EventNotificationPayload { + EventNotificationPayload( + cursor: Int, + previous_cursor: Int, + id: String, + stream: String, + revision: Int, + type_: String, + version: Int, + tags: List(String), + correlation_id: String, + username: String, + ) } pub fn compatibility_migration_enforces_uuidv4_identity_contract_test() { @@ -74,40 +128,37 @@ pub fn dbmate_migrations_enforce_uuidv4_identity_contract_test() { assert_event_notification_objects(connection) } -pub fn jsonb_migration_preserves_existing_json_events_test() { - use connection <- with_test_connection() +pub fn dbmate_upgrade_and_v2_rollback_preserve_contract_test() { + use config, connection <- with_test_database() drop_schema(connection) let assert Ok(priv_directory) = application.priv_directory("factos_pog") - execute_dbmate_up( - connection, - priv_directory <> "/dbmate/20260703000100_factos_pog_event_store.sql", - ) + let v1_migration = + priv_directory <> "/dbmate/20260703000100_factos_pog_event_store.sql" + let v2_migration = + priv_directory <> "/dbmate/20260816000100_factos_pog_v2.sql" + execute_dbmate_up(connection, v1_migration) let assert Ok(_) = pog.query( " - insert into factos_events ( - id, stream, revision, type, version, tags, metadata, data - ) - values ( - 'b3b12f1d-6d85-4f1f-9c2a-94766a34f006', - 'migrated-user-renata', - 0, - 'UserRegistered', - 1, - '', - json_build_object('correlation_id', 'migration-correlation')::text, - convert_to(to_json('renata'::text)::text, 'UTF8') - ) - ", + insert into factos_events ( + id, stream, revision, type, version, tags, metadata, data + ) + values ( + 'b3b12f1d-6d85-4f1f-9c2a-94766a34f006', + 'migrated-user-renata', + 0, + 'UserRegistered', + 1, + '', + json_build_object('correlation_id', 'migration-correlation')::text, + convert_to(to_json('renata'::text)::text, 'UTF8') + ) + ", ) |> pog.execute(on: connection) - execute_dbmate_up( - connection, - priv_directory <> "/dbmate/20260816000100_factos_pog_v2.sql", - ) - + execute_dbmate_up(connection, v2_migration) let assert Ok([factos.Recorded(event:, descriptor:, ..)]) = factos_pog.read_after( connection, @@ -121,10 +172,92 @@ pub fn jsonb_migration_preserves_existing_json_events_test() { == Ok("migration-correlation") assert_uuidv4_identity_contract(connection) assert_event_notification_objects(connection) - execute_dbmate_down( - connection, - priv_directory <> "/dbmate/20260816000100_factos_pog_v2.sql", - ) + + let assert Ok(_) = + pog.query( + " + insert into factos_subscriptions (name, cursor) + values ('existing-v2-subscription', -1) + ", + ) + |> pog.execute(on: connection) + let assert Ok(generation_result) = + pog.query( + " + select generation + from factos_subscriptions + where name = 'existing-v2-subscription' + ", + ) + |> pog.returning(int_column_decoder()) + |> pog.execute(on: connection) + assert generation_result.rows == [0] + + let listener_config = + pog.Config( + ..config, + pool_name: process.new_name(prefix: "factos_pog_test_v2_listener"), + ) + let assert Ok(actor.Started(pid: notifications_pid, data: notifications)) = + pog.start_notifications(listener_config) + let assert Ok(listener) = factos_pog.listen(notifications) + let selector = + process.new_selector() + |> pog.select_notifications(function.identity) + process.sleep(100) + + let assert Ok(v2_dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "migrated-v2-notification", + decider: decider(), + codec: codec(), + ) + |> factos_pog.dispatch( + RegisterUser(username: "v2-notification"), + event_id: uuid.v4_string, + ) + let assert Ok(pog.Notify(payload: v2_payload, ..)) = + process.selector_receive(selector, 5000) + let _ = assert_event_notification_matches(v2_payload, v2_dispatch) + + factos_pog.unlisten(notifications, listener) + process.send_exit(notifications_pid) + execute_dbmate_down(connection, v2_migration) + + let assert Ok(legacy_event) = + pog.query( + " + select + metadata::jsonb ->> 'correlation_id', + convert_from(data, 'UTF8')::jsonb #>> '{}' + from factos_events + where id = 'b3b12f1d-6d85-4f1f-9c2a-94766a34f006' + ", + ) + |> pog.returning(string_pair_decoder()) + |> pog.execute(on: connection) + assert legacy_event.rows == [#("migration-correlation", "renata")] + + let assert Ok(removed_objects) = + pog.query( + " + select coalesce(to_regclass('factos_subscriptions')::text, '') + union all + select coalesce( + to_regprocedure('factos_pog_notify_event_appended()')::text, + '' + ) + union all + select coalesce( + to_regprocedure('factos_pog_lock_event_append()')::text, + '' + ) + ", + ) + |> pog.returning(string_column_decoder()) + |> pog.execute(on: connection) + assert removed_objects.rows == ["", "", ""] } pub fn event_notifications_wake_durable_catch_up_test() { @@ -152,43 +285,1015 @@ pub fn event_notifications_wake_durable_catch_up_test() { limit: 10, codec: codec(), ) - process.sleep(100) + process.sleep(100) + + let assert Ok(dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "user-renata", + decider: decider(), + codec: codec(), + ) + |> factos_pog.dispatch(RegisterUser("renata"), event_id: uuid.v4_string) + + let assert Ok(pog.Notify( + pid: notification_pid, + reference: notification_reference, + channel: "factos_pog_events", + payload:, + )) = process.selector_receive(selector, 5000) + assert notification_pid == notifications_pid + assert notification_reference == listener + let EventNotificationPayload(previous_cursor:, ..) = + assert_event_notification_matches(payload, dispatch) + assert previous_cursor == -1 + + let assert Ok(events) = + factos_pog.read_after( + connection, + query: factos.AllEvents, + after: factos.NoPosition, + limit: 10, + codec: codec(), + ) + assert events == dispatch.events + assert factos.react_all(welcome_reactor(), events) == [SendWelcome("renata")] + + let assert [last] = list.reverse(events) + let assert Ok([]) = + factos_pog.read_after( + connection, + query: factos.AllEvents, + after: last.position, + limit: 10, + codec: codec(), + ) + + // NOTIFY follows the transaction: rolling the event back exposes neither + // a durable event nor a wake-up edge. + let assert Error(pog.TransactionRolledBack(Nil)) = + pog.transaction(connection, fn(transaction_connection) { + let assert Ok(_) = + pog.query( + " + insert into factos_events ( + id, stream, revision, type, version, tags, metadata, data + ) + values ( + 'b3b12f1d-6d85-4f1f-9c2a-94766a34f099', + 'rolled-back', + 0, + 'UserRegistered', + 1, + '', + '{}'::jsonb, + '\"rolled-back\"'::jsonb + ) + ", + ) + |> pog.execute(on: transaction_connection) + Error(Nil) + }) + let assert Error(Nil) = process.selector_receive(selector, 200) + + factos_pog.unlisten(notifications, listener) + Nil +} + +pub fn v2_notification_payload_falls_back_to_ordered_read_test() { + use config, connection <- with_test_database() + reset_schema(connection) + let name = "v2-notification-fallback" + let deliveries = process.new_subject() + let subscription = + test_subscription( + connection, + config, + name, + deliveries, + handle: deliver_managed_subscription_event, + ) + let assert Ok(actor.Started(data: handle, ..)) = + factos_pog.start(subscription) + + let assert Ok(_) = + pog.query( + "alter table factos_events disable trigger factos_pog_event_insert_notify", + ) + |> pog.execute(on: connection) + let assert Ok(dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "v2-notification-fallback", + decider: decider(), + codec: codec(), + ) + |> factos_pog.dispatch( + RegisterUser(username: "v2-fallback"), + event_id: uuid.v4_string, + ) + let assert [factos.Recorded(id:, ..)] = dispatch.events + let assert Ok(_) = + pog.query( + "alter table factos_events enable trigger factos_pog_event_insert_notify", + ) + |> pog.execute(on: connection) + let assert Ok(_) = + pog.query( + " + with notified as materialized ( + select pg_catalog.pg_notify( + 'factos_pog_events', + jsonb_build_object( + 'cursor', position, + 'previous_cursor', -1, + 'event', jsonb_build_object( + 'type', type, + 'version', version, + 'tags', tags::jsonb, + 'metadata', metadata, + 'data', data + ) + )::text + ) + from factos_events + where id = $1 + ) + select 1 + from notified + ", + ) + |> pog.parameter(pog.text(id)) + |> pog.returning(int_column_decoder()) + |> pog.execute(on: connection) + + let assert Ok(ManagedSubscriptionEvent(event: recorded)) = + process.receive(deliveries, within: 10_000) + assert recorded == recorded_dispatch_event(dispatch) + let assert Ok(Nil) = factos_pog.stop(handle) + Nil +} + +pub fn dbmate_subscription_callback_preserves_recorded_envelope_test() { + use config, connection <- with_test_database() + reset_schema_from_dbmate(connection) + let deliveries = process.new_subject() + let subscription = + test_subscription( + connection, + config, + "dbmate-recorded-envelope", + deliveries, + handle: deliver_managed_subscription_event, + ) + let assert Ok(actor.Started(data: handle, ..)) = + factos_pog.start(subscription) + + let assert Ok(dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "dbmate-recorded-envelope", + decider: decider(), + codec: codec(), + ) + |> factos_pog.dispatch( + RegisterUser(username: "dbmate-envelope"), + event_id: uuid.v4_string, + ) + let assert Ok(ManagedSubscriptionEvent(event: recorded)) = + process.receive(deliveries, within: 10_000) + assert recorded == recorded_dispatch_event(dispatch) + + let assert Ok(Nil) = factos_pog.stop(handle) + Nil +} + +pub fn subscription_start_modes_initialize_new_rows_once_test() { + use config, connection <- with_test_database() + reset_schema(connection) + let current_deliveries = process.new_subject() + let assert Ok(first_dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "start-first", + decider: decider(), + codec: codec(), + ) + |> factos_pog.dispatch( + RegisterUser(username: "first"), + event_id: uuid.v4_string, + ) + let assert Ok(second_dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "start-second", + decider: decider(), + codec: codec(), + ) + |> factos_pog.dispatch( + RegisterUser(username: "second"), + event_id: uuid.v4_string, + ) + + let current_subscription = + test_subscription_from( + connection, + config, + "start-current", + current_deliveries, + start_from: factos_pog.Current, + handle: deliver_managed_subscription_event, + ) + let assert Ok(actor.Started(data: current_handle, ..)) = + factos_pog.start(current_subscription) + let assert Error(Nil) = process.receive(current_deliveries, within: 200) + let assert Ok(current_cursor) = + load_subscription_cursor(connection, "start-current") + assert current_cursor == second_dispatch.append.position + let assert Ok(Nil) = factos_pog.stop(current_handle) + + // Once the durable row exists, a new constructor start does not reposition it. + let existing_subscription = + test_subscription_from( + connection, + config, + "start-current", + current_deliveries, + start_from: factos_pog.Origin, + handle: deliver_managed_subscription_event, + ) + let assert Ok(actor.Started(data: existing_handle, ..)) = + factos_pog.start(existing_subscription) + let assert Error(Nil) = process.receive(current_deliveries, within: 200) + let assert Ok(third_dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "start-third", + decider: decider(), + codec: codec(), + ) + |> factos_pog.dispatch( + RegisterUser(username: "third"), + event_id: uuid.v4_string, + ) + let assert Ok(ManagedSubscriptionEvent(event: third_event)) = + process.receive(current_deliveries, within: 10_000) + assert third_event == recorded_dispatch_event(third_dispatch) + let assert Ok(Nil) = factos_pog.stop(existing_handle) + + let after_deliveries = process.new_subject() + let after_subscription = + test_subscription_from( + connection, + config, + "start-after", + after_deliveries, + start_from: factos_pog.After(position: first_dispatch.append.position), + handle: deliver_managed_subscription_event, + ) + let assert Ok(actor.Started(data: after_handle, ..)) = + factos_pog.start(after_subscription) + let assert Ok(ManagedSubscriptionEvent(event: after_second)) = + process.receive(after_deliveries, within: 10_000) + let assert Ok(ManagedSubscriptionEvent(event: after_third)) = + process.receive(after_deliveries, within: 10_000) + assert after_second == recorded_dispatch_event(second_dispatch) + assert after_third == recorded_dispatch_event(third_dispatch) + let assert Ok(Nil) = factos_pog.stop(after_handle) + + let origin_deliveries = process.new_subject() + let origin_subscription = + test_subscription_from( + connection, + config, + "start-origin", + origin_deliveries, + start_from: factos_pog.Origin, + handle: deliver_managed_subscription_event, + ) + let assert Ok(actor.Started(data: origin_handle, ..)) = + factos_pog.start(origin_subscription) + let assert Ok(ManagedSubscriptionEvent(event: origin_first)) = + process.receive(origin_deliveries, within: 10_000) + let assert Ok(ManagedSubscriptionEvent(event: origin_second)) = + process.receive(origin_deliveries, within: 10_000) + let assert Ok(ManagedSubscriptionEvent(event: origin_third)) = + process.receive(origin_deliveries, within: 10_000) + assert origin_first == recorded_dispatch_event(first_dispatch) + assert origin_second == recorded_dispatch_event(second_dispatch) + assert origin_third == recorded_dispatch_event(third_dispatch) + let assert Ok(Nil) = factos_pog.stop(origin_handle) + + Nil +} + +pub fn subscription_status_and_invalid_reset_preserve_checkpoint_test() { + use config, connection <- with_test_database() + reset_schema(connection) + let deliveries = process.new_subject() + let assert Ok(first_dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "status-first", + decider: decider(), + codec: codec(), + ) + |> factos_pog.dispatch( + RegisterUser(username: "status-first"), + event_id: uuid.v4_string, + ) + let subscription = + test_subscription_from( + connection, + config, + "status-subscription", + deliveries, + start_from: factos_pog.Current, + handle: deliver_managed_subscription_event, + ) + let assert Ok(actor.Started(data: handle, ..)) = + factos_pog.start(subscription) + let assert Ok(Nil) = factos_pog.stop(handle) + + let assert Ok(second_dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "status-second", + decider: decider(), + codec: codec(), + ) + |> factos_pog.dispatch( + RegisterUser(username: "status-second"), + event_id: uuid.v4_string, + ) + let assert Ok(third_dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "status-third", + decider: decider(), + codec: codec(), + ) + |> factos_pog.dispatch( + RegisterUser(username: "status-third"), + event_id: uuid.v4_string, + ) + + let assert Ok(status) = + factos_pog.subscription_status(connection, name: "status-subscription") + let assert factos_pog.SubscriptionStatus( + name: "status-subscription", + cursor: first_cursor, + generation: 0, + event_log_position: third_cursor, + events_behind: 2, + ) = status + assert first_cursor == first_dispatch.append.position + assert third_cursor == third_dispatch.append.position + + let negative = factos.SequencePosition(-2) + let assert Error(factos_pog.InvalidSubscriptionPosition(position:)) = + factos_pog.reset_subscription( + connection, + name: "status-subscription", + start_from: factos_pog.After(position: negative), + ) + assert position == negative + let assert Ok(after_negative) = + factos_pog.subscription_status(connection, name: "status-subscription") + assert after_negative == status + + let assert factos.SequencePosition(head) = third_dispatch.append.position + let ahead = factos.SequencePosition(head + 1) + let assert Error(factos_pog.SubscriptionStartAfterEventLog( + name: error_name, + requested:, + event_log_position:, + )) = + factos_pog.reset_subscription( + connection, + name: "status-subscription", + start_from: factos_pog.After(position: ahead), + ) + assert error_name == "status-subscription" + assert requested == ahead + assert event_log_position == third_dispatch.append.position + let assert Ok(after_ahead) = + factos_pog.subscription_status(connection, name: "status-subscription") + assert after_ahead == status + + let assert Error(factos_pog.SubscriptionNotFound(name: missing_name)) = + factos_pog.reset_subscription( + connection, + name: "missing-subscription", + start_from: factos_pog.Origin, + ) + assert missing_name == "missing-subscription" + let assert Error(factos_pog.InvalidSubscriptionName(name: blank_name)) = + factos_pog.subscription_status(connection, name: " ") + assert blank_name == " " + + let assert Ok(Nil) = + factos_pog.reset_subscription( + connection, + name: "status-subscription", + start_from: factos_pog.Current, + ) + let assert Ok(factos_pog.SubscriptionStatus( + cursor: reset_cursor, + generation: 1, + event_log_position: reset_head, + events_behind: 0, + .., + )) = factos_pog.subscription_status(connection, name: "status-subscription") + assert reset_cursor == third_dispatch.append.position + assert reset_head == third_dispatch.append.position + let _ = second_dispatch + + Nil +} + +pub fn online_reset_replays_origin_after_and_current_ranges_test() { + use config, connection <- with_test_database() + reset_schema(connection) + let deliveries = process.new_subject() + let assert Ok(first_dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "reset-first", + decider: decider(), + codec: codec(), + ) + |> factos_pog.dispatch( + RegisterUser(username: "reset-first"), + event_id: uuid.v4_string, + ) + let assert Ok(second_dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "reset-second", + decider: decider(), + codec: codec(), + ) + |> factos_pog.dispatch( + RegisterUser(username: "reset-second"), + event_id: uuid.v4_string, + ) + let assert Ok(third_dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "reset-third", + decider: decider(), + codec: codec(), + ) + |> factos_pog.dispatch( + RegisterUser(username: "reset-third"), + event_id: uuid.v4_string, + ) + let subscription = + test_subscription_from( + connection, + config, + "online-reset", + deliveries, + start_from: factos_pog.Current, + handle: deliver_managed_subscription_event, + ) + let assert Ok(actor.Started(data: handle, ..)) = + factos_pog.start(subscription) + let assert Error(Nil) = process.receive(deliveries, within: 200) + + let assert Ok(Nil) = + factos_pog.reset_subscription( + connection, + name: "online-reset", + start_from: factos_pog.Origin, + ) + let assert Ok(ManagedSubscriptionEvent(event: origin_first)) = + process.receive(deliveries, within: 10_000) + let assert Ok(ManagedSubscriptionEvent(event: origin_second)) = + process.receive(deliveries, within: 10_000) + let assert Ok(ManagedSubscriptionEvent(event: origin_third)) = + process.receive(deliveries, within: 10_000) + assert origin_first == recorded_dispatch_event(first_dispatch) + assert origin_second == recorded_dispatch_event(second_dispatch) + assert origin_third == recorded_dispatch_event(third_dispatch) + let assert Ok(Nil) = + wait_for_subscription_cursor( + connection, + "online-reset", + third_dispatch.append.position, + 50, + ) + let assert Ok(factos_pog.SubscriptionStatus(generation: 1, ..)) = + factos_pog.subscription_status(connection, name: "online-reset") + + let assert Ok(Nil) = + factos_pog.reset_subscription( + connection, + name: "online-reset", + start_from: factos_pog.After(position: first_dispatch.append.position), + ) + let assert Ok(ManagedSubscriptionEvent(event: after_second)) = + process.receive(deliveries, within: 10_000) + let assert Ok(ManagedSubscriptionEvent(event: after_third)) = + process.receive(deliveries, within: 10_000) + assert after_second == recorded_dispatch_event(second_dispatch) + assert after_third == recorded_dispatch_event(third_dispatch) + let assert Ok(Nil) = + wait_for_subscription_cursor( + connection, + "online-reset", + third_dispatch.append.position, + 50, + ) + let assert Ok(factos_pog.SubscriptionStatus(generation: 2, ..)) = + factos_pog.subscription_status(connection, name: "online-reset") + + let assert Ok(Nil) = + factos_pog.reset_subscription( + connection, + name: "online-reset", + start_from: factos_pog.Current, + ) + let assert Error(Nil) = process.receive(deliveries, within: 200) + let assert Ok(factos_pog.SubscriptionStatus( + cursor: current_cursor, + generation: 3, + events_behind: 0, + .., + )) = factos_pog.subscription_status(connection, name: "online-reset") + assert current_cursor == third_dispatch.append.position + + let assert Ok(fourth_dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "reset-fourth", + decider: decider(), + codec: codec(), + ) + |> factos_pog.dispatch( + RegisterUser(username: "reset-fourth"), + event_id: uuid.v4_string, + ) + let assert Ok(ManagedSubscriptionEvent(event: fourth_event)) = + process.receive(deliveries, within: 10_000) + assert fourth_event == recorded_dispatch_event(fourth_dispatch) + let assert Ok(Nil) = factos_pog.stop(handle) + + Nil +} + +pub fn stale_generation_checkpoint_cannot_overwrite_online_reset_test() { + use config, connection <- with_test_database() + reset_schema(connection) + let handler_started = process.new_subject() + let assert Ok(subscription) = + factos_pog.new_subscription( + connection:, + config: pog.Config( + ..config, + pool_name: process.new_name(prefix: "stale-generation-pool"), + ), + name: "stale-generation", + query: factos.AllEvents, + start_from: factos_pog.Current, + codec: codec(), + handle: fn(recorded) { + process.send(handler_started, BlockingHandlerStarted(event: recorded)) + process.sleep(500) + Ok(Nil) + }, + ) + let assert Ok(actor.Started(data: handle, ..)) = + factos_pog.start(subscription) + let assert Ok(dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "stale-generation", + decider: decider(), + codec: codec(), + ) + |> factos_pog.dispatch( + RegisterUser(username: "stale-generation"), + event_id: uuid.v4_string, + ) + let assert Ok(BlockingHandlerStarted(event: first_attempt)) = + process.receive(handler_started, within: 10_000) + assert first_attempt == recorded_dispatch_event(dispatch) + + let assert Ok(Nil) = + factos_pog.reset_subscription( + connection, + name: "stale-generation", + start_from: factos_pog.Origin, + ) + let assert Ok(BlockingHandlerStarted(event: second_attempt)) = + process.receive(handler_started, within: 10_000) + assert second_attempt == first_attempt + + // The stale generation completed its callback but could not advance the row. + let assert Ok(factos_pog.SubscriptionStatus( + cursor: factos.NoPosition, + generation: 1, + events_behind: 1, + .., + )) = factos_pog.subscription_status(connection, name: "stale-generation") + + let assert Ok(Nil) = + wait_for_subscription_cursor( + connection, + "stale-generation", + dispatch.append.position, + 50, + ) + let assert Ok(factos_pog.SubscriptionStatus( + cursor: final_cursor, + generation: 1, + events_behind: 0, + .., + )) = factos_pog.subscription_status(connection, name: "stale-generation") + assert final_cursor == dispatch.append.position + let assert Ok(Nil) = factos_pog.stop(handle) + + Nil +} + +pub fn projection_subscription_rolls_back_failed_attempt_and_checkpoint_test() { + use config, connection <- with_test_database() + reset_schema(connection) + reset_managed_subscription_state(connection) + let assert Ok(dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "atomic-projection", + decider: decider(), + codec: codec(), + ) + |> factos_pog.dispatch( + RegisterUser(username: "atomic-projection"), + event_id: uuid.v4_string, + ) + let assert Ok(subscription) = + factos_pog.new_projection_subscription( + connection:, + config: pog.Config( + ..config, + pool_name: process.new_name(prefix: "atomic-projection-pool"), + ), + name: "atomic-projection", + query: factos.AllEvents, + start_from: factos_pog.Origin, + codec: codec(), + project: fn(transaction_connection, recorded) { + use attempt <- result.try(increment_managed_subscription_attempts( + connection, + )) + use _ <- result.try(insert_test_projection( + transaction_connection, + recorded, + )) + case attempt { + 1 -> Error("fail the first projection attempt") + _ -> Ok(Nil) + } + }, + ) + let assert Ok(actor.Started(data: handle, ..)) = + factos_pog.start(subscription) + let assert Ok(Nil) = + wait_for_subscription_cursor( + connection, + "atomic-projection", + dispatch.append.position, + 50, + ) + + let factos.Recorded(id: event_id, event: UserRegistered(username:), ..) = + recorded_dispatch_event(dispatch) + assert projected_users(connection) == [#(event_id, username)] + assert managed_subscription_attempts(connection) == 2 + let assert Ok(factos_pog.SubscriptionStatus( + cursor: cursor, + generation: 0, + events_behind: 0, + .., + )) = factos_pog.subscription_status(connection, name: "atomic-projection") + assert cursor == dispatch.append.position + let assert Ok(Nil) = factos_pog.stop(handle) + + Nil +} + +pub fn projection_reset_race_rolls_back_stale_generation_test() { + use config, connection <- with_test_database() + reset_schema(connection) + reset_managed_subscription_state(connection) + let attempts = process.new_subject() + let assert Ok(subscription) = + factos_pog.new_projection_subscription( + connection:, + config: pog.Config( + ..config, + pool_name: process.new_name(prefix: "projection-reset-pool"), + ), + name: "projection-reset", + query: factos.AllEvents, + start_from: factos_pog.Current, + codec: codec(), + project: fn(transaction_connection, recorded) { + use attempt <- result.try(increment_managed_subscription_attempts( + connection, + )) + use _ <- result.try(insert_test_projection( + transaction_connection, + recorded, + )) + process.send( + attempts, + ProjectionAttemptStarted(attempt:, event: recorded), + ) + process.sleep(500) + Ok(Nil) + }, + ) + let assert Ok(actor.Started(data: handle, ..)) = + factos_pog.start(subscription) + let assert Ok(dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "projection-reset", + decider: decider(), + codec: codec(), + ) + |> factos_pog.dispatch( + RegisterUser(username: "projection-reset"), + event_id: uuid.v4_string, + ) + let assert Ok(ProjectionAttemptStarted(attempt: 1, event: first_attempt)) = + process.receive(attempts, within: 10_000) + assert first_attempt == recorded_dispatch_event(dispatch) + + let assert Ok(Nil) = + factos_pog.reset_subscription( + connection, + name: "projection-reset", + start_from: factos_pog.Origin, + ) + let assert Ok(ProjectionAttemptStarted(attempt: 2, event: second_attempt)) = + process.receive(attempts, within: 10_000) + assert second_attempt == first_attempt + + // The first transaction was rolled back, and the retry is not visible before + // its generation-checked cursor update commits. + assert projected_users(connection) == [] + let assert Ok(factos_pog.SubscriptionStatus( + cursor: factos.NoPosition, + generation: 1, + events_behind: 1, + .., + )) = factos_pog.subscription_status(connection, name: "projection-reset") + + let assert Ok(Nil) = + wait_for_subscription_cursor( + connection, + "projection-reset", + dispatch.append.position, + 50, + ) + let factos.Recorded(id: event_id, event: UserRegistered(username:), ..) = + recorded_dispatch_event(dispatch) + assert projected_users(connection) == [#(event_id, username)] + assert managed_subscription_attempts(connection) == 2 + let assert Ok(factos_pog.SubscriptionStatus( + cursor: final_cursor, + generation: 1, + events_behind: 0, + .., + )) = factos_pog.subscription_status(connection, name: "projection-reset") + assert final_cursor == dispatch.append.position + let assert Ok(Nil) = factos_pog.stop(handle) + + Nil +} + +pub fn dispatch_and_wait_returns_after_projection_commit_test() { + use config, connection <- with_test_database() + reset_schema(connection) + reset_managed_subscription_state(connection) + let barriers = process.new_subject() + let results = process.new_subject() + let assert Ok(subscription) = + factos_pog.new_projection_subscription( + connection:, + config: pog.Config( + ..config, + pool_name: process.new_name(prefix: "dispatch-wait-projection-pool"), + ), + name: "dispatch-wait-projection", + query: factos.AllEvents, + start_from: factos_pog.Current, + codec: codec(), + project: fn(transaction_connection, recorded) { + use _ <- result.try(insert_test_projection( + transaction_connection, + recorded, + )) + block_projection("dispatch-wait-projection", barriers, recorded) + }, + ) + let assert Ok(actor.Started(data: handle, ..)) = + factos_pog.start(subscription) + let _worker = + start_dispatch_wait_worker( + connection, + stream: "dispatch-wait-projection", + username: "dispatch-wait-projection", + subscriptions: ["dispatch-wait-projection"], + timeout: duration.seconds(5), + results:, + ) + + let assert Ok(ProjectionBarrierReady(event: projected_event, release:, ..)) = + process.receive(barriers, within: 10_000) + let assert Error(Nil) = process.receive(results, within: 200) + assert projected_users(connection) == [] + + process.send(release, Nil) + let assert Ok(DispatchWaitCompleted(result: Ok(dispatch))) = + process.receive(results, within: 10_000) + assert recorded_dispatch_event(dispatch) == projected_event + let factos.Recorded(id:, event: UserRegistered(username:), ..) = + projected_event + assert projected_users(connection) == [#(id, username)] + let assert Ok(factos_pog.SubscriptionStatus( + cursor: cursor, + events_behind: 0, + .., + )) = + factos_pog.subscription_status(connection, name: "dispatch-wait-projection") + assert cursor == dispatch.append.position + let assert Ok(Nil) = factos_pog.stop(handle) + + Nil +} + +pub fn dispatch_and_wait_waits_for_every_unique_subscription_test() { + use config, connection <- with_test_database() + reset_schema(connection) + let barriers = process.new_subject() + let results = process.new_subject() + let first_subscription = + blocking_projection_subscription(connection, config, "wait-first", barriers) + let second_subscription = + blocking_projection_subscription( + connection, + config, + "wait-second", + barriers, + ) + let assert Ok(actor.Started(data: first_handle, ..)) = + factos_pog.start(first_subscription) + let assert Ok(actor.Started(data: second_handle, ..)) = + factos_pog.start(second_subscription) + let selected = ["wait-second", "wait-first", "wait-second"] + let _worker = + start_dispatch_wait_worker( + connection, + stream: "wait-every-subscription", + username: "wait-every-subscription", + subscriptions: selected, + timeout: duration.seconds(5), + results:, + ) + let assert Ok(first_ready) = process.receive(barriers, within: 10_000) + let assert Ok(second_ready) = process.receive(barriers, within: 10_000) + let ProjectionBarrierReady(event: first_event, release: first_release, ..) = + first_ready + let ProjectionBarrierReady(event: second_event, release: second_release, ..) = + second_ready + assert first_event == second_event + + let assert Error(factos_pog.SubscriptionWaitTimedOut( + through:, + pending: [first_pending, second_pending], + )) = + factos_pog.wait_for_subscriptions( + connection, + subscriptions: selected, + through: first_event.position, + timeout: duration.nanoseconds(500_000), + ) + assert through == first_event.position + let factos_pog.SubscriptionStatus(name: first_pending_name, ..) = + first_pending + let factos_pog.SubscriptionStatus(name: second_pending_name, ..) = + second_pending + assert first_pending_name == "wait-second" + assert second_pending_name == "wait-first" + let assert Error(factos_pog.SubscriptionWaitTimedOut(..)) = + factos_pog.wait_for_subscriptions( + connection, + subscriptions: selected, + through: first_event.position, + timeout: duration.milliseconds(-1), + ) + + process.send(first_release, Nil) + let assert Error(Nil) = process.receive(results, within: 200) + process.send(second_release, Nil) + let assert Ok(DispatchWaitCompleted(result: Ok(dispatch))) = + process.receive(results, within: 10_000) + assert recorded_dispatch_event(dispatch) == first_event + + let assert Ok(Nil) = factos_pog.stop(first_handle) + let assert Ok(Nil) = factos_pog.stop(second_handle) + Nil +} + +pub fn dispatch_and_wait_timeout_carries_committed_dispatch_test() { + use config, connection <- with_test_database() + reset_schema(connection) + let barriers = process.new_subject() + let subscription = + blocking_projection_subscription( + connection, + config, + "timeout-projection", + barriers, + ) + let assert Ok(actor.Started(data: handle, ..)) = + factos_pog.start(subscription) + + let result = + factos_pog.new_dispatch( + connection:, + stream: "timeout-projection", + decider: decider(), + codec: codec(), + ) + |> factos_pog.dispatch_and_wait( + RegisterUser(username: "timeout-projection"), + event_id: uuid.v4_string, + subscriptions: ["timeout-projection"], + timeout: duration.milliseconds(50), + ) + let assert Error(factos_pog.DispatchCommittedButNotObserved( + dispatch:, + reason: factos_pog.SubscriptionWaitTimedOut(through:, pending: [pending]), + )) = result + assert through == dispatch.append.position + let assert factos_pog.SubscriptionStatus( + name: pending_name, + cursor: factos.NoPosition, + events_behind: 1, + .., + ) = pending + assert pending_name == "timeout-projection" + let assert Ok(events) = + factos_pog.read_after( + connection, + query: factos.AllEvents, + after: factos.NoPosition, + limit: 10, + codec: codec(), + ) + assert events == dispatch.events - let assert Ok(dispatch) = + let assert Ok(ProjectionBarrierReady(release:, ..)) = + process.receive(barriers, within: 10_000) + process.send(release, Nil) + let assert Ok(Nil) = + wait_for_subscription_cursor( + connection, + "timeout-projection", + dispatch.append.position, + 50, + ) + let assert Ok(Nil) = factos_pog.stop(handle) + Nil +} + +pub fn dispatch_and_wait_missing_subscription_carries_commit_test() { + use connection <- with_test_connection() + reset_schema(connection) + let result = factos_pog.new_dispatch( connection:, - stream: "user-renata", + stream: "missing-wait-subscription", decider: decider(), codec: codec(), ) - |> factos_pog.dispatch(RegisterUser("renata"), event_id: uuid.v4_string) - - let assert Ok(pog.Notify( - pid: notification_pid, - reference: notification_reference, - channel: "factos_pog_events", - payload:, - )) = process.selector_receive(selector, 5000) - assert notification_pid == notifications_pid - assert notification_reference == listener - let assert Ok(#( - cursor, - previous_cursor, - type_, - version, - tags, - correlation_id, - username, - )) = json.parse(payload, using: event_notification_payload_decoder()) - let assert factos.SequencePosition(position) = dispatch.append.position - assert cursor == position - assert previous_cursor == -1 - assert type_ == "UserRegistered" - assert version == 1 - assert tags == ["username:renata"] - assert correlation_id == "event-renata" - assert username == "renata" - + |> factos_pog.dispatch_and_wait( + RegisterUser(username: "missing-wait-subscription"), + event_id: uuid.v4_string, + subscriptions: ["missing-subscription"], + timeout: duration.seconds(5), + ) + let assert Error(factos_pog.DispatchCommittedButNotObserved( + dispatch:, + reason: factos_pog.SubscriptionNotFound(name: missing_name), + )) = result + assert missing_name == "missing-subscription" let assert Ok(events) = factos_pog.read_after( connection, @@ -198,46 +1303,149 @@ pub fn event_notifications_wake_durable_catch_up_test() { codec: codec(), ) assert events == dispatch.events - assert factos.react_all(welcome_reactor(), events) == [SendWelcome("renata")] +} - let assert [last] = list.reverse(events) - let assert Ok([]) = +pub fn dispatch_and_wait_wraps_precommit_failure_test() { + use connection <- with_test_connection() + reset_schema(connection) + let assert Ok(first_dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "precommit-failure", + decider: decider(), + codec: codec(), + ) + |> factos_pog.dispatch( + RegisterUser(username: "already-registered"), + event_id: uuid.v4_string, + ) + let result = + factos_pog.new_dispatch( + connection:, + stream: "precommit-failure", + decider: decider(), + codec: codec(), + ) + |> factos_pog.dispatch_and_wait( + RegisterUser(username: "already-registered"), + event_id: uuid.v4_string, + subscriptions: ["missing-subscription"], + timeout: duration.seconds(5), + ) + let assert Error(factos_pog.DispatchNotCommitted(error: factos_pog.DomainError( + AlreadyTaken, + ))) = result + let assert Ok(events) = factos_pog.read_after( connection, query: factos.AllEvents, - after: last.position, + after: factos.NoPosition, limit: 10, codec: codec(), ) + assert events == first_dispatch.events +} - // NOTIFY follows the transaction: rolling the event back exposes neither - // a durable event nor a wake-up edge. - let assert Error(pog.TransactionRolledBack(Nil)) = - pog.transaction(connection, fn(transaction_connection) { - let assert Ok(_) = - pog.query( - " - insert into factos_events ( - id, stream, revision, type, version, tags, metadata, data - ) - values ( - 'b3b12f1d-6d85-4f1f-9c2a-94766a34f099', - 'rolled-back', - 0, - 'UserRegistered', - 1, - '', - '{}'::jsonb, - '\"rolled-back\"'::jsonb - ) - ", - ) - |> pog.execute(on: transaction_connection) - Error(Nil) - }) - let assert Error(Nil) = process.selector_receive(selector, 200) +pub fn dispatch_and_wait_bypasses_lookup_without_barrier_test() { + use connection <- with_test_connection() + reset_schema(connection) + let assert Ok(no_event_dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "no-event-wait", + decider: counter_decider(), + codec: counter_codec(), + ) + |> factos_pog.dispatch_and_wait( + DoNothing, + event_id: uuid.v4_string, + subscriptions: ["missing-subscription"], + timeout: duration.seconds(5), + ) + assert no_event_dispatch.append.position == factos.NoPosition + assert no_event_dispatch.events == [] + + let assert Ok(event_dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "empty-subscription-wait", + decider: counter_decider(), + codec: counter_codec(), + ) + |> factos_pog.dispatch_and_wait( + Increment, + event_id: uuid.v4_string, + subscriptions: [], + timeout: duration.seconds(5), + ) + let assert factos.SequencePosition(_) = event_dispatch.append.position + let assert [_] = event_dispatch.events + + let assert Error(factos_pog.InvalidSubscriptionName(name: blank_name)) = + factos_pog.wait_for_subscriptions( + connection, + subscriptions: [" "], + through: event_dispatch.append.position, + timeout: duration.milliseconds(0), + ) + assert blank_name == " " + let negative = factos.SequencePosition(-2) + let assert Error(factos_pog.InvalidSubscriptionPosition(position:)) = + factos_pog.wait_for_subscriptions( + connection, + subscriptions: ["missing-subscription"], + through: negative, + timeout: duration.milliseconds(0), + ) + assert position == negative + Nil +} + +pub fn invalid_initial_start_leaves_no_durable_checkpoint_test() { + use config, connection <- with_test_database() + reset_schema(connection) + let deliveries = process.new_subject() + let assert Ok(dispatch) = + factos_pog.new_dispatch( + connection:, + stream: "invalid-start-head", + decider: decider(), + codec: codec(), + ) + |> factos_pog.dispatch( + RegisterUser(username: "invalid-start-head"), + event_id: uuid.v4_string, + ) + + let negative_subscription = + test_subscription_from( + connection, + config, + "negative-initial-start", + deliveries, + start_from: factos_pog.After(position: factos.SequencePosition(-2)), + handle: deliver_managed_subscription_event, + ) + wait_for_subscription_start_failure(negative_subscription) + let assert Error(factos_pog.SubscriptionNotFound(name: negative_name)) = + factos_pog.subscription_status(connection, name: "negative-initial-start") + assert negative_name == "negative-initial-start" + + let assert factos.SequencePosition(head) = dispatch.append.position + let ahead_subscription = + test_subscription_from( + connection, + config, + "ahead-initial-start", + deliveries, + start_from: factos_pog.After(position: factos.SequencePosition(head + 1)), + handle: deliver_managed_subscription_event, + ) + wait_for_subscription_start_failure(ahead_subscription) + let assert Error(factos_pog.SubscriptionNotFound(name: ahead_name)) = + factos_pog.subscription_status(connection, name: "ahead-initial-start") + assert ahead_name == "ahead-initial-start" - factos_pog.unlisten(notifications, listener) Nil } @@ -271,7 +1479,7 @@ pub fn managed_subscription_starts_catches_up_and_stops_test() { let assert Ok(ManagedSubscriptionEvent(event: first_event)) = process.receive(deliveries, within: 10_000) - assert first_event == decoded_dispatch_event(first_dispatch) + assert first_event == recorded_dispatch_event(first_dispatch) let assert Ok(Nil) = wait_for_subscription_cursor( connection, @@ -293,7 +1501,7 @@ pub fn managed_subscription_starts_catches_up_and_stops_test() { ) let assert Ok(ManagedSubscriptionEvent(event: second_event)) = process.receive(deliveries, within: 10_000) - assert second_event == decoded_dispatch_event(second_dispatch) + assert second_event == recorded_dispatch_event(second_dispatch) let assert Ok(Nil) = wait_for_subscription_cursor( connection, @@ -328,7 +1536,7 @@ pub fn managed_subscription_starts_catches_up_and_stops_test() { |> pog.returning(int_column_decoder()) |> pog.execute(on: connection) let assert [oversized_cursor] = returned.rows - let assert Ok(ManagedSubscriptionEvent(event: factos.Decoded( + let assert Ok(ManagedSubscriptionEvent(event: factos.Recorded( event: UserRegistered(username:), .., ))) = process.receive(deliveries, within: 10_000) @@ -383,10 +1591,11 @@ pub fn managed_subscription_checkpoints_each_successful_event_test() { ), name:, query: factos.AllEvents, + start_from: factos_pog.Origin, codec: codec(), handle: fn(event) { process.send(deliveries, ManagedSubscriptionEvent(event:)) - let factos.Decoded(event: UserRegistered(username:), ..) = event + let factos.Recorded(event: UserRegistered(username:), ..) = event case username == "lucy" { True -> Error("expected retry") False -> Ok(Nil) @@ -400,8 +1609,8 @@ pub fn managed_subscription_checkpoints_each_successful_event_test() { process.receive(deliveries, within: 10_000) let assert Ok(ManagedSubscriptionEvent(event: second_delivery)) = process.receive(deliveries, within: 10_000) - assert first_delivery == decoded_dispatch_event(first_dispatch) - assert second_delivery == decoded_dispatch_event(second_dispatch) + assert first_delivery == recorded_dispatch_event(first_dispatch) + assert second_delivery == recorded_dispatch_event(second_dispatch) let assert Ok(cursor) = load_subscription_cursor(connection, name) assert cursor == first_dispatch.append.position @@ -448,7 +1657,7 @@ pub fn supervised_subscription_restarts_from_durable_cursor_test() { let assert Ok(ManagedSubscriptionEvent(event:)) = process.receive(deliveries, within: 10_000) - assert event == decoded_dispatch_event(dispatch) + assert event == recorded_dispatch_event(dispatch) let assert Ok(Nil) = wait_for_subscription_cursor(connection, name, dispatch.append.position, 50) assert managed_subscription_attempts(connection) == 2 @@ -1039,39 +2248,104 @@ fn is_already_taken_dispatch( } } -fn with_test_connection( - body: fn(pog.Connection) -> Nil, -) -> Result(Nil, testcontainer_error.Error) { - with_test_database(fn(_config, connection) { body(connection) }) +fn shared_postgres() -> SharedPostgres { + global_value.create_with_unique_name( + "factos_pog_test.shared_postgres", + start_shared_postgres, + ) } -fn with_test_database( - body: fn(pog.Config, pog.Connection) -> Nil, -) -> Result(Nil, testcontainer_error.Error) { - use postgres_container <- testcontainer.with_formula( - postgres.new() |> postgres.formula(), +fn start_shared_postgres() -> SharedPostgres { + let host = environment_variable("FACTOS_POG_TEST_HOST", default: "127.0.0.1") + let port = test_postgres_port() + let database = + environment_variable("FACTOS_POG_TEST_DATABASE", default: "factos_pog") + let username = + environment_variable("FACTOS_POG_TEST_USERNAME", default: "postgres") + let password = + environment_variable("FACTOS_POG_TEST_PASSWORD", default: "postgres") + let #(admin_pool_pid, _, admin_connection) = + start_test_connection(host:, port:, database:, username:, password:) + SharedPostgres( + admin_pool_pid:, + admin_connection:, + host:, + port:, + username:, + password:, ) +} + +fn environment_variable(name: String, default default_value: String) -> String { + envoy.get(name) + |> result.unwrap(default_value) +} + +fn test_postgres_port() -> Int { + let value = environment_variable("FACTOS_POG_TEST_PORT", default: "55432") + let assert Ok(port) = int.parse(value) + port +} + +fn with_test_connection(body: fn(pog.Connection) -> Nil) -> Nil { + with_test_database(fn(_config, connection) { body(connection) }) +} + +fn with_test_database(body: fn(pog.Config, pog.Connection) -> Nil) -> Nil { + let SharedPostgres(admin_connection:, host:, port:, username:, password:, ..) = + shared_postgres() + let database_suffix = uuid.v4_string() |> string.replace("-", "_") + let database_name = "factos_pog_test_" <> database_suffix + create_test_database(admin_connection, database_name) let #(pool_pid, config, connection) = - start_test_connection(postgres_container) + start_test_connection( + host:, + port:, + database: database_name, + username:, + password:, + ) body(config, connection) process.send_exit(pool_pid) - process.sleep(100) - Ok(Nil) + drop_test_database(admin_connection, database_name) +} + +fn create_test_database( + admin_connection: pog.Connection, + database_name: String, +) -> Nil { + // The UUID-derived name contains only ASCII letters, digits, and underscores. + let assert Ok(_) = + pog.query("create database " <> database_name) + |> pog.execute(on: admin_connection) + Nil +} + +fn drop_test_database( + admin_connection: pog.Connection, + database_name: String, +) -> Nil { + let assert Ok(_) = + pog.query("drop database if exists " <> database_name <> " with (force)") + |> pog.execute(on: admin_connection) + Nil } fn start_test_connection( - postgres_container: postgres.PostgresContainer, + host host_name: String, + port port_number: Int, + database database_name: String, + username username: String, + password password: String, ) -> #(process.Pid, pog.Config, pog.Connection) { - let postgres.PostgresContainer(host:, port:, database:, username:, ..) = - postgres_container let pool_name = process.new_name("factos_pog_test") let config = pog.default_config(pool_name) - |> pog.host(host) - |> pog.port(port) - |> pog.database(database) + |> pog.host(host_name) + |> pog.port(port_number) + |> pog.database(database_name) |> pog.user(username) - |> pog.password(Some("postgres")) + |> pog.password(Some(password)) |> pog.ssl(pog.SslDisabled) let assert Ok(actor.Started(pid:, ..)) = pog.start(config) @@ -1121,6 +2395,23 @@ fn drop_schema(connection: pog.Connection) -> Nil { Nil } +fn wait_for_subscription_start_failure( + subscription: factos_pog.Subscription(Event, String), +) -> Nil { + let subscription_process = + process.spawn_unlinked(fn() { + let _ = factos_pog.start(subscription) + Nil + }) + let monitor = process.monitor(subscription_process) + let assert Ok(Nil) = + process.new_selector() + |> process.select_specific_monitor(monitor, fn(_) { Nil }) + |> process.selector_receive(5000) + process.demonitor_process(monitor) + Nil +} + fn reset_managed_subscription_state(connection: pog.Connection) -> Nil { let assert Ok(_) = pog.query("drop table if exists factos_pog_test_subscription") @@ -1143,6 +2434,19 @@ fn reset_managed_subscription_state(connection: pog.Connection) -> Nil { ", ) |> pog.execute(on: connection) + let assert Ok(_) = + pog.query("drop table if exists factos_pog_test_projection") + |> pog.execute(on: connection) + let assert Ok(_) = + pog.query( + " + create table factos_pog_test_projection ( + event_id text primary key, + username text not null + ) + ", + ) + |> pog.execute(on: connection) Nil } @@ -1154,7 +2458,29 @@ fn test_subscription( handle handle: fn( pog.Connection, process.Subject(ManagedSubscriptionMessage), - factos.Decoded(Event), + factos.Recorded(Event), + ) -> Result(Nil, String), +) -> factos_pog.Subscription(Event, String) { + test_subscription_from( + connection, + config, + name, + deliveries, + start_from: factos_pog.Origin, + handle:, + ) +} + +fn test_subscription_from( + connection: pog.Connection, + config: pog.Config, + name: String, + deliveries: process.Subject(ManagedSubscriptionMessage), + start_from start_from: factos_pog.SubscriptionStart, + handle handle: fn( + pog.Connection, + process.Subject(ManagedSubscriptionMessage), + factos.Recorded(Event), ) -> Result(Nil, String), ) -> factos_pog.Subscription(Event, String) { let assert Ok(subscription) = @@ -1166,17 +2492,18 @@ fn test_subscription( ), name:, query: factos.AllEvents, + start_from:, codec: codec(), handle: fn(event) { handle(connection, deliveries, event) }, ) subscription } -fn decoded_dispatch_event( +fn recorded_dispatch_event( dispatch: factos_pog.Dispatch(Event), -) -> factos.Decoded(Event) { - let assert [factos.Recorded(event:, descriptor:, ..)] = dispatch.events - factos.Decoded(event:, descriptor:) +) -> factos.Recorded(Event) { + let assert [recorded] = dispatch.events + recorded } fn load_subscription_cursor( @@ -1240,7 +2567,7 @@ fn wait_for_subscription_cursor( fn deliver_managed_subscription_event( _connection: pog.Connection, deliveries: process.Subject(ManagedSubscriptionMessage), - event: factos.Decoded(Event), + event: factos.Recorded(Event), ) -> Result(Nil, String) { process.send(deliveries, ManagedSubscriptionEvent(event:)) Ok(Nil) @@ -1249,8 +2576,18 @@ fn deliver_managed_subscription_event( fn retry_managed_subscription_event_once( connection: pog.Connection, deliveries: process.Subject(ManagedSubscriptionMessage), - event: factos.Decoded(Event), + event: factos.Recorded(Event), ) -> Result(Nil, String) { + use attempts <- result.try(increment_managed_subscription_attempts(connection)) + case attempts { + 1 -> Error("retry managed subscription event once") + _ -> deliver_managed_subscription_event(connection, deliveries, event) + } +} + +fn increment_managed_subscription_attempts( + connection: pog.Connection, +) -> Result(Int, String) { case pog.query( " @@ -1266,14 +2603,117 @@ fn retry_managed_subscription_event_once( Error(error) -> Error(string.inspect(error)) Ok(returned) -> case returned.rows { - [1] -> Error("retry managed subscription event once") - [_] -> deliver_managed_subscription_event(connection, deliveries, event) + [attempts] -> Ok(attempts) [] -> Error("managed subscription attempt row is missing") [_, _, ..] -> Error("managed subscription attempt row is duplicated") } } } +fn block_projection( + name: String, + barriers: process.Subject(ProjectionBarrierMessage), + recorded: factos.Recorded(Event), +) -> Result(Nil, String) { + let release = process.new_subject() + process.send( + barriers, + ProjectionBarrierReady(name:, event: recorded, release:), + ) + case process.receive(release, within: 10_000) { + Ok(Nil) -> Ok(Nil) + Error(Nil) -> Error("projection release timed out") + } +} + +fn blocking_projection_subscription( + connection: pog.Connection, + config: pog.Config, + name: String, + barriers: process.Subject(ProjectionBarrierMessage), +) -> factos_pog.Subscription(Event, String) { + let assert Ok(subscription) = + factos_pog.new_projection_subscription( + connection:, + config: pog.Config( + ..config, + pool_name: process.new_name(prefix: "blocking-projection-pool"), + ), + name:, + query: factos.AllEvents, + start_from: factos_pog.Current, + codec: codec(), + project: fn(_transaction_connection, recorded) { + block_projection(name, barriers, recorded) + }, + ) + subscription +} + +fn start_dispatch_wait_worker( + connection: pog.Connection, + stream stream_name: String, + username username: String, + subscriptions subscriptions: List(String), + timeout timeout: duration.Duration, + results results: process.Subject(DispatchWaitMessage), +) -> process.Pid { + process.spawn(fn() { + let result = + factos_pog.new_dispatch( + connection:, + stream: stream_name, + decider: decider(), + codec: codec(), + ) + |> factos_pog.dispatch_and_wait( + RegisterUser(username:), + event_id: uuid.v4_string, + subscriptions:, + timeout:, + ) + process.send(results, DispatchWaitCompleted(result:)) + }) +} + +fn insert_test_projection( + connection: pog.Connection, + recorded: factos.Recorded(Event), +) -> Result(Nil, String) { + let factos.Recorded(id:, event: UserRegistered(username:), ..) = recorded + pog.query( + " + insert into factos_pog_test_projection (event_id, username) + values ($1, $2) + ", + ) + |> pog.parameter(pog.text(id)) + |> pog.parameter(pog.text(username)) + |> pog.execute(on: connection) + |> result.map(fn(_) { Nil }) + |> result.map_error(string.inspect) +} + +fn projected_users(connection: pog.Connection) -> List(#(String, String)) { + let assert Ok(returned) = + pog.query( + " + select event_id, username + from factos_pog_test_projection + order by event_id + ", + ) + |> pog.returning(string_pair_decoder()) + |> pog.execute(on: connection) + returned.rows +} + +fn string_pair_decoder() -> decode.Decoder(#(String, String)) { + use first <- decode.field(0, decode.string) + use second <- decode.field(1, decode.string) + decode.success(#(first, second)) +} + fn managed_subscription_attempts(connection: pog.Connection) -> Int { let assert Ok(returned) = pog.query( @@ -1573,11 +3013,56 @@ fn username_conformance_query() -> factos.Query { ]) } +fn assert_event_notification_matches( + payload: String, + dispatch: factos_pog.Dispatch(Event), +) -> EventNotificationPayload { + let assert Ok(notification) = + json.parse(payload, using: event_notification_payload_decoder()) + let EventNotificationPayload( + cursor:, + id:, + stream:, + revision:, + type_:, + version:, + tags:, + correlation_id:, + username:, + .., + ) = notification + let assert [ + factos.Recorded( + id: recorded_id, + stream: recorded_stream, + revision: recorded_revision, + position: factos.SequencePosition(position), + event: UserRegistered(username: recorded_username), + descriptor:, + ), + ] = dispatch.events + assert cursor == position + assert id == recorded_id + assert stream == recorded_stream + assert revision == recorded_revision + assert type_ == factos.event_type_name(descriptor.type_) + assert version == descriptor.version + assert tags == list.map(descriptor.tags, factos.tag_value) + let assert Ok(recorded_correlation_id) = + factos.metadata_get(descriptor.metadata, factos.correlation_id) + assert correlation_id == recorded_correlation_id + assert username == recorded_username + notification +} + fn event_notification_payload_decoder() -> decode.Decoder( - #(Int, Int, String, Int, List(String), String, String), + EventNotificationPayload, ) { use cursor <- decode.field("cursor", decode.int) use previous_cursor <- decode.field("previous_cursor", decode.int) + use id <- decode.subfield(["event", "id"], decode.string) + use stream <- decode.subfield(["event", "stream"], decode.string) + use revision <- decode.subfield(["event", "revision"], decode.int) use type_ <- decode.subfield(["event", "type"], decode.string) use version <- decode.subfield(["event", "version"], decode.int) use tags <- decode.subfield(["event", "tags"], decode.string |> decode.list) @@ -1586,14 +3071,17 @@ fn event_notification_payload_decoder() -> decode.Decoder( decode.string, ) use username <- decode.subfield(["event", "data"], decode.string) - decode.success(#( - cursor, - previous_cursor, - type_, - version, - tags, - correlation_id, - username, + decode.success(EventNotificationPayload( + cursor:, + previous_cursor:, + id:, + stream:, + revision:, + type_:, + version:, + tags:, + correlation_id:, + username:, )) } @@ -1694,6 +3182,7 @@ fn counter_decide( let CounterState(total) = state case command { Increment -> Ok([Incremented(total + 1)]) + DoNothing -> Ok([]) } } -- 2.51.2