From 0fbf019899dc4e884801a3ac455f985d0f4d5ca3 Mon Sep 17 00:00:00 2001 From: Renatillas Date: Fri, 26 Jun 2026 05:44:32 +0200 Subject: [PATCH] Postgres backend and multiple examples --- README.md | 479 ++++++---- backends/factos_kurrentdb_erlang/gleam.toml | 6 + ...db.gleam => factos_kurrentdb_erlang.gleam} | 125 ++- .../test/factos_kurrentdb_erlang_test.gleam | 28 +- backends/factos_pog/README.md | 24 + backends/factos_pog/compose.yml | 14 + backends/factos_pog/gleam.toml | 18 + backends/factos_pog/manifest.toml | 31 + .../factos_pog/src/factos/factos_pog.gleam | 619 +++++++++++++ .../factos_pog/test/factos_pog_test.gleam | 335 +++++++ backends/factos_sqlight/gleam.toml | 6 + .../{sqlight.gleam => factos_sqlight.gleam} | 93 +- .../test/factos_sqlight_test.gleam | 31 +- examples/orders_sqlight/gleam.toml | 13 + examples/orders_sqlight/manifest.toml | 28 + .../orders_sqlight/src/order_workflow.gleam | 866 ++++++++++++++++++ .../orders_sqlight/src/orders_sqlight.gleam | 5 + .../test/orders_sqlight_test.gleam | 16 + examples/tickets_pog/compose.yml | 14 + examples/tickets_pog/gleam.toml | 13 + examples/tickets_pog/manifest.toml | 33 + examples/tickets_pog/src/ticket_sale.gleam | 317 +++++++ examples/tickets_pog/src/tickets_pog.gleam | 5 + .../tickets_pog/test/tickets_pog_test.gleam | 15 + gleam.toml | 5 + src/factos.gleam | 141 ++- 26 files changed, 3041 insertions(+), 239 deletions(-) rename backends/factos_kurrentdb_erlang/src/factos/{kurrentdb.gleam => factos_kurrentdb_erlang.gleam} (77%) create mode 100644 backends/factos_pog/README.md create mode 100644 backends/factos_pog/compose.yml create mode 100644 backends/factos_pog/gleam.toml create mode 100644 backends/factos_pog/manifest.toml create mode 100644 backends/factos_pog/src/factos/factos_pog.gleam create mode 100644 backends/factos_pog/test/factos_pog_test.gleam rename backends/factos_sqlight/src/factos/{sqlight.gleam => factos_sqlight.gleam} (77%) create mode 100644 examples/orders_sqlight/gleam.toml create mode 100644 examples/orders_sqlight/manifest.toml create mode 100644 examples/orders_sqlight/src/order_workflow.gleam create mode 100644 examples/orders_sqlight/src/orders_sqlight.gleam create mode 100644 examples/orders_sqlight/test/orders_sqlight_test.gleam create mode 100644 examples/tickets_pog/compose.yml create mode 100644 examples/tickets_pog/gleam.toml create mode 100644 examples/tickets_pog/manifest.toml create mode 100644 examples/tickets_pog/src/ticket_sale.gleam create mode 100644 examples/tickets_pog/src/tickets_pog.gleam create mode 100644 examples/tickets_pog/test/tickets_pog_test.gleam diff --git a/README.md b/README.md index 3b67cff..ae049e8 100644 --- a/README.md +++ b/README.md @@ -1,49 +1,81 @@ # Factos -Prototype context-first event-sourcing helpers for Gleam. +Factos is a set of prototype Gleam libraries for context-first Event Sourcing. -This package deliberately does not model Event Sourcing as aggregates. A command -capability reads the facts relevant to one decision, folds a temporary decision -model, decides which new facts to record, and records them only when the relevant -context can be protected. +The libraries are based on the interpretation described in Rico Fritzsche's +[Simply Event Sourcing](https://ricofritzsche.me/simply-event-sourcing/): Event +Sourcing is not defined by aggregates, aggregate roots, CQRS, message brokers, +microservices, or stream-per-object storage. Event Sourcing means accepted facts +are persisted as the authoritative history of the system, and that relevant +history is used when deciding whether new facts may be accepted. -## Opinion +Factos models that idea directly: -Event Sourcing is the persistence idea: accepted facts are the authoritative -state of the system. Aggregates, CQRS, projections, message brokers, and stream -versioning are implementation choices. +1. A command arrives with an intention. +2. A domain capability chooses the facts relevant to that decision. +3. Those facts are folded into a temporary decision state. +4. The decision either rejects the command with a domain error or produces new facts. +5. The store appends the new facts only if the relevant context has remained stable. -This prototype follows the shape described by Command Context Consistency and -Dynamic Consistency Boundaries: +The consistency boundary follows the command decision. It is not forced to be a +predefined `User`, `Order`, or `Customer` aggregate stream. -1. A command defines the context it needs. -2. The context is expressed as a query over event types and tags. -3. The application folds only those facts into a decision model. -4. The application produces new facts. -5. The store must reject the append if matching facts appeared after the context - was observed. +## Libraries -That last step requires store support. KurrentDB's normal append API supports -expected stream revision checks. That is useful, but it is not the same as a -DCB-style query-conditioned append. +This repository contains three Gleam libraries: -## Domain Model +1. `factos`: store-independent domain primitives. +2. `factos_sqlight`: SQLite backend implemented with the `sqlight` package. +3. `factos_kurrentdb_erlang`: KurrentDB backend for the Erlang target. -Keep commands, events, state, decisions, evolution, and errors in your app. +The core library is intentionally small. It knows about facts, event types, tags, +queries, contexts, deciders, views, recorded events, loaded streams, and append +conditions. It does not know how bytes are encoded, where events are stored, how +subscriptions work, whether projections are synchronous, or which transport is +used. -```gleam -import factos +Backend libraries own storage details. They define storage codecs, persistence +errors, migrations, and dispatch functions for their storage technology. -pub type Command { - RegisterUser(username: String) -} +## Concepts + +### Events Are Facts + +An event is a fact that has been accepted by the application. The event history is +the source of truth. Derived state can be rebuilt by folding events with an +evolution function. + +Factos does not require a base `Event` interface. Your application defines its own +event type: +```gleam pub type Event { UsernameReserved(username: String) UserRegistered(username: String) + DisplayNameChanged(user_id: String, name: String) +} +``` + +### Deciders Are Pure Domain Capabilities + +A `Decider` is a pure command-handling component made from: + +1. an initial state, +2. a decision function, and +3. an evolution function. + +The decision function receives the temporary state needed for one command and +returns either new events or a domain error. The evolution function folds accepted +events into that state. + +```gleam +import factos + +pub type Command { + RegisterUser(username: String) } -pub type UsernameState { +pub type State { UsernameAvailable UsernameTaken } @@ -52,37 +84,34 @@ pub type DomainError { UsernameAlreadyTaken } -pub fn evolve(state: UsernameState, event: Event) -> UsernameState { +pub fn evolve(state: State, event: Event) -> State { case state, event { UsernameAvailable, UsernameReserved(_) -> UsernameTaken UsernameAvailable, UserRegistered(_) -> UsernameTaken + UsernameAvailable, DisplayNameChanged(_, _) -> state UsernameTaken, UsernameReserved(_) -> state UsernameTaken, UserRegistered(_) -> state + UsernameTaken, DisplayNameChanged(_, _) -> state } } -pub fn decide(state: UsernameState, command: Command) { +pub fn decide(state: State, command: Command) -> Result(List(Event), DomainError) { case state, command { UsernameAvailable, RegisterUser(username) -> Ok([UserRegistered(username)]) UsernameTaken, RegisterUser(_) -> Error(UsernameAlreadyTaken) } } -pub fn registration_decider() { +pub fn registration_decider() -> factos.Decider(Command, State, Event, DomainError) { factos.decider( initial: UsernameAvailable, - decide: decide, - evolve: evolve, + decide:, + evolve:, ) } ``` -`Decider` is inspired by FModel: it is a small pure domain component made from -three things only: initial state, a decision function, and an evolution function. -It does not know about SQLite, KurrentDB, projections, subscriptions, HTTP, retries, or -repositories. - -You can test a decider without any storage: +Deciders are easy to test without any storage: ```gleam factos.compute_events( @@ -92,14 +121,15 @@ factos.compute_events( ) ``` -## Command Context +### Command Context Consistency -The command context is not a `User` aggregate. It is the facts relevant to the -decision: has this username been reserved or registered? +The command context is the set of facts required to make one decision. -```gleam -import factos +For registering a username, the command does not need every event for a `User` +object. It only needs facts that can make that username unavailable, such as +`UsernameReserved` and `UserRegistered` for the same username. +```gleam pub fn username_context(username: String) -> factos.Query { factos.query([ factos.query_item( @@ -113,110 +143,187 @@ pub fn username_context(username: String) -> factos.Query { } ``` -Query items are OR-combined. Within one item, event types are OR-combined and -tags are AND-combined. +Factos query semantics are deliberately simple: + +1. `factos.query([])` becomes `AllEvents`. +2. Query items are OR-combined. +3. Within one query item, event types are OR-combined. +4. Within one query item, tags are AND-combined. +5. Empty types in an item match any event type. +6. Empty tags in an item match any tags. + +When a backend reads a context it returns a `factos.Context` containing: + +1. the query that defined the context, +2. the folded decision state, +3. the matching recorded events, +4. the highest observed sequence position, and +5. an append condition: `FailIfEventsMatch(query, after: position)`. -## Tags +That append condition captures Command Context Consistency: append the newly +decided facts only if no facts matching the command context appeared after the +position used for the decision. -Tags are an explicit query contract. If a future command needs to select events -by username, account, invoice number, or product, that value must be exposed as a -tag when the event is written. +### Dynamic Consistency Boundary Tags + +Dynamic Consistency Boundary (DCB) applies the same context-first consistency +principle through a tag-based event-store contract. Event data is opaque to the +store, so anything that must be queryable for context reads or consistency checks +has to be exposed as an event type or tag when writing the event. ```gleam factos.tag("username:renata") factos.tag("account:abc123") +factos.tag("restaurant") +factos.tag("sku:burger") ``` -This is intentionally opinionated. Tags duplicate selected payload information, -but they make consistency and query needs visible at the event-store boundary. - -## Backends +Tags intentionally duplicate selected payload information. That duplication is +the contract: it makes future command-context queries visible at the event-store +boundary instead of hiding them inside opaque payloads. -Factos is split into a small core package and backend packages: +### Stream Consistency Is Still Supported -1. `factos` contains the store-independent domain primitives: `Decider`, `View`, `Query`, `Tag`, `Recorded`, and `Context`. -2. `backends/factos_sqlight` provides module `factos/sqlight` for SQLite via the `sqlight` package. -3. `backends/factos_kurrentdb_erlang` provides module `factos/kurrentdb` for KurrentDB on Erlang. +Factos also supports stream-based workflows through `load_stream` and +`dispatch_stream` in the backends. This is useful when a single stream really is +the right boundary for a decision. -Backends own their storage codecs and storage errors. The core package does not depend on either SQLite or KurrentDB. +Stream revision checks are not the definition of Event Sourcing. They are one +possible consistency strategy. They can over-conflict when unrelated events share +the same stream, and they can under-model rules that require facts from multiple +streams. -## Codecs +## Core Library: `factos` -The app owns encoding and decoding. Backend codecs return both the domain event -and the event-store metadata needed by the context API. +Import the core package when you want pure domain components and shared event +metadata types. ```gleam -import factos/sqlight as factos_sqlight -import sqlight - -pub fn codec() -> factos_sqlight.EventCodec(Event, DecodeError) { - factos_sqlight.EventCodec(encode: encode, decode: decode_event) -} +import factos ``` -For SQLite, `encode` returns a `factos_sqlight.Proposed` with an id, domain event, -event type, tags, and encoded bytes. `decode` returns a `factos.Decoded` with the -domain event, event type, and tags read from the stored event. +The core library provides: + +1. `EventType` and `Tag` wrappers for store-visible event metadata. +2. `Query` and `QueryItem` for command contexts. +3. `SequencePosition` for global event-log positions. +4. `AppendCondition` for context-stability requirements. +5. `Decider` for command-side decisions. +6. `View` for query-side projection folds. +7. `Decoded`, `Recorded`, `Context`, and `LoadedStream` records used by backends. -The library does not force one JSON shape for tags. In a real app, store tags in -custom metadata or payload fields consistently, and decode them at the boundary. +### Pure Command Computation -## SQLite Backend +Use `compute_events` when you already have relevant event history and want to test +or run a decider without storage: ```gleam -import factos/sqlight as factos_sqlight +factos.compute_events( + decider: registration_decider(), + events: [UsernameReserved("renata")], + command: RegisterUser("renata"), +) +``` -use connection <- sqlight.with_connection("file:events.sqlite3") -let assert Ok(Nil) = factos_sqlight.migrate(connection) +Use `compute_state` when you want to apply the events produced by a decision to +an existing state: -factos_sqlight.dispatch_context( - connection, - stream: "facts", - query: username_context("renata"), +```gleam +factos.compute_state( decider: registration_decider(), - codec: codec(), + current: option.None, command: RegisterUser("renata"), ) ``` -The SQLite backend stores an append-only `factos_events` table and uses `BEGIN IMMEDIATE` while dispatching commands. It can enforce `FailIfEventsMatch(query, after)` inside the same SQLite transaction. +### Projection Computation -## Reading A Context +`View` is the projection-side equivalent of a decider's `evolve` function. It is +also pure and store-independent. ```gleam -factos_sqlight.read_context( - connection, - query: username_context("renata"), - decider: registration_decider(), - codec: codec(), -) +let registrations = + factos.view(initial: 0, evolve: fn(count, event) { + case event { + UserRegistered(_) -> count + 1 + UsernameReserved(_) -> count + DisplayNameChanged(_, _) -> count + } + }) + +factos.project(view: registrations, events: [ + UserRegistered("renata"), + UserRegistered("lucy"), +]) +``` + +Views can be merged when they consume the same event type: + +```gleam +let dashboard = factos.merge_views(registrations, display_name_changes) +``` + +Factos intentionally stops at pure projection computation. Materialized view +storage, catch-up subscriptions, delivery retries, and read-model rebuilds belong +to application or backend-specific code. + +## SQLite Backend: `factos_sqlight` + +The SQLite backend stores events in an append-only table named `factos_events` and +uses `BEGIN IMMEDIATE` while dispatching commands. That lets it enforce +`FailIfEventsMatch(query, after)` transactionally in the same database that stores +events. + +```gleam +import factos/factos_sqlight +import sqlight + +use connection <- sqlight.with_connection("events.sqlite3") +let assert Ok(Nil) = factos_sqlight.migrate(connection) ``` -Backend `read_context` functions load events, decode them, filter them by the -full query, fold state with the decider, and return a `factos.Context`. +The schema contains: -The returned context includes: +1. `position`: monotonically increasing SQLite row position. +2. `id`: application-provided event id. +3. `stream`: stream name used for stream-based dispatch. +4. `revision`: per-stream revision. +5. `type`: event type name. +6. `tags`: newline-separated tag text. +7. `data`: opaque application-encoded bytes. -1. The folded decision state. -2. The matching recorded events. -3. The highest observed sequence position. -4. A DCB-style append condition: `FailIfEventsMatch(query, after: position)`. +The table enforces `unique(stream, revision)` and indexes stream revisions and +positions. -## Appending +### SQLite Codecs -The ideal append condition is: +Your application owns encoding and decoding. `factos_sqlight.EventCodec` keeps the +backend generic over event payloads and domain event types. ```gleam -factos.FailIfEventsMatch(query, after: position) +pub fn codec() -> factos_sqlight.EventCodec(Event, DecodeError) { + factos_sqlight.EventCodec(encode: encode, decode: decode) +} + +fn encode(event: Event) -> factos_sqlight.Proposed(Event) { + factos_sqlight.Proposed( + id: "event-" <> event.username, + event: event, + type_: factos.event_type("UserRegistered"), + tags: [factos.tag("username:" <> event.username)], + data: bit_array.from_string(event.username), + ) +} ``` -That means: append these new facts only if no facts matching the command context -appeared after the position used for the decision. +The encoder returns the domain event, event type, tags, and bytes to persist. The +decoder receives the stored row and must return `factos.Decoded(event)` with the +domain event, event type, and tags that should participate in query matching. + +### SQLite Context Dispatch -The SQLite backend enforces this condition transactionally. The KurrentDB backend -models the same condition, but returns `UnsupportedAppendCondition` for it because -regular KurrentDB append checks stream revisions, not arbitrary event-type/tag -queries. +Use `dispatch_context` when a command's consistency boundary is a query over event +types and tags rather than a single stream. ```gleam factos_sqlight.dispatch_context( @@ -229,12 +336,19 @@ factos_sqlight.dispatch_context( ) ``` -Use this shape for stores that support DCB-style atomic append conditions. +`dispatch_context` performs the full read-decide-append flow inside a transaction: -## Stream Consistency +1. begin an immediate SQLite transaction, +2. read matching events, +3. fold the decision state, +4. run the decider, +5. check whether matching events appeared after the observed position, +6. append produced events to the target stream, and +7. commit or roll back. -KurrentDB can safely protect a single stream with expected revision checks. This -is still useful when a stream is the right consistency boundary. +### SQLite Stream Dispatch + +Use `dispatch_stream` when the stream is the intended consistency boundary. ```gleam factos_sqlight.dispatch_stream( @@ -246,71 +360,126 @@ factos_sqlight.dispatch_stream( ) ``` -This is aggregate-stream style consistency. It is not the definition of Event -Sourcing, and it may over-conflict when unrelated events share the same stream. +The backend loads the stream, folds state, decides, and appends only if the +stream revision still matches the revision that was loaded. -## Views +## KurrentDB Erlang Backend: `factos_kurrentdb_erlang` -`View` is the projection-side equivalent of a decider's `evolve` function. It is -also pure and store-independent. +The KurrentDB backend integrates Factos with the Erlang-target KurrentDB client. +It supports stream reads, stream appends with expected revisions, and context +reads from `$all` using event-type filters. ```gleam -let registrations = - factos.view(initial: 0, evolve: fn(count, event) { - case event { - UserRegistered(_) -> count + 1 - UsernameReserved(_) -> count - } - }) +import factos/factos_kurrentdb_erlang +import kurrentdb +import kurrentdb_erlang -factos.project(view: registrations, events: [ - UserRegistered("renata"), - UserRegistered("lucy"), -]) +let assert Ok(client) = + kurrentdb.from_connection_string( + "kurrentdb://admin:changeit@localhost:2113?tls=true", + ) + +let assert Ok(connection) = + kurrentdb_erlang.new(client) + |> kurrentdb_erlang.verify_ca_certificate_file("certs/ca.crt") + |> kurrentdb_erlang.start(option.None) ``` -Views can be merged when they consume the same event type: +### KurrentDB Codecs + +`factos_kurrentdb_erlang.EventCodec` adapts between domain events and +`append_to_stream.Event` values from the KurrentDB client. ```gleam -let dashboard = factos.merge_views(registrations, reservations) +pub fn codec() -> factos_kurrentdb_erlang.EventCodec(Event, DecodeError) { + factos_kurrentdb_erlang.EventCodec(encode: encode, decode: decode) +} ``` -Factos intentionally stops at pure projection computation. It does not provide a -materialized-view repository abstraction yet; persistence and delivery choices -belong outside the domain component. +The encoder returns `Proposed(event, type_, tags, message)`. The `message` is the +actual KurrentDB append event. The decoder receives a KurrentDB recorded event and +returns a `factos.Decoded(event)`. -## KurrentDB Backend Tradeoffs +### KurrentDB Stream Dispatch -KurrentDB support available through `factos_kurrentdb_erlang`: +KurrentDB's regular append API can protect a stream revision. Use +`dispatch_stream` for that flow. -1. Read a single stream. -2. Append to a stream with `NoStream`, `Revision(n)`, `StreamExists`, or `Any`. -3. Read `$all` with event-type or stream-name filters. +```gleam +factos_kurrentdb_erlang.dispatch_stream( + connection, + stream: "user-renata", + decider: registration_decider(), + codec: codec(), + command: RegisterUser("renata"), + timeout: 10_000, +) +``` -Newer KurrentDB versions also support secondary and user-defined indexes that -can be consumed through `$all` stream-prefix filters such as `$idx-et-...` or -`$idx-user-...`. Those improve reads, but the docs describe secondary indexes as -eventually consistent. They should not be treated as a command-decision -consistency guarantee unless the write path can atomically enforce the same -condition. +Empty streams map to `factos.NoEvents`; loaded streams map to +`factos.CurrentRevision(n)`. Appends use KurrentDB expected-revision checks. -## Inspired By FModel +### KurrentDB Context Reads -Factos borrows FModel's useful core idea: model behavior as pure data structures -that hold functions (`Decider`, `View`) and keep infrastructure outside them. +`read_context` can read from `$all`. It translates the event types in a +`factos.Query` into a KurrentDB `$all` event-type prefix filter, decodes events, +then applies full Factos query matching locally, including tags. -Factos intentionally does not copy these FModel parts yet: +```gleam +factos_kurrentdb_erlang.read_context( + connection, + query: username_context("renata"), + decider: registration_decider(), + codec: codec(), + timeout: 10_000, +) +``` + +The returned context still contains `FailIfEventsMatch(query, after: position)`. +However, KurrentDB's regular append operation cannot atomically enforce arbitrary +event-type/tag query conditions. For that reason `dispatch_context` returns +`UnsupportedAppendCondition` for `FailIfEventsMatch`. + +This is intentional documentation of the tradeoff: KurrentDB stream revision +checks are useful, but they are not the same as a DCB-style query-conditioned +append. If your consistency rule is genuinely context-based across streams, you +need a write path that can atomically enforce that context condition. -1. `Aggregate` wrappers, because the package is trying not to make aggregates the - center of the model. -2. Generic repository traits, because Gleam code can pass concrete functions and - records without committing to one application architecture. -3. Sagas/process managers, because they are event-driven messaging/workflow - concerns and should remain separate from the event-sourcing core until a real - use case needs them. +## Example -## Development +The `examples/src/order_workflow.gleam` file contains a restaurant order workflow +using `factos_sqlight`. It demonstrates: + +1. domain commands and events, +2. a custom state machine, +3. domain-specific errors, +4. stream dispatch for one order, +5. application-owned encoding and decoding, and +6. a projection view for kitchen summary data. + +Run it from the examples package: ```sh -gleam test +cd examples +gleam run ``` + +## Tradeoffs + +Factos is a prototype. It deliberately leaves many production concerns outside the +core package: + +1. event schema evolution, +2. snapshots, +3. subscriptions, +4. projection repositories, +5. retry policies, +6. side-effect orchestration, +7. idempotency policies beyond event ids, +8. serialization format choices, and +9. distributed deployment concerns. + +Those are real engineering problems, but they are separate from the core Event +Sourcing definition. Factos keeps the starting point simple: persist accepted +facts, derive temporary decision state from relevant history, and record new facts +only if that relevant history is still valid. diff --git a/backends/factos_kurrentdb_erlang/gleam.toml b/backends/factos_kurrentdb_erlang/gleam.toml index b94ce90..8e66d06 100644 --- a/backends/factos_kurrentdb_erlang/gleam.toml +++ b/backends/factos_kurrentdb_erlang/gleam.toml @@ -1,5 +1,11 @@ name = "factos_kurrentdb_erlang" version = "1.0.0" +description = "KurrentDB Erlang backend for Factos context-first Event Sourcing." +licences = ["Apache-2.0"] +links = [ + { title = "Factos", href = "https://factos.hexdocs.pm" }, + { title = "Simply Event Sourcing", href = "https://ricofritzsche.me/simply-event-sourcing/" }, +] [dependencies] factos = { path = "../.." } diff --git a/backends/factos_kurrentdb_erlang/src/factos/kurrentdb.gleam b/backends/factos_kurrentdb_erlang/src/factos/factos_kurrentdb_erlang.gleam similarity index 77% rename from backends/factos_kurrentdb_erlang/src/factos/kurrentdb.gleam rename to backends/factos_kurrentdb_erlang/src/factos/factos_kurrentdb_erlang.gleam index f43fede..e501221 100644 --- a/backends/factos_kurrentdb_erlang/src/factos/kurrentdb.gleam +++ b/backends/factos_kurrentdb_erlang/src/factos/factos_kurrentdb_erlang.gleam @@ -1,8 +1,19 @@ //// KurrentDB Erlang backend for Factos. //// -//// This backend supports stream revision consistency. It can read command -//// contexts from `$all`, but KurrentDB's regular append operation cannot -//// atomically enforce Factos' DCB-style `FailIfEventsMatch` append condition. +//// This backend integrates Factos with the Erlang-target KurrentDB client. It +//// supports stream revision consistency and context reads from `$all`. +//// +//// KurrentDB's regular append API can protect a stream with expected revision +//// checks. That is useful when the stream is the real consistency boundary. +//// However, a Factos context append condition is query-based: +//// `FailIfEventsMatch(query, after)`. KurrentDB's regular append operation cannot +//// atomically enforce an arbitrary event-type/tag query condition, so this backend +//// reports `UnsupportedAppendCondition` for context dispatch. +//// +//// This distinction is intentional. Stream revision consistency is one valid +//// Event Sourcing implementation strategy; Command Context Consistency and DCB- +//// style tag contracts require stores or write paths that can protect the actual +//// command context. import factos import gleam/list @@ -14,6 +25,11 @@ import kurrentdb_erlang import youid/uuid pub type Proposed(event) { + /// A domain event prepared for KurrentDB persistence. + /// + /// The application codec creates this value. `event` is kept for domain-level + /// typing, `type_` and `tags` are Factos query metadata, and `message` is the + /// KurrentDB append event sent to the client. Proposed( event: event, type_: factos.EventType, @@ -23,6 +39,11 @@ pub type Proposed(event) { } pub type EventCodec(event, decode_error) { + /// Application-owned KurrentDB event codec. + /// + /// `encode` converts a domain event into a KurrentDB append message plus Factos + /// metadata. `decode` converts a KurrentDB recorded event into a + /// `factos.Decoded` domain event. Decode failures are wrapped as `DecodeError`. EventCodec( encode: fn(event) -> Proposed(event), decode: fn(read_stream.RecordedEvent) -> @@ -31,14 +52,35 @@ pub type EventCodec(event, decode_error) { } pub type Error(domain_error, decode_error) { + /// The decider rejected the command with a domain error. DomainError(domain_error) + + /// The application codec could not decode a stored event. DecodeError(decode_error) + + /// KurrentDB returned an error while reading a stream or `$all`. ReadError(kurrentdb_erlang.Error(read_stream.ResponseError)) + + /// KurrentDB returned an error while appending to a stream. AppendError(kurrentdb_erlang.Error(append_to_stream.ResponseError)) + + /// A read stream did not produce a message before the configured timeout. ReadTimedOut + + /// The requested append condition cannot be enforced by this backend. + /// + /// This is expected for `FailIfEventsMatch` because the regular KurrentDB append + /// API protects stream revisions, not arbitrary Factos context queries. UnsupportedAppendCondition(factos.AppendCondition) } +/// Read and fold a command context from KurrentDB `$all`. +/// +/// Event types in the query are translated into a KurrentDB event-type prefix +/// filter where possible. Decoded events are then filtered locally with the full +/// Factos query, including tags. The returned context contains a +/// `FailIfEventsMatch(query, after)` append condition, but this backend cannot +/// enforce that condition during append. pub fn read_context( connection: kurrentdb_erlang.Connection, query query: factos.Query, @@ -51,6 +93,15 @@ pub fn read_context( read_context_events(connection, query, initial, evolve, codec, timeout) } +/// Attempt a full context-first dispatch flow. +/// +/// This function reads the context and runs the decider, then delegates to +/// `append_with_condition`. Because normal KurrentDB appends cannot enforce +/// `FailIfEventsMatch`, context dispatch returns `UnsupportedAppendCondition` for +/// the context condition produced by `read_context`. +/// +/// Use this function to make the storage limitation explicit. Prefer +/// `dispatch_stream` when a stream revision is the intended consistency boundary. pub fn dispatch_context( connection: kurrentdb_erlang.Connection, stream stream_name: String, @@ -62,10 +113,10 @@ pub fn dispatch_context( ) -> Result(append_to_stream.Append, Error(domain_error, decode_error)) { use context <- result.try(read_context( connection, - query: query, - decider: decider, - codec: codec, - timeout: timeout, + query:, + decider:, + codec:, + timeout:, )) use pair <- result.try( factos.decide_context(context, command, decider) @@ -83,6 +134,11 @@ pub fn dispatch_context( ) } +/// Load and fold one KurrentDB stream. +/// +/// Missing streams are treated as empty streams and returned with +/// `factos.NoEvents`. Existing streams return their latest observed revision as +/// `factos.CurrentRevision(n)`. pub fn load_stream( connection: kurrentdb_erlang.Connection, stream stream_name: String, @@ -98,6 +154,11 @@ pub fn load_stream( load_stream_events(connection, stream_name, initial, evolve, codec, timeout) } +/// Run a stream-based read-decide-append command flow. +/// +/// The backend loads the target stream, folds it into state, runs the decider, and +/// appends produced events with KurrentDB expected-revision checks. Use this when +/// one KurrentDB stream is the correct consistency boundary for the command. pub fn dispatch_stream( connection: kurrentdb_erlang.Connection, stream stream_name: String, @@ -236,7 +297,7 @@ fn append_to_stream_with_config( events: List(event), codec: EventCodec(event, decode_error), config: append_to_stream.Configuration, - timeout: Int, + within: Int, ) -> Result(append_to_stream.Append, Error(domain_error, decode_error)) { case events { [] -> @@ -250,14 +311,11 @@ fn append_to_stream_with_config( kurrentdb_erlang.append_to_stream( connection, stream: stream_name, - events: list.map(events, fn(event) { - let Proposed(message: message, ..) = encode(event) - message - }), + events: list.map(events, fn(event) { encode(event).message }), config: config, ) - kurrentdb_erlang.await(task, within: timeout) + kurrentdb_erlang.await(task, within:) |> result.map_error(AppendError) } } @@ -271,9 +329,9 @@ fn receive_context( codec: EventCodec(event, decode_error), events: List(factos.Recorded(event)), position: factos.SequencePosition, - timeout: Int, + within: Int, ) -> Result(factos.Context(event, state), Error(domain_error, decode_error)) { - case kurrentdb_erlang.receive(stream, within: timeout) { + case kurrentdb_erlang.receive(stream, within:) { Error(Nil) -> { kurrentdb_erlang.close(stream) Error(ReadTimedOut) @@ -292,18 +350,6 @@ fn receive_context( kurrentdb_erlang.close(stream) Error(ReadError(error)) } - Ok(kurrentdb_erlang.ReadEvent(event)) -> - receive_context_event( - stream, - query, - state, - evolve, - codec, - events, - position, - timeout, - event, - ) Ok(kurrentdb_erlang.ReadMessage(read_stream.ReadEvent(event))) -> receive_context_event( stream, @@ -313,7 +359,7 @@ fn receive_context( codec, events, position, - timeout, + within, event, ) Ok(kurrentdb_erlang.ReadMessage(read_stream.LastAllStreamPosition(read_stream.Position( @@ -331,7 +377,7 @@ fn receive_context( position, factos.SequencePosition(commit_position), ), - timeout, + within, ) Ok(kurrentdb_erlang.ReadMessage(_)) -> receive_context( @@ -342,7 +388,7 @@ fn receive_context( codec, events, position, - timeout, + within, ) } } @@ -395,12 +441,12 @@ fn receive_stream( codec: EventCodec(event, decode_error), events: List(factos.Recorded(event)), revision: factos.Revision, - timeout: Int, + within: Int, ) -> Result( factos.LoadedStream(event, state), Error(domain_error, decode_error), ) { - case kurrentdb_erlang.receive(stream, within: timeout) { + case kurrentdb_erlang.receive(stream, within:) { Error(Nil) -> { kurrentdb_erlang.close(stream) Error(ReadTimedOut) @@ -427,17 +473,6 @@ fn receive_stream( _ -> Error(ReadError(error)) } } - Ok(kurrentdb_erlang.ReadEvent(event)) -> - receive_stream_event( - stream, - stream_name, - state, - evolve, - codec, - events, - timeout, - event, - ) Ok(kurrentdb_erlang.ReadMessage(read_stream.ReadEvent(event))) -> receive_stream_event( stream, @@ -446,7 +481,7 @@ fn receive_stream( evolve, codec, events, - timeout, + within, event, ) Ok(kurrentdb_erlang.ReadMessage(_)) -> @@ -458,7 +493,7 @@ fn receive_stream( codec, events, revision, - timeout, + within, ) } } diff --git a/backends/factos_kurrentdb_erlang/test/factos_kurrentdb_erlang_test.gleam b/backends/factos_kurrentdb_erlang/test/factos_kurrentdb_erlang_test.gleam index e8c339f..017ee46 100644 --- a/backends/factos_kurrentdb_erlang/test/factos_kurrentdb_erlang_test.gleam +++ b/backends/factos_kurrentdb_erlang/test/factos_kurrentdb_erlang_test.gleam @@ -1,5 +1,5 @@ import factos -import factos/kurrentdb as factos_kurrentdb +import factos/factos_kurrentdb_erlang import gleam/bit_array import gleam/int import gleam/list @@ -51,7 +51,7 @@ pub fn dispatch_stream_handles_many_events_integration_test() { dispatch_counter_stream_many(stream_name, event_type, 100) let assert Ok(loaded) = - factos_kurrentdb.load_stream( + factos_kurrentdb_erlang.load_stream( connection(), stream: stream_name, decider: counter_decider(), @@ -77,7 +77,7 @@ pub fn read_context_handles_many_streams_integration_test() { dispatch_counter_context_streams_many(event_type, 50) let assert Ok(context) = - factos_kurrentdb.read_context( + factos_kurrentdb_erlang.read_context( connection(), query: query, decider: counter_decider(), @@ -111,9 +111,12 @@ fn dispatch_counter_stream_many( stream_name: String, event_type: String, remaining: Int, -) -> Result(append_to_stream.Append, factos_kurrentdb.Error(Nil, DecodeError)) { +) -> Result( + append_to_stream.Append, + factos_kurrentdb_erlang.Error(Nil, DecodeError), +) { let result = - factos_kurrentdb.dispatch_stream( + factos_kurrentdb_erlang.dispatch_stream( connection(), stream: stream_name, decider: counter_decider(), @@ -133,9 +136,12 @@ fn dispatch_counter_stream_many( fn dispatch_counter_context_streams_many( event_type: String, remaining: Int, -) -> Result(append_to_stream.Append, factos_kurrentdb.Error(Nil, DecodeError)) { +) -> Result( + append_to_stream.Append, + factos_kurrentdb_erlang.Error(Nil, DecodeError), +) { let result = - factos_kurrentdb.dispatch_stream( + factos_kurrentdb_erlang.dispatch_stream( connection(), stream: unique_name("counter-context"), decider: counter_decider(), @@ -183,8 +189,8 @@ fn counter_evolve(state: CounterState, event: CounterEvent) -> CounterState { fn counter_codec( event_type: String, -) -> factos_kurrentdb.EventCodec(CounterEvent, DecodeError) { - factos_kurrentdb.EventCodec( +) -> factos_kurrentdb_erlang.EventCodec(CounterEvent, DecodeError) { + factos_kurrentdb_erlang.EventCodec( encode: encode_counter_event(_, event_type), decode: decode_counter_event(_, event_type), ) @@ -193,10 +199,10 @@ fn counter_codec( fn encode_counter_event( event: CounterEvent, event_type: String, -) -> factos_kurrentdb.Proposed(CounterEvent) { +) -> factos_kurrentdb_erlang.Proposed(CounterEvent) { case event { Incremented(value) -> - factos_kurrentdb.Proposed( + factos_kurrentdb_erlang.Proposed( event: event, type_: factos.event_type(event_type), tags: [factos.tag("counter:load")], diff --git a/backends/factos_pog/README.md b/backends/factos_pog/README.md new file mode 100644 index 0000000..785cef7 --- /dev/null +++ b/backends/factos_pog/README.md @@ -0,0 +1,24 @@ +# factos_pog + +PostgreSQL backend for Factos using [`pog`](https://hex.pm/packages/pog). + +This backend follows the "Simply Event Sourcing" shape used by Factos: accepted facts are stored in an append-only event table, command handlers select the facts relevant to their decision, and appends are accepted only when that command context is still stable. + +## Design Notes + +The table stores opaque event bytes plus store-visible query metadata: event type and tags. This is a DCB-style tradeoff. PostgreSQL does not need to understand payloads, but any payload value needed by future context queries must be exposed as a tag when the event is written. + +`dispatch_context` runs inside a PostgreSQL transaction and locks the event table before reading, deciding, checking, and appending. This is intentionally conservative. It makes arbitrary `FailIfEventsMatch(query, after)` checks correct without trying to infer lock keys from dynamic query metadata. A higher-throughput backend could replace the table lock with advisory locks or more granular query-specific locks, but only if it preserves the same context-stability guarantee. + +`dispatch_stream` is also available for applications where one stream revision really is the intended consistency boundary. It is an implementation strategy, not the definition of Event Sourcing. + +## Usage + +Start a `pog` pool in your application supervision tree, run `migrate`, then call `dispatch_context` or `dispatch_stream` with your domain decider and codec. + +```gleam +let connection = pog.named_connection(pool_name) +let assert Ok(Nil) = factos_pog.migrate(connection) +``` + +Your codec owns event serialization. The backend only persists bytes and query metadata. diff --git a/backends/factos_pog/compose.yml b/backends/factos_pog/compose.yml new file mode 100644 index 0000000..504ee59 --- /dev/null +++ b/backends/factos_pog/compose.yml @@ -0,0 +1,14 @@ +services: + postgres: + image: postgres:18 + environment: + POSTGRES_DB: factos_pog + POSTGRES_USER: postgres + POSTGRES_PASSWORD: postgres + ports: + - "5432:5432" + healthcheck: + test: ["CMD-SHELL", "pg_isready -U postgres -d factos_pog"] + interval: 1s + timeout: 5s + retries: 20 diff --git a/backends/factos_pog/gleam.toml b/backends/factos_pog/gleam.toml new file mode 100644 index 0000000..06d83e0 --- /dev/null +++ b/backends/factos_pog/gleam.toml @@ -0,0 +1,18 @@ +name = "factos_pog" +version = "1.0.0" +description = "PostgreSQL backend for Factos context-first Event Sourcing using pog." +licences = ["Apache-2.0"] +links = [ + { title = "Factos", href = "https://factos.hexdocs.pm" }, + { title = "Simply Event Sourcing", href = "https://ricofritzsche.me/simply-event-sourcing/" }, +] + +[dependencies] +factos = { path = "../.." } +gleam_stdlib = ">= 1.0.0 and < 2.0.0" +pog = ">= 4.1.0 and < 5.0.0" + +[dev_dependencies] +gleeunit = ">= 1.0.0 and < 2.0.0" +gleam_erlang = ">= 1.0.0 and < 2.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 new file mode 100644 index 0000000..310c988 --- /dev/null +++ b/backends/factos_pog/manifest.toml @@ -0,0 +1,31 @@ +# Do not manually edit this file, it is managed by Gleam. +# +# This file locks the dependency versions used, to make your build +# deterministic and to prevent unexpected versions from being included +# in your application. +# +# You should check this file into your source control repository. + +packages = [ + { name = "backoff", version = "1.1.6", build_tools = ["rebar3"], requirements = [], otp_app = "backoff", source = "hex", outer_checksum = "CF0CFFF8995FB20562F822E5CC47D8CCF664C5ECDC26A684CBE85C225F9D7C39" }, + { name = "exception", version = "2.1.1", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "exception", source = "hex", outer_checksum = "6BDEA95248093599391C3B5DF1835C5C6A86C353C2F99CE539B450E3432FE117" }, + { name = "factos", version = "1.0.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], source = "local", path = "../.." }, + { name = "gleam_erlang", version = "1.3.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_erlang", source = "hex", outer_checksum = "1124AD3AA21143E5AF0FC5CF3D9529F6DB8CA03E43A55711B60B6B7B3874375C" }, + { name = "gleam_otp", version = "1.2.0", build_tools = ["gleam"], requirements = ["gleam_erlang", "gleam_stdlib"], otp_app = "gleam_otp", source = "hex", outer_checksum = "BA6A294E295E428EC1562DC1C11EA7530DCB981E8359134BEABC8493B7B2258E" }, + { name = "gleam_stdlib", version = "1.0.3", build_tools = ["gleam"], requirements = [], otp_app = "gleam_stdlib", source = "hex", outer_checksum = "1F543AFBA5D33DA493E6087F4E4C4F20D899411343512686C98A8ABB2963CF22" }, + { name = "gleam_time", version = "1.8.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_time", source = "hex", outer_checksum = "533D8723774D61AD4998324F5DD1DABDCDBFABAFB9E87CB5D03C6955448FC97D" }, + { name = "gleeunit", version = "1.11.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleeunit", source = "hex", outer_checksum = "EC31ABA74256AEA531EDF8169931D775BBB384FED0A8A1BDC4DD9354E3E21826" }, + { 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" }, + { name = "pog", version = "4.1.0", build_tools = ["gleam"], requirements = ["exception", "gleam_erlang", "gleam_otp", "gleam_stdlib", "gleam_time", "pgo"], otp_app = "pog", source = "hex", outer_checksum = "E4AFBA39A5FAA2E77291836C9683ADE882E65A06AB28CA7D61AE7A3AD61EBBD5" }, +] + +[requirements] +factos = { path = "../.." } +gleam_erlang = { version = ">= 1.0.0 and < 2.0.0" } +gleam_stdlib = { version = ">= 1.0.0 and < 2.0.0" } +gleeunit = { version = ">= 1.0.0 and < 2.0.0" } +global_value = { version = ">= 1.0.0 and < 2.0.0" } +pog = { version = ">= 4.1.0 and < 5.0.0" } diff --git a/backends/factos_pog/src/factos/factos_pog.gleam b/backends/factos_pog/src/factos/factos_pog.gleam new file mode 100644 index 0000000..eeea368 --- /dev/null +++ b/backends/factos_pog/src/factos/factos_pog.gleam @@ -0,0 +1,619 @@ +//// PostgreSQL backend for Factos using the `pog` package. +//// +//// This backend stores accepted facts in an append-only `factos_events` table. +//// The event history is the source of truth; projections and stream-shaped reads +//// are derived views over that history. +//// +//// The context dispatch flow follows the Command Context Consistency idea from +//// "Simply Event Sourcing": a command selects the facts required for its decision, +//// folds them into temporary state, decides new facts, and appends those facts +//// only when no relevant facts appeared after the observed context position. +//// +//// The query contract is intentionally tag-based. PostgreSQL stores opaque event +//// bytes, an event type, and tags. This keeps domain serialization outside the +//// backend, but it means any payload value needed for a selective consistency +//// query must be written as a tag. + +import factos +import gleam/dynamic/decode +import gleam/list +import gleam/result +import gleam/string +import pog + +pub type Proposed(event) { + /// A domain event prepared for PostgreSQL persistence. + /// + /// The application codec creates this value. `id` should identify the event for + /// the application. `type_` and `tags` are store-visible query metadata. `data` + /// is opaque bytes owned by the application codec. + Proposed( + id: String, + event: event, + type_: factos.EventType, + tags: List(factos.Tag), + data: BitArray, + ) +} + +pub type StoredEvent { + /// A raw event row read from PostgreSQL before domain decoding. + /// + /// Decoders receive this value so they can inspect stored metadata and bytes. + /// `position` is the global append order. `revision` is the per-stream revision. + StoredEvent( + position: Int, + id: String, + stream: String, + revision: Int, + type_: factos.EventType, + tags: List(factos.Tag), + data: BitArray, + ) +} + +pub type EventCodec(event, decode_error) { + /// Application-owned PostgreSQL event codec. + /// + /// `encode` converts a domain event into bytes and metadata. `decode` converts a + /// stored row back into a `factos.Decoded` domain event. Decode failures are kept + /// in the application's own error type and wrapped as `DecodeError`. + EventCodec( + encode: fn(event) -> Proposed(event), + decode: fn(StoredEvent) -> Result(factos.Decoded(event), decode_error), + ) +} + +pub type Append { + /// Result of a successful append. + /// + /// `current_revision` is the latest revision of the target stream after the + /// append. `position` is the global position of the last inserted event, or + /// `NoPosition` when no events were produced. + Append(current_revision: Int, position: factos.SequencePosition) +} + +pub type Error(domain_error, decode_error) { + /// The decider rejected the command with a domain error. + DomainError(domain_error) + + /// The application codec could not decode a stored event. + DecodeError(decode_error) + + /// PostgreSQL or `pog` returned an error while running a query. + StoreError(pog.QueryError) + + /// A stream revision or context append condition failed. + AppendConditionFailed(factos.AppendCondition) +} + +/// Create or update the PostgreSQL schema required by this backend. +/// +/// The schema is an append-only `factos_events` table with a global identity +/// `position`, per-stream `revision`, event `type`, newline-encoded `tags`, and +/// opaque `data` bytes. The `(stream, revision)` uniqueness constraint supports +/// stream-revision consistency. Position indexes support context reads and checks. +pub fn migrate(connection: pog.Connection) -> Result(Nil, Error(_, _)) { + // `pog` uses prepared statements, and PostgreSQL does not allow multiple SQL + // commands in one prepared statement. Keep migrations split into individual + // statements rather than relying on client-side SQL script execution. + use _ <- result.try(execute_migration( + connection, + " + create table if not exists factos_events ( + position bigint generated always as identity primary key, + id text not null, + stream text not null, + revision integer not null, + type text not null, + tags text not null, + data bytea not null, + unique(stream, revision) + ) + ", + )) + use _ <- result.try(execute_migration( + connection, + " + create index if not exists factos_events_stream_revision + on factos_events(stream, revision) + ", + )) + use _ <- result.try(execute_migration( + connection, + " + create index if not exists factos_events_position + on factos_events(position) + ", + )) + Ok(Nil) +} + +fn execute_migration( + connection: pog.Connection, + sql: String, +) -> Result(Nil, Error(_, _)) { + pog.query(sql) + |> pog.execute(on: connection) + |> result.map(fn(_) { Nil }) + |> result.map_error(StoreError) +} + +/// Read and fold the facts selected by a command-context query. +/// +/// The backend reads stored rows, decodes them with the supplied codec, filters +/// them with `factos.matches_query`, folds matching events with the decider's +/// `evolve` function, and returns a `factos.Context` with a +/// `FailIfEventsMatch(query, after)` append condition. +pub fn read_context( + connection: pog.Connection, + query query: factos.Query, + decider decider: factos.Decider(command, state, event, domain_error), + codec codec: EventCodec(event, decode_error), +) -> Result(factos.Context(event, state), Error(domain_error, decode_error)) { + let factos.Decider(initial, _, evolve) = decider + + use events <- result.try(read_matching_events(connection, query, codec)) + let position = highest_recorded_position(events) + + Ok(factos.Context( + query:, + state: factos.evolve_recorded( + initial: initial, + events: events, + evolve: evolve, + ), + events: events, + position: position, + append_condition: factos.FailIfEventsMatch(query, position), + )) +} + +/// Run a full context-first read-decide-append command flow. +/// +/// PostgreSQL does not have a native primitive for "append if no row matching this +/// arbitrary event-type/tag query appeared after position N". This backend uses a +/// transaction plus `lock table factos_events in exclusive mode` to make the read, +/// context check, and append atomic for all Factos queries. +/// +/// That lock is the main throughput tradeoff: unrelated writers queue behind each +/// other even if their contexts do not overlap. It is deliberately simple and +/// correct. A future backend can use advisory locks or query-specific lock keys, +/// but only if it keeps the same context-stability guarantee. +pub fn dispatch_context( + connection: pog.Connection, + stream stream_name: String, + query query: factos.Query, + decider decider: factos.Decider(command, state, event, domain_error), + codec codec: EventCodec(event, decode_error), + command command: command, +) -> Result(Append, Error(domain_error, decode_error)) { + use transaction_connection <- run_locked_transaction(connection) + use context <- result.try(read_context( + transaction_connection, + query: query, + decider: decider, + codec: codec, + )) + use pair <- result.try( + factos.decide_context(context, command, decider) + |> result.map_error(DomainError), + ) + let #(context, events) = pair + + append_with_condition( + transaction_connection, + stream_name, + events, + codec, + context.append_condition, + ) +} + +/// Load and fold one stream. +/// +/// This supports classic stream-revision consistency. The returned +/// `factos.LoadedStream` contains the folded state, decoded recorded events, and +/// current stream revision. +pub fn load_stream( + connection: pog.Connection, + stream stream_name: String, + decider decider: factos.Decider(command, state, event, domain_error), + codec codec: EventCodec(event, decode_error), +) -> Result( + factos.LoadedStream(event, state), + Error(domain_error, decode_error), +) { + let factos.Decider(initial, _, evolve) = decider + use events <- result.try(read_stream_events(connection, stream_name, codec)) + + Ok(factos.LoadedStream( + stream: stream_name, + state: factos.evolve_recorded( + initial: initial, + events: events, + evolve: evolve, + ), + events: events, + revision: stream_revision(events), + )) +} + +/// Run a stream-based read-decide-append command flow. +/// +/// Use this when one stream is intentionally the consistency boundary. It remains +/// useful, but it is not required by Event Sourcing. For command-specific rules, +/// prefer `dispatch_context` so the protected boundary follows the decision. +pub fn dispatch_stream( + connection: pog.Connection, + stream stream_name: String, + decider decider: factos.Decider(command, state, event, domain_error), + codec codec: EventCodec(event, decode_error), + command command: command, +) -> Result(Append, Error(domain_error, decode_error)) { + use transaction_connection <- run_locked_transaction(connection) + use loaded <- result.try(load_stream( + transaction_connection, + stream: stream_name, + decider: decider, + codec: codec, + )) + let factos.Decider(_, decide, _) = decider + use events <- result.try( + decide(loaded.state, command) + |> result.map_error(DomainError), + ) + + append_stream_events( + transaction_connection, + stream_name, + events, + codec, + loaded.revision, + ) +} + +fn run_locked_transaction( + connection: pog.Connection, + work: fn(pog.Connection) -> Result(Append, Error(domain_error, decode_error)), +) -> Result(Append, Error(domain_error, decode_error)) { + case + { + use transaction_connection <- pog.transaction(connection) + use _ <- result.try( + pog.query("lock table factos_events in exclusive mode") + |> pog.execute(on: transaction_connection) + |> result.map(fn(_) { Nil }) + |> result.map_error(StoreError), + ) + work(transaction_connection) + } + { + Ok(append) -> Ok(append) + Error(pog.TransactionQueryError(error)) -> Error(StoreError(error)) + Error(pog.TransactionRolledBack(error)) -> Error(error) + } +} + +fn append_with_condition( + connection: pog.Connection, + stream_name: String, + events: List(event), + codec: EventCodec(event, decode_error), + condition: factos.AppendCondition, +) -> Result(Append, Error(domain_error, decode_error)) { + case condition { + factos.NoAppendCondition -> + append_current_stream(connection, stream_name, events, codec) + factos.FailIfEventsMatch(query, after) -> + case has_matching_events_after(connection, query, after) { + Error(error) -> Error(StoreError(error)) + Ok(True) -> Error(AppendConditionFailed(condition)) + Ok(False) -> + append_current_stream(connection, stream_name, events, codec) + } + } +} + +fn append_current_stream( + connection: pog.Connection, + stream_name: String, + events: List(event), + codec: EventCodec(event, decode_error), +) -> Result(Append, Error(domain_error, decode_error)) { + use revision <- result.try( + current_revision(connection, stream_name) + |> result.map_error(StoreError), + ) + append_stream_events( + connection, + stream_name, + events, + codec, + factos.CurrentRevision(revision), + ) +} + +fn append_stream_events( + connection: pog.Connection, + stream_name: String, + events: List(event), + codec: EventCodec(event, decode_error), + expected: factos.Revision, +) -> Result(Append, Error(domain_error, decode_error)) { + case events { + [] -> + Ok(Append( + current_revision: revision_to_int(expected), + position: factos.NoPosition, + )) + [_, ..] -> { + use current <- result.try( + current_revision(connection, stream_name) + |> result.map_error(StoreError), + ) + case expected_matches(expected, current) { + False -> Error(AppendConditionFailed(factos.NoAppendCondition)) + True -> + insert_events( + connection, + stream_name, + events, + codec, + current + 1, + factos.NoPosition, + ) + } + } + } +} + +fn insert_events( + connection: pog.Connection, + stream_name: String, + events: List(event), + codec: EventCodec(event, decode_error), + revision: Int, + position: factos.SequencePosition, +) -> Result(Append, Error(domain_error, decode_error)) { + case events { + [] -> Ok(Append(current_revision: revision - 1, position: position)) + [event, ..rest] -> { + let EventCodec(encode, _) = codec + let Proposed(id, _, type_, tags, data) = encode(event) + use returned <- result.try( + pog.query( + " + insert into factos_events (id, stream, revision, type, tags, data) + values ($1, $2, $3, $4, $5, $6) + returning position + ", + ) + |> pog.parameter(pog.text(id)) + |> pog.parameter(pog.text(stream_name)) + |> pog.parameter(pog.int(revision)) + |> pog.parameter(pog.text(factos.event_type_name(type_))) + |> pog.parameter(pog.text(tags_to_text(tags))) + |> pog.parameter(pog.bytea(data)) + |> pog.returning(int_field_decoder()) + |> pog.execute(on: connection) + |> result.map_error(StoreError), + ) + let position = case returned.rows { + [position, ..] -> factos.SequencePosition(position) + [] -> position + } + insert_events( + connection, + stream_name, + rest, + codec, + revision + 1, + position, + ) + } + } +} + +fn read_matching_events( + connection: pog.Connection, + query: factos.Query, + codec: EventCodec(event, decode_error), +) -> Result(List(factos.Recorded(event)), Error(domain_error, decode_error)) { + use rows <- result.try( + pog.query( + "select position, id, stream, revision, type, tags, data from factos_events order by position", + ) + |> pog.returning(stored_event_decoder()) + |> pog.execute(on: connection) + |> result.map(fn(returned) { returned.rows }) + |> result.map_error(StoreError), + ) + decode_rows(rows, codec) + |> result.map(list.filter(_, factos.matches_query(_, query))) +} + +fn read_stream_events( + connection: pog.Connection, + stream_name: String, + codec: EventCodec(event, decode_error), +) -> Result(List(factos.Recorded(event)), Error(domain_error, decode_error)) { + use rows <- result.try( + pog.query( + "select position, id, stream, revision, type, tags, data from factos_events where stream = $1 order by revision", + ) + |> pog.parameter(pog.text(stream_name)) + |> pog.returning(stored_event_decoder()) + |> pog.execute(on: connection) + |> result.map(fn(returned) { returned.rows }) + |> result.map_error(StoreError), + ) + decode_rows(rows, codec) +} + +fn decode_rows( + rows: List(StoredEvent), + codec: EventCodec(event, decode_error), +) -> Result(List(factos.Recorded(event)), Error(domain_error, decode_error)) { + case rows { + [] -> Ok([]) + [row, ..rest] -> { + use recorded <- result.try(decode_row(row, codec)) + use rest <- result.try(decode_rows(rest, codec)) + Ok([recorded, ..rest]) + } + } +} + +fn decode_row( + row: StoredEvent, + codec: EventCodec(event, decode_error), +) -> Result(factos.Recorded(event), Error(domain_error, decode_error)) { + let EventCodec(_, decode_event) = codec + use decoded <- result.try(decode_event(row) |> result.map_error(DecodeError)) + let factos.Decoded(event, type_, tags) = decoded + let StoredEvent(position, id, stream, revision, _, _, _) = row + + Ok(factos.Recorded( + id: id, + stream: stream, + revision: revision, + position: factos.SequencePosition(position), + type_: type_, + tags: tags, + event: event, + )) +} + +fn stored_event_decoder() -> decode.Decoder(StoredEvent) { + use position <- decode.field(0, decode.int) + use id <- decode.field(1, decode.string) + use stream <- decode.field(2, decode.string) + use revision <- decode.field(3, decode.int) + use type_name <- decode.field(4, decode.string) + use tags <- decode.field(5, decode.string) + use data <- decode.field(6, decode.bit_array) + decode.success(StoredEvent( + position: position, + id: id, + stream: stream, + revision: revision, + type_: factos.event_type(type_name), + tags: tags_from_text(tags), + data: data, + )) +} + +fn current_revision( + connection: pog.Connection, + stream_name: String, +) -> Result(Int, pog.QueryError) { + use returned <- result.try( + pog.query( + "select coalesce(max(revision), -1) from factos_events where stream = $1", + ) + |> pog.parameter(pog.text(stream_name)) + |> pog.returning(int_field_decoder()) + |> pog.execute(on: connection), + ) + + case returned.rows { + [revision, ..] -> Ok(revision) + [] -> Ok(-1) + } +} + +fn has_matching_events_after( + connection: pog.Connection, + query: factos.Query, + after: factos.SequencePosition, +) -> Result(Bool, pog.QueryError) { + let after_position = case after { + factos.NoPosition -> -1 + factos.SequencePosition(position) -> position + } + pog.query("select type, tags from factos_events where position > $1") + |> pog.parameter(pog.int(after_position)) + |> pog.returning(query_match_decoder()) + |> pog.execute(on: connection) + |> result.map(fn(returned) { + use pair <- list.any(returned.rows) + let #(type_, tags) = pair + matches_query_parts(type_, tags, query) + }) +} + +fn query_match_decoder() -> decode.Decoder( + #(factos.EventType, List(factos.Tag)), +) { + use type_name <- decode.field(0, decode.string) + use tags <- decode.field(1, decode.string) + decode.success(#(factos.event_type(type_name), tags_from_text(tags))) +} + +fn int_field_decoder() -> decode.Decoder(Int) { + use value <- decode.field(0, decode.int) + decode.success(value) +} + +fn matches_query_parts( + type_: factos.EventType, + tags: List(factos.Tag), + query: factos.Query, +) -> Bool { + factos.matches_query( + factos.Recorded( + id: "", + stream: "", + revision: 0, + position: factos.NoPosition, + type_: type_, + tags: tags, + event: Nil, + ), + query, + ) +} + +fn stream_revision(events: List(factos.Recorded(event))) -> factos.Revision { + case list.reverse(events) { + [] -> factos.NoEvents + [event, ..] -> factos.CurrentRevision(event.revision) + } +} + +fn highest_recorded_position( + events: List(factos.Recorded(event)), +) -> factos.SequencePosition { + case list.reverse(events) { + [] -> factos.NoPosition + [event, ..] -> event.position + } +} + +fn expected_matches(expected: factos.Revision, current: Int) -> Bool { + case expected { + factos.NoEvents -> current == -1 + factos.CurrentRevision(revision) -> current == revision + } +} + +fn revision_to_int(revision: factos.Revision) -> Int { + case revision { + factos.NoEvents -> -1 + factos.CurrentRevision(revision) -> revision + } +} + +fn tags_to_text(tags: List(factos.Tag)) -> String { + tags + |> list.map(factos.tag_value) + |> string.join(with: "\n") +} + +fn tags_from_text(tags: String) -> List(factos.Tag) { + case string.is_empty(tags) { + True -> [] + False -> tags |> string.split(on: "\n") |> list.map(factos.tag) + } +} diff --git a/backends/factos_pog/test/factos_pog_test.gleam b/backends/factos_pog/test/factos_pog_test.gleam new file mode 100644 index 0000000..3d4bb01 --- /dev/null +++ b/backends/factos_pog/test/factos_pog_test.gleam @@ -0,0 +1,335 @@ +import factos +import factos/factos_pog +import gleam/bit_array +import gleam/erlang/process +import gleam/int +import gleam/list +import gleam/option.{Some} +import gleam/result +import gleeunit +import global_value +import pog + +pub fn main() -> Nil { + gleeunit.main() +} + +type Command { + RegisterUser(username: String) +} + +type Event { + UserRegistered(username: String) +} + +type State { + Available + Taken +} + +type DomainError { + AlreadyTaken +} + +type DecodeError { + UnknownEvent + InvalidData +} + +type CounterCommand { + Increment +} + +type CounterEvent { + Incremented(value: Int) +} + +type CounterState { + CounterState(total: Int) +} + +type TestGlobalData { + TestGlobalData(connection: pog.Connection) +} + +pub fn dispatch_stream_persists_events_test() { + let TestGlobalData(connection) = global_data() + reset_schema(connection) + + let assert Ok(factos_pog.Append(current_revision: 0, position: _)) = + factos_pog.dispatch_stream( + connection, + stream: "user-renata", + decider: decider(), + codec: codec(), + command: RegisterUser("renata"), + ) + + let assert Ok(loaded) = + factos_pog.load_stream( + connection, + stream: "user-renata", + decider: decider(), + codec: codec(), + ) + + assert loaded.state == Taken + assert loaded.revision == factos.CurrentRevision(0) +} + +pub fn dispatch_context_reads_by_event_type_and_tags_test() { + let TestGlobalData(connection) = global_data() + reset_schema(connection) + + let query = username_query("renata") + + let assert Ok(factos_pog.Append( + current_revision: 0, + position: factos.SequencePosition(_), + )) = + factos_pog.dispatch_context( + connection, + stream: "user-renata", + query: query, + decider: decider(), + codec: codec(), + command: RegisterUser("renata"), + ) + + let assert Ok(context) = + factos_pog.read_context( + connection, + query: query, + decider: decider(), + codec: codec(), + ) + + assert context.state == Taken + assert list.length(context.events) == 1 + assert context.position != factos.NoPosition +} + +pub fn dispatch_context_handles_many_streams_test() { + let TestGlobalData(connection) = global_data() + reset_schema(connection) + + let query = + factos.query([ + factos.query_item(types: [factos.event_type("Incremented")], tags: [ + factos.tag("counter:load"), + ]), + ]) + + let assert Ok(factos_pog.Append( + current_revision: 0, + position: factos.SequencePosition(_), + )) = dispatch_counter_context_many(connection, query, 25) + + let assert Ok(context) = + factos_pog.read_context( + connection, + query: query, + decider: counter_decider(), + codec: counter_codec(), + ) + + assert context.state == CounterState(25) + assert list.length(context.events) == 25 +} + +fn global_data() -> TestGlobalData { + global_value.create_with_unique_name("factos_pog_test.global.data", fn() { + TestGlobalData(connection: start_test_connection()) + }) +} + +fn start_test_connection() -> pog.Connection { + let pool_name = process.new_name("factos_pog_test") + let config = + pog.default_config(pool_name) + |> pog.host("127.0.0.1") + |> pog.port(5432) + |> pog.database("factos_pog") + |> pog.user("postgres") + |> pog.password(Some("postgres")) + |> pog.ssl(pog.SslDisabled) + + let assert Ok(_) = pog.start(config) + process.sleep(100) + pog.named_connection(pool_name) +} + +fn reset_schema(connection: pog.Connection) -> Nil { + let assert Ok(_) = + pog.query("drop table if exists factos_events") + |> pog.execute(on: connection) + let assert Ok(Nil) = factos_pog.migrate(connection) + Nil +} + +fn username_query(username: String) -> factos.Query { + factos.query([ + factos.query_item(types: [factos.event_type("UserRegistered")], tags: [ + factos.tag("username:" <> username), + ]), + ]) +} + +fn decider() -> factos.Decider(Command, State, Event, DomainError) { + factos.decider(initial: Available, decide:, evolve:) +} + +fn decide(state: State, command: Command) -> Result(List(Event), DomainError) { + case state, command { + Available, RegisterUser(username) -> Ok([UserRegistered(username)]) + Taken, RegisterUser(_) -> Error(AlreadyTaken) + } +} + +fn evolve(_state: State, _event: Event) -> State { + Taken +} + +fn codec() -> factos_pog.EventCodec(Event, DecodeError) { + factos_pog.EventCodec(encode:, decode:) +} + +fn encode(event: Event) -> factos_pog.Proposed(Event) { + factos_pog.Proposed( + id: "event-" <> event.username, + event: event, + type_: factos.event_type("UserRegistered"), + tags: [factos.tag("username:" <> event.username)], + data: bit_array.from_string(event.username), + ) +} + +fn decode( + stored: factos_pog.StoredEvent, +) -> Result(factos.Decoded(Event), DecodeError) { + case factos.event_type_name(stored.type_) { + "UserRegistered" -> { + use username <- result.try( + bit_array.to_string(stored.data) + |> result.replace_error(InvalidData), + ) + Ok(factos.Decoded( + event: UserRegistered(username), + type_: stored.type_, + tags: stored.tags, + )) + } + _ -> Error(UnknownEvent) + } +} + +fn dispatch_counter_context_many( + connection: pog.Connection, + query: factos.Query, + remaining: Int, +) -> Result(factos_pog.Append, factos_pog.Error(Nil, DecodeError)) { + case remaining { + 0 -> + factos_pog.dispatch_context( + connection, + stream: "counter-context-0", + query: query, + decider: counter_decider(), + codec: counter_codec(), + command: Increment, + ) + _ -> { + let stream_name = "counter-context-" <> int.to_string(remaining) + let result = + factos_pog.dispatch_context( + connection, + stream: stream_name, + query: query, + decider: counter_decider(), + codec: counter_codec(), + command: Increment, + ) + case remaining, result { + 1, _ -> result + _, Ok(_) -> + dispatch_counter_context_many(connection, query, remaining - 1) + _, Error(error) -> Error(error) + } + } + } +} + +fn counter_decider() -> factos.Decider( + CounterCommand, + CounterState, + CounterEvent, + Nil, +) { + factos.decider( + initial: CounterState(0), + decide: counter_decide, + evolve: counter_evolve, + ) +} + +fn counter_decide( + state: CounterState, + command: CounterCommand, +) -> Result(List(CounterEvent), Nil) { + let CounterState(total) = state + case command { + Increment -> Ok([Incremented(total + 1)]) + } +} + +fn counter_evolve(state: CounterState, event: CounterEvent) -> CounterState { + let CounterState(total) = state + case event { + Incremented(_) -> CounterState(total + 1) + } +} + +fn counter_codec() -> factos_pog.EventCodec(CounterEvent, DecodeError) { + factos_pog.EventCodec( + encode: encode_counter_event, + decode: decode_counter_event, + ) +} + +fn encode_counter_event( + event: CounterEvent, +) -> factos_pog.Proposed(CounterEvent) { + case event { + Incremented(value) -> + factos_pog.Proposed( + id: "counter-event-" <> int.to_string(value), + event: event, + type_: factos.event_type("Incremented"), + tags: [factos.tag("counter:load")], + data: bit_array.from_string(int.to_string(value)), + ) + } +} + +fn decode_counter_event( + stored: factos_pog.StoredEvent, +) -> Result(factos.Decoded(CounterEvent), DecodeError) { + case factos.event_type_name(stored.type_) { + "Incremented" -> { + use text <- result.try( + bit_array.to_string(stored.data) + |> result.replace_error(InvalidData), + ) + use value <- result.try( + int.parse(text) + |> result.replace_error(InvalidData), + ) + Ok(factos.Decoded( + event: Incremented(value), + type_: stored.type_, + tags: stored.tags, + )) + } + _ -> Error(UnknownEvent) + } +} diff --git a/backends/factos_sqlight/gleam.toml b/backends/factos_sqlight/gleam.toml index cea193f..e2faf16 100644 --- a/backends/factos_sqlight/gleam.toml +++ b/backends/factos_sqlight/gleam.toml @@ -1,5 +1,11 @@ name = "factos_sqlight" version = "1.0.0" +description = "SQLite backend for Factos context-first Event Sourcing using sqlight." +licences = ["Apache-2.0"] +links = [ + { title = "Factos", href = "https://factos.hexdocs.pm" }, + { title = "Simply Event Sourcing", href = "https://ricofritzsche.me/simply-event-sourcing/" }, +] [dependencies] factos = { path = "../.." } diff --git a/backends/factos_sqlight/src/factos/sqlight.gleam b/backends/factos_sqlight/src/factos/factos_sqlight.gleam similarity index 77% rename from backends/factos_sqlight/src/factos/sqlight.gleam rename to backends/factos_sqlight/src/factos/factos_sqlight.gleam index 722677d..290efc1 100644 --- a/backends/factos_sqlight/src/factos/sqlight.gleam +++ b/backends/factos_sqlight/src/factos/factos_sqlight.gleam @@ -1,4 +1,21 @@ //// SQLite backend for Factos using the `sqlight` package. +//// +//// This backend stores events in an append-only `factos_events` table and uses +//// SQLite transactions to implement both supported dispatch styles: +//// +//// 1. `dispatch_stream` protects one stream with a per-stream revision check. +//// 2. `dispatch_context` protects a command context with +//// `FailIfEventsMatch(query, after)`. +//// +//// The context flow is the important part for Command Context Consistency. The +//// command reads the facts selected by a `factos.Query`, folds them into a +//// temporary decision state, decides new facts, and appends those facts only when +//// no matching facts appeared after the observed position. SQLite can enforce +//// that condition transactionally because the context check and append happen in +//// the same database transaction. +//// +//// Event payload encoding is deliberately application-owned. The backend only +//// stores bytes plus query metadata (`EventType` and `Tag`). import factos import gleam/dynamic/decode @@ -8,6 +25,11 @@ import gleam/string import sqlight pub type Proposed(event) { + /// A domain event prepared for SQLite persistence. + /// + /// The application codec creates this value. `id` should identify the event for + /// the application. `type_` and `tags` are store-visible query metadata. `data` + /// is opaque bytes owned by the application codec. Proposed( id: String, event: event, @@ -18,6 +40,11 @@ pub type Proposed(event) { } pub type StoredEvent { + /// A raw event row read from SQLite before domain decoding. + /// + /// Decoders receive this value so they can inspect the stored event type, tags, + /// and bytes. `position` is the global append order. `revision` is the + /// per-stream revision. StoredEvent( position: Int, id: String, @@ -30,6 +57,11 @@ pub type StoredEvent { } pub type EventCodec(event, decode_error) { + /// Application-owned SQLite event codec. + /// + /// `encode` converts a domain event into bytes and metadata. `decode` converts a + /// stored row back into a `factos.Decoded` domain event. Decode failures are kept + /// in the application's own error type and wrapped as `DecodeError`. EventCodec( encode: fn(event) -> Proposed(event), decode: fn(StoredEvent) -> Result(factos.Decoded(event), decode_error), @@ -37,17 +69,35 @@ pub type EventCodec(event, decode_error) { } pub type Append { + /// Result of a successful append. + /// + /// `current_revision` is the latest revision of the target stream after the + /// append. `position` is the global position of the last inserted event, or + /// `NoPosition` when no events were produced. Append(current_revision: Int, position: factos.SequencePosition) } pub type Error(domain_error, decode_error) { + /// The decider rejected the command with a domain error. DomainError(domain_error) + + /// The application codec could not decode a stored event. DecodeError(decode_error) + + /// SQLite returned an error. StoreError(sqlight.Error) + + /// A stream revision or context append condition failed. AppendConditionFailed(factos.AppendCondition) } -pub fn migrate(connection: sqlight.Connection) -> Result(Nil, sqlight.Error) { +/// Create or update the SQLite schema required by this backend. +/// +/// The schema is an append-only `factos_events` table with a global autoincrement +/// `position`, per-stream `revision`, event `type`, newline-encoded `tags`, and +/// opaque `data` bytes. It also creates indexes for stream/revision reads and +/// position-based context checks. +pub fn migrate(connection: sqlight.Connection) -> Result(Nil, Error(_, _)) { sqlight.exec( " create table if not exists factos_events ( @@ -67,8 +117,15 @@ pub fn migrate(connection: sqlight.Connection) -> Result(Nil, sqlight.Error) { ", on: connection, ) + |> result.map_error(StoreError) } +/// Read and fold the facts selected by a command-context query. +/// +/// The backend reads stored rows, decodes them with the supplied codec, filters +/// them with `factos.matches_query`, folds matching events with the decider's +/// `evolve` function, and returns a `factos.Context` with a +/// `FailIfEventsMatch(query, after)` append condition. pub fn read_context( connection: sqlight.Connection, query query: factos.Query, @@ -93,6 +150,15 @@ pub fn read_context( )) } +/// Run a full context-first read-decide-append command flow. +/// +/// This function starts `BEGIN IMMEDIATE`, reads the query context, runs the +/// decider, verifies that no matching events appeared after the context position, +/// appends produced events to `stream`, and commits. Any error rolls the +/// transaction back. +/// +/// Use this when the command's real consistency boundary is the selected event +/// context rather than one predefined stream. pub fn dispatch_context( connection: sqlight.Connection, stream stream_name: String, @@ -131,6 +197,11 @@ pub fn dispatch_context( finish_transaction(connection, result) } +/// Load and fold one stream. +/// +/// This supports classic stream-revision consistency. The returned +/// `factos.LoadedStream` contains the folded state, decoded recorded events, and +/// current stream revision. pub fn load_stream( connection: sqlight.Connection, stream stream_name: String, @@ -155,6 +226,11 @@ pub fn load_stream( )) } +/// Run a stream-based read-decide-append command flow. +/// +/// The backend loads the target stream, folds it into state, runs the decider, and +/// appends produced events only if the stream revision still matches the loaded +/// revision. Use this when one stream is the intended consistency boundary. pub fn dispatch_stream( connection: sqlight.Connection, stream stream_name: String, @@ -423,18 +499,17 @@ fn current_revision( connection: sqlight.Connection, stream_name: String, ) -> Result(Int, sqlight.Error) { - sqlight.query( + use rows <- result.map(sqlight.query( "select coalesce(max(revision), -1) from factos_events where stream = ?", on: connection, with: [sqlight.text(stream_name)], expecting: int_field_decoder(), - ) - |> result.map(fn(rows) { - case rows { - [revision, ..] -> revision - [] -> -1 - } - }) + )) + + case rows { + [revision, ..] -> revision + [] -> -1 + } } fn has_matching_events_after( diff --git a/backends/factos_sqlight/test/factos_sqlight_test.gleam b/backends/factos_sqlight/test/factos_sqlight_test.gleam index 7198c09..4886bad 100644 --- a/backends/factos_sqlight/test/factos_sqlight_test.gleam +++ b/backends/factos_sqlight/test/factos_sqlight_test.gleam @@ -1,5 +1,5 @@ import factos -import factos/sqlight as factos_sqlight +import factos/factos_sqlight import gleam/bit_array import gleam/int import gleam/list @@ -122,7 +122,7 @@ pub fn dispatch_context_handles_many_streams_test() { } fn decider() -> factos.Decider(Command, State, Event, DomainError) { - factos.decider(initial: Available, decide: decide, evolve: evolve) + factos.decider(initial: Available, decide:, evolve:) } fn decide(state: State, command: Command) -> Result(List(Event), DomainError) { @@ -137,30 +137,27 @@ fn evolve(_state: State, _event: Event) -> State { } fn codec() -> factos_sqlight.EventCodec(Event, DecodeError) { - factos_sqlight.EventCodec(encode: encode, decode: decode_event) + factos_sqlight.EventCodec(encode:, decode:) } fn encode(event: Event) -> factos_sqlight.Proposed(Event) { - case event { - UserRegistered(username) -> - factos_sqlight.Proposed( - id: "event-" <> username, - event: event, - type_: factos.event_type("UserRegistered"), - tags: [factos.tag("username:" <> username)], - data: bit_array.from_string(username), - ) - } + factos_sqlight.Proposed( + id: "event-" <> event.username, + event: event, + type_: factos.event_type("UserRegistered"), + tags: [factos.tag("username:" <> event.username)], + data: bit_array.from_string(event.username), + ) } -fn decode_event( +fn decode( stored: factos_sqlight.StoredEvent, ) -> Result(factos.Decoded(Event), DecodeError) { case factos.event_type_name(stored.type_) { "UserRegistered" -> { use username <- result.try( bit_array.to_string(stored.data) - |> result.map_error(fn(_) { InvalidData }), + |> result.replace_error(InvalidData), ) Ok(factos.Decoded( event: UserRegistered(username), @@ -298,11 +295,11 @@ fn decode_counter_event( "Incremented" -> { use text <- result.try( bit_array.to_string(stored.data) - |> result.map_error(fn(_) { InvalidData }), + |> result.replace_error(InvalidData), ) use value <- result.try( int.parse(text) - |> result.map_error(fn(_) { InvalidData }), + |> result.replace_error(InvalidData), ) Ok(factos.Decoded( event: Incremented(value), diff --git a/examples/orders_sqlight/gleam.toml b/examples/orders_sqlight/gleam.toml new file mode 100644 index 0000000..eda8d63 --- /dev/null +++ b/examples/orders_sqlight/gleam.toml @@ -0,0 +1,13 @@ +name = "orders_sqlight" +version = "1.0.0" + +[dependencies] +factos = { path = "../.." } +factos_sqlight = { path = "../../backends/factos_sqlight" } +gleam_erlang = ">= 1.0.0 and < 2.0.0" +gleam_stdlib = ">= 1.0.0 and < 2.0.0" +simplifile = ">= 2.0.0 and < 3.0.0" +sqlight = ">= 1.1.0 and < 2.0.0" + +[dev_dependencies] +gleeunit = ">= 1.0.0 and < 2.0.0" diff --git a/examples/orders_sqlight/manifest.toml b/examples/orders_sqlight/manifest.toml new file mode 100644 index 0000000..649b039 --- /dev/null +++ b/examples/orders_sqlight/manifest.toml @@ -0,0 +1,28 @@ +# Do not manually edit this file, it is managed by Gleam. +# +# This file locks the dependency versions used, to make your build +# deterministic and to prevent unexpected versions from being included +# in your application. +# +# You should check this file into your source control repository. + +packages = [ + { name = "esqlite", version = "0.9.0", build_tools = ["rebar3"], requirements = [], otp_app = "esqlite", source = "hex", outer_checksum = "CCF72258A4EE152EC7AD92AA9A03552EB6CA1B06B65C93AD5B6E55C302E05855" }, + { name = "factos", version = "1.0.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], source = "local", path = "../.." }, + { name = "factos_sqlight", version = "1.0.0", build_tools = ["gleam"], requirements = ["factos", "gleam_stdlib", "sqlight"], source = "local", path = "../../backends/factos_sqlight" }, + { name = "filepath", version = "1.1.2", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "filepath", source = "hex", outer_checksum = "B06A9AF0BF10E51401D64B98E4B627F1D2E48C154967DA7AF4D0914780A6D40A" }, + { name = "gleam_erlang", version = "1.3.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_erlang", source = "hex", outer_checksum = "1124AD3AA21143E5AF0FC5CF3D9529F6DB8CA03E43A55711B60B6B7B3874375C" }, + { name = "gleam_stdlib", version = "1.0.3", build_tools = ["gleam"], requirements = [], otp_app = "gleam_stdlib", source = "hex", outer_checksum = "1F543AFBA5D33DA493E6087F4E4C4F20D899411343512686C98A8ABB2963CF22" }, + { name = "gleeunit", version = "1.11.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleeunit", source = "hex", outer_checksum = "EC31ABA74256AEA531EDF8169931D775BBB384FED0A8A1BDC4DD9354E3E21826" }, + { name = "simplifile", version = "2.4.0", build_tools = ["gleam"], requirements = ["filepath", "gleam_stdlib"], otp_app = "simplifile", source = "hex", outer_checksum = "7C18AFA4FED0B4CE1FA5B0B4BAC1FA1744427054EA993565F6F3F82E5453170D" }, + { name = "sqlight", version = "1.1.0", build_tools = ["gleam"], requirements = ["esqlite", "gleam_stdlib"], otp_app = "sqlight", source = "hex", outer_checksum = "ECA1A4B45C35EB9EFCEEB7FAAC7BF5D8B2C777A7C1FC8A9C12CB67D54CED42E7" }, +] + +[requirements] +factos = { path = "../.." } +factos_sqlight = { path = "../../backends/factos_sqlight" } +gleam_erlang = { version = ">= 1.0.0 and < 2.0.0" } +gleam_stdlib = { version = ">= 1.0.0 and < 2.0.0" } +gleeunit = { version = ">= 1.0.0 and < 2.0.0" } +simplifile = { version = ">= 2.0.0 and < 3.0.0" } +sqlight = { version = ">= 1.1.0 and < 2.0.0" } diff --git a/examples/orders_sqlight/src/order_workflow.gleam b/examples/orders_sqlight/src/order_workflow.gleam new file mode 100644 index 0000000..1b56b81 --- /dev/null +++ b/examples/orders_sqlight/src/order_workflow.gleam @@ -0,0 +1,866 @@ +import factos +import factos/factos_sqlight +import gleam/bit_array +import gleam/erlang/process +import gleam/int +import gleam/io +import gleam/list +import gleam/result +import gleam/string +import simplifile +import sqlight + +const order_count = 40 + +const concurrency = 8 + +const receive_timeout = 30_000 + +const write_retries = 200 + +pub type Command { + OpenOrder(table: Int) + AddItem(sku: String, name: String, price: Int) + RemoveItem(sku: String) + SubmitOrder + StartPreparing + MarkReady + Serve + Pay(amount: Int) + CancelOrder(reason: String) +} + +pub type Event { + OrderOpened(table: Int) + ItemAdded(sku: String, name: String, price: Int) + ItemRemoved(sku: String) + OrderSubmitted + PreparationStarted + OrderMarkedReady + OrderServed + PaymentReceived(amount: Int) + OrderCancelled(reason: String) +} + +pub type Item { + Item(sku: String, name: String, price: Int) +} + +pub type State { + NoOrder + Draft(table: Int, items: List(Item)) + Submitted(table: Int, items: List(Item)) + Preparing(table: Int, items: List(Item)) + Ready(table: Int, items: List(Item)) + Served(table: Int, items: List(Item)) + Paid(table: Int, items: List(Item), amount: Int) + Cancelled(reason: String) +} + +pub type DomainError { + OrderAlreadyOpen + OrderNotOpen + DuplicateItem(sku: String) + ItemNotFound(sku: String) + EmptyOrder + PaymentTooLow(required: Int, paid: Int) + WorkerTimedOut(remaining: Int) + InvalidTransition(action: String, state: String) +} + +pub type DecodeError { + UnknownEventType(String) + InvalidPayload(String) +} + +pub type KitchenSummary { + KitchenSummary( + opened: Int, + submitted: Int, + preparing: Int, + ready: Int, + served: Int, + paid: Int, + cancelled: Int, + revenue: Int, + ) +} + +pub type ExampleResult { + ExampleResult( + order_id: String, + final_state: State, + kitchen_summary: KitchenSummary, + recorded_events: Int, + ) +} + +pub type StressResult { + StressResult( + orders: Int, + paid_orders: Int, + cancelled_orders: Int, + recorded_events: Int, + revenue: Int, + ) +} + +type WorkerResult { + WorkerResult(final_state: State, recorded_events: Int, revenue: Int) +} + +type WorkerMessage { + WorkerFinished( + order_number: Int, + result: Result(WorkerResult, factos_sqlight.Error(DomainError, DecodeError)), + ) +} + +pub fn main() -> Nil { + case run() { + Ok(result) -> + io.println( + "restaurant stress workflow completed: " + <> int.to_string(result.orders) + <> " orders, " + <> int.to_string(result.recorded_events) + <> " events", + ) + Error(_) -> io.println("restaurant stress workflow failed") + } +} + +pub fn run() -> Result( + StressResult, + factos_sqlight.Error(DomainError, DecodeError), +) { + let database_path = "/tmp/factos_examples_stress.sqlite3" + log("reset database " <> database_path) + let _ = simplifile.delete_file(database_path) + + log("prepare database") + use _ <- result.try(prepare_database(database_path)) + + let workers = process.new_subject() + + let _ = + int.range(from: 1, to: concurrency + 1, with: Nil, run: fn(_, order_number) { + spawn_order(workers, database_path, order_number) + Nil + }) + + collect_workers( + workers, + database_path: database_path, + remaining: order_count, + next_order: concurrency + 1, + summary: StressResult(0, 0, 0, 0, 0), + ) +} + +fn spawn_order( + workers: process.Subject(WorkerMessage), + database_path: String, + order_number: Int, +) -> Nil { + log("spawn order " <> int.to_string(order_number)) + let _ = + process.spawn(fn() { + log("order " <> int.to_string(order_number) <> " started") + let result = run_order(database_path, order_number) + log( + "order " + <> int.to_string(order_number) + <> " finished with " + <> result_to_string(result), + ) + process.send(workers, WorkerFinished(order_number, result)) + }) + Nil +} + +fn prepare_database( + database_path: String, +) -> Result(Nil, factos_sqlight.Error(DomainError, DecodeError)) { + use connection <- sqlight.with_connection(database_path) + use _ <- result.try(configure_connection(connection)) + factos_sqlight.migrate(connection) +} + +fn configure_connection( + connection: sqlight.Connection, +) -> Result(Nil, factos_sqlight.Error(DomainError, DecodeError)) { + sqlight.exec( + "pragma journal_mode = wal; pragma busy_timeout = 50", + on: connection, + ) + |> result.map_error(factos_sqlight.StoreError) +} + +fn run_order( + database_path: String, + order_number: Int, +) -> Result(WorkerResult, factos_sqlight.Error(DomainError, DecodeError)) { + use connection <- sqlight.with_connection(database_path) + use _ <- result.try(configure_connection(connection)) + + let order_id = "stress-" <> int.to_string(order_number) + use _ <- result.try(dispatch_commands( + connection, + order_number, + order_id, + workflow(order_number), + )) + + log("order " <> int.to_string(order_number) <> " loading stream") + use loaded <- result.try(factos_sqlight.load_stream( + connection, + stream: order_stream(order_id), + decider: order_decider(), + codec: order_codec(), + )) + + Ok(WorkerResult( + final_state: loaded.state, + recorded_events: list.length(loaded.events), + revenue: revenue(loaded.state), + )) +} + +fn collect_workers( + workers: process.Subject(WorkerMessage), + database_path database_path: String, + remaining remaining: Int, + next_order next_order: Int, + summary summary: StressResult, +) -> Result(StressResult, factos_sqlight.Error(DomainError, DecodeError)) { + case remaining { + 0 -> Ok(summary) + _ -> + case process.receive(workers, within: receive_timeout) { + Ok(WorkerFinished(order_number, Ok(result))) -> { + log("collector received order " <> int.to_string(order_number)) + case next_order <= order_count { + True -> spawn_order(workers, database_path, next_order) + False -> Nil + } + collect_workers( + workers, + database_path: database_path, + remaining: remaining - 1, + next_order: next_order + 1, + summary: add_worker_result(summary, result), + ) + } + Ok(WorkerFinished(order_number, Error(error))) -> { + log( + "collector received error from order " + <> int.to_string(order_number) + <> ": " + <> store_error_to_string(error), + ) + Error(error) + } + Error(Nil) -> { + log( + "collector timed out with " + <> int.to_string(remaining) + <> " remaining", + ) + Error(factos_sqlight.DomainError(WorkerTimedOut(remaining))) + } + } + } +} + +fn add_worker_result( + summary: StressResult, + result: WorkerResult, +) -> StressResult { + let #(paid_orders, cancelled_orders) = case result.final_state { + Paid(_, _, _) -> #(summary.paid_orders + 1, summary.cancelled_orders) + Cancelled(_) -> #(summary.paid_orders, summary.cancelled_orders + 1) + NoOrder + | Draft(_, _) + | Submitted(_, _) + | Preparing(_, _) + | Ready(_, _) + | Served(_, _) -> #(summary.paid_orders, summary.cancelled_orders) + } + + StressResult( + orders: summary.orders + 1, + paid_orders: paid_orders, + cancelled_orders: cancelled_orders, + recorded_events: summary.recorded_events + result.recorded_events, + revenue: summary.revenue + result.revenue, + ) +} + +fn workflow(order_number: Int) -> List(Command) { + let draft_commands = [ + OpenOrder(table: order_number), + AddItem(sku: "burger", name: "House Burger", price: 16), + AddItem(sku: "fries", name: "Fries", price: 6), + AddItem(sku: "shake", name: "Vanilla Shake", price: 8), + RemoveItem(sku: "shake"), + SubmitOrder, + ] + + case should_cancel(order_number) { + True -> + list.append(draft_commands, [ + CancelOrder(reason: "guest left before kitchen started"), + ]) + False -> + list.append(draft_commands, [ + StartPreparing, + MarkReady, + Serve, + Pay(amount: 25), + ]) + } +} + +fn should_cancel(order_number: Int) -> Bool { + int.modulo(order_number, by: 5) == Ok(0) +} + +fn dispatch_commands( + connection: sqlight.Connection, + order_number: Int, + order_id: String, + commands: List(Command), +) -> Result(Nil, factos_sqlight.Error(DomainError, DecodeError)) { + case commands { + [] -> Ok(Nil) + [command, ..rest] -> { + log( + "order " + <> int.to_string(order_number) + <> " dispatch " + <> command_to_string(command), + ) + use _ <- result.try(dispatch_with_retry( + connection, + order_number, + order_id, + command, + attempts: write_retries, + )) + dispatch_commands(connection, order_number, order_id, rest) + } + } +} + +fn dispatch_with_retry( + connection: sqlight.Connection, + order_number: Int, + order_id: String, + command: Command, + attempts attempts: Int, +) -> Result( + factos_sqlight.Append, + factos_sqlight.Error(DomainError, DecodeError), +) { + let result = dispatch(connection, order_id, command) + + case attempts > 0, result { + _, Ok(append) -> Ok(append) + True, Error(factos_sqlight.StoreError(_)) -> { + log( + "order " + <> int.to_string(order_number) + <> " retry " + <> command_to_string(command) + <> " after SQLite store error; attempts left " + <> int.to_string(attempts - 1), + ) + process.sleep(retry_delay(order_number, attempts)) + dispatch_with_retry( + connection, + order_number, + order_id, + command, + attempts: attempts - 1, + ) + } + _, Error(error) -> { + log( + "order " + <> int.to_string(order_number) + <> " failed " + <> command_to_string(command) + <> " with " + <> store_error_to_string(error), + ) + Error(error) + } + } +} + +fn retry_delay(order_number: Int, attempts: Int) -> Int { + case int.modulo(order_number + attempts, by: 10) { + Ok(offset) -> 5 + offset + Error(Nil) -> 5 + } +} + +fn dispatch( + connection: sqlight.Connection, + order_id: String, + command: Command, +) -> Result( + factos_sqlight.Append, + factos_sqlight.Error(DomainError, DecodeError), +) { + factos_sqlight.dispatch_stream( + connection, + stream: order_stream(order_id), + decider: order_decider(), + codec: order_codec(), + command:, + ) +} + +fn order_stream(order_id: String) -> String { + "restaurant-order-" <> order_id +} + +fn revenue(state: State) -> Int { + case state { + Paid(_, _, amount) -> amount + NoOrder + | Draft(_, _) + | Submitted(_, _) + | Preparing(_, _) + | Ready(_, _) + | Served(_, _) + | Cancelled(_) -> 0 + } +} + +pub fn order_decider() -> factos.Decider(Command, State, Event, DomainError) { + factos.decider(initial: NoOrder, decide:, evolve:) +} + +fn decide(state: State, command: Command) -> Result(List(Event), DomainError) { + case state, command { + NoOrder, OpenOrder(table) -> Ok([OrderOpened(table)]) + NoOrder, _ -> Error(OrderNotOpen) + + Draft(_, _), OpenOrder(_) -> Error(OrderAlreadyOpen) + Draft(_, items), AddItem(sku, name, price) -> + case has_item(items, sku) { + True -> Error(DuplicateItem(sku)) + False -> Ok([ItemAdded(sku, name, price)]) + } + Draft(_, items), RemoveItem(sku) -> + case has_item(items, sku) { + True -> Ok([ItemRemoved(sku)]) + False -> Error(ItemNotFound(sku)) + } + Draft(_, items), SubmitOrder -> + case list.is_empty(items) { + True -> Error(EmptyOrder) + False -> Ok([OrderSubmitted]) + } + Draft(_, _), CancelOrder(reason) -> Ok([OrderCancelled(reason)]) + Draft(_, _), StartPreparing + | Draft(_, _), MarkReady + | Draft(_, _), Serve + | Draft(_, _), Pay(_) + -> invalid(command, state) + + Submitted(_, _), StartPreparing -> Ok([PreparationStarted]) + Submitted(_, _), CancelOrder(reason) -> Ok([OrderCancelled(reason)]) + Submitted(_, _), OpenOrder(_) + | Submitted(_, _), AddItem(_, _, _) + | Submitted(_, _), RemoveItem(_) + | Submitted(_, _), SubmitOrder + | Submitted(_, _), MarkReady + | Submitted(_, _), Serve + | Submitted(_, _), Pay(_) + -> invalid(command, state) + + Preparing(_, _), MarkReady -> Ok([OrderMarkedReady]) + Preparing(_, _), OpenOrder(_) + | Preparing(_, _), AddItem(_, _, _) + | Preparing(_, _), RemoveItem(_) + | Preparing(_, _), SubmitOrder + | Preparing(_, _), StartPreparing + | Preparing(_, _), Serve + | Preparing(_, _), Pay(_) + | Preparing(_, _), CancelOrder(_) + -> invalid(command, state) + + Ready(_, _), Serve -> Ok([OrderServed]) + Ready(_, _), OpenOrder(_) + | Ready(_, _), AddItem(_, _, _) + | Ready(_, _), RemoveItem(_) + | Ready(_, _), SubmitOrder + | Ready(_, _), StartPreparing + | Ready(_, _), MarkReady + | Ready(_, _), Pay(_) + | Ready(_, _), CancelOrder(_) + -> invalid(command, state) + + Served(_, items), Pay(amount) -> { + let required = total(items) + case amount >= required { + True -> Ok([PaymentReceived(amount)]) + False -> Error(PaymentTooLow(required: required, paid: amount)) + } + } + Served(_, _), OpenOrder(_) + | Served(_, _), AddItem(_, _, _) + | Served(_, _), RemoveItem(_) + | Served(_, _), SubmitOrder + | Served(_, _), StartPreparing + | Served(_, _), MarkReady + | Served(_, _), Serve + | Served(_, _), CancelOrder(_) + -> invalid(command, state) + + Paid(_, _, _), _ -> invalid(command, state) + Cancelled(_), _ -> invalid(command, state) + } +} + +fn evolve(state: State, event: Event) -> State { + case event { + OrderOpened(table) -> Draft(table: table, items: []) + ItemAdded(sku, name, price) -> + add_item_to_state(state, Item(sku, name, price)) + ItemRemoved(sku) -> remove_item_from_state(state, sku) + OrderSubmitted -> move_to_submitted(state) + PreparationStarted -> move_to_preparing(state) + OrderMarkedReady -> move_to_ready(state) + OrderServed -> move_to_served(state) + PaymentReceived(amount) -> move_to_paid(state, amount) + OrderCancelled(reason) -> Cancelled(reason) + } +} + +fn add_item_to_state(state: State, item: Item) -> State { + case state { + Draft(table, items) -> Draft(table: table, items: [item, ..items]) + NoOrder + | Submitted(_, _) + | Preparing(_, _) + | Ready(_, _) + | Served(_, _) + | Paid(_, _, _) + | Cancelled(_) -> state + } +} + +fn remove_item_from_state(state: State, sku: String) -> State { + case state { + Draft(table, items) -> + Draft( + table: table, + items: list.filter(items, fn(item) { item.sku != sku }), + ) + NoOrder + | Submitted(_, _) + | Preparing(_, _) + | Ready(_, _) + | Served(_, _) + | Paid(_, _, _) + | Cancelled(_) -> state + } +} + +fn move_to_submitted(state: State) -> State { + case state { + Draft(table, items) -> Submitted(table: table, items: items) + NoOrder + | Submitted(_, _) + | Preparing(_, _) + | Ready(_, _) + | Served(_, _) + | Paid(_, _, _) + | Cancelled(_) -> state + } +} + +fn move_to_preparing(state: State) -> State { + case state { + Submitted(table, items) -> Preparing(table: table, items: items) + NoOrder + | Draft(_, _) + | Preparing(_, _) + | Ready(_, _) + | Served(_, _) + | Paid(_, _, _) + | Cancelled(_) -> state + } +} + +fn move_to_ready(state: State) -> State { + case state { + Preparing(table, items) -> Ready(table: table, items: items) + NoOrder + | Draft(_, _) + | Submitted(_, _) + | Ready(_, _) + | Served(_, _) + | Paid(_, _, _) + | Cancelled(_) -> state + } +} + +fn move_to_served(state: State) -> State { + case state { + Ready(table, items) -> Served(table: table, items: items) + NoOrder + | Draft(_, _) + | Submitted(_, _) + | Preparing(_, _) + | Served(_, _) + | Paid(_, _, _) + | Cancelled(_) -> state + } +} + +fn move_to_paid(state: State, amount: Int) -> State { + case state { + Served(table, items) -> Paid(table: table, items: items, amount: amount) + NoOrder + | Draft(_, _) + | Submitted(_, _) + | Preparing(_, _) + | Ready(_, _) + | Paid(_, _, _) + | Cancelled(_) -> state + } +} + +pub fn kitchen_summary_view() -> factos.View(KitchenSummary, Event) { + factos.view( + initial: KitchenSummary(0, 0, 0, 0, 0, 0, 0, 0), + evolve: evolve_summary, + ) +} + +fn evolve_summary(summary: KitchenSummary, event: Event) -> KitchenSummary { + case event { + OrderOpened(_) -> KitchenSummary(..summary, opened: summary.opened + 1) + ItemAdded(_, _, _) | ItemRemoved(_) -> summary + OrderSubmitted -> + KitchenSummary(..summary, submitted: summary.submitted + 1) + PreparationStarted -> + KitchenSummary(..summary, preparing: summary.preparing + 1) + OrderMarkedReady -> KitchenSummary(..summary, ready: summary.ready + 1) + OrderServed -> KitchenSummary(..summary, served: summary.served + 1) + PaymentReceived(amount) -> + KitchenSummary( + ..summary, + paid: summary.paid + 1, + revenue: summary.revenue + amount, + ) + OrderCancelled(_) -> + KitchenSummary(..summary, cancelled: summary.cancelled + 1) + } +} + +pub fn order_codec() -> factos_sqlight.EventCodec(Event, DecodeError) { + factos_sqlight.EventCodec(encode: encode_event, decode: decode_event) +} + +fn encode_event(event: Event) -> factos_sqlight.Proposed(Event) { + case event { + OrderOpened(table) -> + proposed(event, "OrderOpened", [int.to_string(table)], []) + ItemAdded(sku, name, price) -> + proposed(event, "ItemAdded", [sku, name, int.to_string(price)], [ + factos.tag("sku:" <> sku), + ]) + ItemRemoved(sku) -> + proposed(event, "ItemRemoved", [sku], [ + factos.tag("sku:" <> sku), + ]) + OrderSubmitted -> proposed(event, "OrderSubmitted", [], []) + PreparationStarted -> proposed(event, "PreparationStarted", [], []) + OrderMarkedReady -> proposed(event, "OrderMarkedReady", [], []) + OrderServed -> proposed(event, "OrderServed", [], []) + PaymentReceived(amount) -> + proposed(event, "PaymentReceived", [int.to_string(amount)], []) + OrderCancelled(reason) -> proposed(event, "OrderCancelled", [reason], []) + } +} + +fn proposed( + event: Event, + type_name: String, + fields: List(String), + tags: List(factos.Tag), +) -> factos_sqlight.Proposed(Event) { + factos_sqlight.Proposed( + id: "example-" <> type_name <> "-" <> fields_to_payload(fields), + event: event, + type_: factos.event_type(type_name), + tags: [factos.tag("restaurant"), ..tags], + data: bit_array.from_string(fields_to_payload(fields)), + ) +} + +fn decode_event( + stored: factos_sqlight.StoredEvent, +) -> Result(factos.Decoded(Event), DecodeError) { + let type_name = factos.event_type_name(stored.type_) + let fields = + stored.data + |> bit_array.to_string + |> result.replace_error(InvalidPayload(type_name)) + |> result.map(payload_to_fields) + + use event <- result.try(decode_fields(type_name, fields)) + Ok(factos.Decoded(event: event, type_: stored.type_, tags: stored.tags)) +} + +fn decode_fields( + type_name: String, + fields_result: Result(List(String), DecodeError), +) -> Result(Event, DecodeError) { + use fields <- result.try(fields_result) + case type_name, fields { + "OrderOpened", [table] -> { + use table <- result.try(parse_int(table, type_name)) + Ok(OrderOpened(table)) + } + "ItemAdded", [sku, name, price] -> { + use price <- result.try(parse_int(price, type_name)) + Ok(ItemAdded(sku, name, price)) + } + "ItemRemoved", [sku] -> Ok(ItemRemoved(sku)) + "OrderSubmitted", [] -> Ok(OrderSubmitted) + "PreparationStarted", [] -> Ok(PreparationStarted) + "OrderMarkedReady", [] -> Ok(OrderMarkedReady) + "OrderServed", [] -> Ok(OrderServed) + "PaymentReceived", [amount] -> { + use amount <- result.try(parse_int(amount, type_name)) + Ok(PaymentReceived(amount)) + } + "OrderCancelled", [reason] -> Ok(OrderCancelled(reason)) + _, _ -> Error(UnknownEventType(type_name)) + } +} + +fn parse_int(value: String, type_name: String) -> Result(Int, DecodeError) { + int.parse(value) + |> result.replace_error(InvalidPayload(type_name)) +} + +fn fields_to_payload(fields: List(String)) -> String { + string.join(fields, with: "|") +} + +fn payload_to_fields(payload: String) -> List(String) { + case string.is_empty(payload) { + True -> [] + False -> string.split(payload, on: "|") + } +} + +fn has_item(items: List(Item), sku: String) -> Bool { + list.any(items, fn(item) { item.sku == sku }) +} + +fn total(items: List(Item)) -> Int { + list.fold(items, 0, fn(total, item) { total + item.price }) +} + +fn invalid(command: Command, state: State) -> Result(List(Event), DomainError) { + Error(InvalidTransition(command_to_string(command), state_to_string(state))) +} + +fn command_to_string(command: Command) -> String { + case command { + OpenOrder(_) -> "OpenOrder" + AddItem(_, _, _) -> "AddItem" + RemoveItem(_) -> "RemoveItem" + SubmitOrder -> "SubmitOrder" + StartPreparing -> "StartPreparing" + MarkReady -> "MarkReady" + Serve -> "Serve" + Pay(_) -> "Pay" + CancelOrder(_) -> "CancelOrder" + } +} + +fn result_to_string( + result: Result(WorkerResult, factos_sqlight.Error(DomainError, DecodeError)), +) -> String { + case result { + Ok(worker_result) -> + "ok " + <> state_to_string(worker_result.final_state) + <> " events=" + <> int.to_string(worker_result.recorded_events) + Error(error) -> "error " <> store_error_to_string(error) + } +} + +fn store_error_to_string( + error: factos_sqlight.Error(DomainError, DecodeError), +) -> String { + case error { + factos_sqlight.DomainError(error) -> + "domain:" <> domain_error_to_string(error) + factos_sqlight.DecodeError(error) -> + "decode:" <> decode_error_to_string(error) + factos_sqlight.StoreError(sqlight.SqlightError(code, message, _)) -> + "sqlite(code=" + <> int.to_string(sqlight.error_code_to_int(code)) + <> ", message=" + <> message + <> ")" + factos_sqlight.AppendConditionFailed(_) -> "append-condition-failed" + } +} + +fn domain_error_to_string(error: DomainError) -> String { + case error { + OrderAlreadyOpen -> "OrderAlreadyOpen" + OrderNotOpen -> "OrderNotOpen" + DuplicateItem(sku) -> "DuplicateItem(" <> sku <> ")" + ItemNotFound(sku) -> "ItemNotFound(" <> sku <> ")" + EmptyOrder -> "EmptyOrder" + PaymentTooLow(required, paid) -> + "PaymentTooLow(required=" + <> int.to_string(required) + <> ", paid=" + <> int.to_string(paid) + <> ")" + WorkerTimedOut(remaining) -> + "WorkerTimedOut(remaining=" <> int.to_string(remaining) <> ")" + InvalidTransition(action, state) -> + "InvalidTransition(" <> action <> ", " <> state <> ")" + } +} + +fn decode_error_to_string(error: DecodeError) -> String { + case error { + UnknownEventType(type_name) -> "UnknownEventType(" <> type_name <> ")" + InvalidPayload(type_name) -> "InvalidPayload(" <> type_name <> ")" + } +} + +fn log(message: String) -> Nil { + io.println("[factos-example] " <> message) +} + +fn state_to_string(state: State) -> String { + case state { + NoOrder -> "NoOrder" + Draft(_, _) -> "Draft" + Submitted(_, _) -> "Submitted" + Preparing(_, _) -> "Preparing" + Ready(_, _) -> "Ready" + Served(_, _) -> "Served" + Paid(_, _, _) -> "Paid" + Cancelled(_) -> "Cancelled" + } +} diff --git a/examples/orders_sqlight/src/orders_sqlight.gleam b/examples/orders_sqlight/src/orders_sqlight.gleam new file mode 100644 index 0000000..09c68b8 --- /dev/null +++ b/examples/orders_sqlight/src/orders_sqlight.gleam @@ -0,0 +1,5 @@ +import order_workflow + +pub fn main() -> Nil { + order_workflow.main() +} diff --git a/examples/orders_sqlight/test/orders_sqlight_test.gleam b/examples/orders_sqlight/test/orders_sqlight_test.gleam new file mode 100644 index 0000000..be92879 --- /dev/null +++ b/examples/orders_sqlight/test/orders_sqlight_test.gleam @@ -0,0 +1,16 @@ +import gleeunit +import order_workflow + +pub fn main() -> Nil { + gleeunit.main() +} + +pub fn restaurant_order_example_runs_concurrently_under_stress_test() { + let assert Ok(order_workflow.StressResult( + orders: 40, + paid_orders: 32, + cancelled_orders: 8, + recorded_events: 376, + revenue: 800, + )) = order_workflow.run() +} diff --git a/examples/tickets_pog/compose.yml b/examples/tickets_pog/compose.yml new file mode 100644 index 0000000..a6570a6 --- /dev/null +++ b/examples/tickets_pog/compose.yml @@ -0,0 +1,14 @@ +services: + postgres: + image: postgres:18 + environment: + POSTGRES_DB: tickets_pog + POSTGRES_USER: postgres + POSTGRES_PASSWORD: postgres + ports: + - "5433:5432" + healthcheck: + test: ["CMD-SHELL", "pg_isready -U postgres -d tickets_pog"] + interval: 1s + timeout: 5s + retries: 20 diff --git a/examples/tickets_pog/gleam.toml b/examples/tickets_pog/gleam.toml new file mode 100644 index 0000000..85a5bc5 --- /dev/null +++ b/examples/tickets_pog/gleam.toml @@ -0,0 +1,13 @@ +name = "tickets_pog" +version = "1.0.0" + +[dependencies] +factos = { path = "../.." } +factos_pog = { path = "../../backends/factos_pog" } +gleam_erlang = ">= 1.0.0 and < 2.0.0" +gleam_stdlib = ">= 1.0.0 and < 2.0.0" +global_value = ">= 1.0.0 and < 2.0.0" +pog = ">= 4.1.0 and < 5.0.0" + +[dev_dependencies] +gleeunit = ">= 1.0.0 and < 2.0.0" diff --git a/examples/tickets_pog/manifest.toml b/examples/tickets_pog/manifest.toml new file mode 100644 index 0000000..964e5cc --- /dev/null +++ b/examples/tickets_pog/manifest.toml @@ -0,0 +1,33 @@ +# Do not manually edit this file, it is managed by Gleam. +# +# This file locks the dependency versions used, to make your build +# deterministic and to prevent unexpected versions from being included +# in your application. +# +# You should check this file into your source control repository. + +packages = [ + { name = "backoff", version = "1.1.6", build_tools = ["rebar3"], requirements = [], otp_app = "backoff", source = "hex", outer_checksum = "CF0CFFF8995FB20562F822E5CC47D8CCF664C5ECDC26A684CBE85C225F9D7C39" }, + { name = "exception", version = "2.1.1", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "exception", source = "hex", outer_checksum = "6BDEA95248093599391C3B5DF1835C5C6A86C353C2F99CE539B450E3432FE117" }, + { name = "factos", version = "1.0.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], source = "local", path = "../.." }, + { name = "factos_pog", version = "1.0.0", build_tools = ["gleam"], requirements = ["factos", "gleam_stdlib", "pog"], source = "local", path = "../../backends/factos_pog" }, + { name = "gleam_erlang", version = "1.3.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_erlang", source = "hex", outer_checksum = "1124AD3AA21143E5AF0FC5CF3D9529F6DB8CA03E43A55711B60B6B7B3874375C" }, + { name = "gleam_otp", version = "1.2.0", build_tools = ["gleam"], requirements = ["gleam_erlang", "gleam_stdlib"], otp_app = "gleam_otp", source = "hex", outer_checksum = "BA6A294E295E428EC1562DC1C11EA7530DCB981E8359134BEABC8493B7B2258E" }, + { name = "gleam_stdlib", version = "1.0.3", build_tools = ["gleam"], requirements = [], otp_app = "gleam_stdlib", source = "hex", outer_checksum = "1F543AFBA5D33DA493E6087F4E4C4F20D899411343512686C98A8ABB2963CF22" }, + { name = "gleam_time", version = "1.8.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_time", source = "hex", outer_checksum = "533D8723774D61AD4998324F5DD1DABDCDBFABAFB9E87CB5D03C6955448FC97D" }, + { name = "gleeunit", version = "1.11.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleeunit", source = "hex", outer_checksum = "EC31ABA74256AEA531EDF8169931D775BBB384FED0A8A1BDC4DD9354E3E21826" }, + { 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" }, + { name = "pog", version = "4.1.0", build_tools = ["gleam"], requirements = ["exception", "gleam_erlang", "gleam_otp", "gleam_stdlib", "gleam_time", "pgo"], otp_app = "pog", source = "hex", outer_checksum = "E4AFBA39A5FAA2E77291836C9683ADE882E65A06AB28CA7D61AE7A3AD61EBBD5" }, +] + +[requirements] +factos = { path = "../.." } +factos_pog = { path = "../../backends/factos_pog" } +gleam_erlang = { version = ">= 1.0.0 and < 2.0.0" } +gleam_stdlib = { version = ">= 1.0.0 and < 2.0.0" } +gleeunit = { version = ">= 1.0.0 and < 2.0.0" } +global_value = { version = ">= 1.0.0 and < 2.0.0" } +pog = { version = ">= 4.1.0 and < 5.0.0" } diff --git a/examples/tickets_pog/src/ticket_sale.gleam b/examples/tickets_pog/src/ticket_sale.gleam new file mode 100644 index 0000000..54f4ac9 --- /dev/null +++ b/examples/tickets_pog/src/ticket_sale.gleam @@ -0,0 +1,317 @@ +import factos +import factos/factos_pog +import gleam/bit_array +import gleam/erlang/process +import gleam/int +import gleam/io +import gleam/list +import gleam/option.{Some} +import gleam/result +import global_value +import pog + +const event_id = "gleamconf-2026" + +const ticket_capacity = 100 + +const purchase_attempts = 300 + +const concurrency = 128 + +const postgres_pool_size = 64 + +const receive_timeout = 60_000 + +pub type Command { + BuyTicket(buyer: String) +} + +pub type Event { + TicketSold(buyer: String) +} + +pub type State { + TicketWindow(capacity: Int, sold: Int) +} + +pub type DomainError { + SoldOut(capacity: Int) +} + +pub type DecodeError { + UnknownEventType(String) + InvalidPayload(String) +} + +pub type SaleSummary { + SaleSummary(attempts: Int, accepted: Int, sold_out: Int, recorded_events: Int) +} + +type TestGlobalData { + TestGlobalData(connection: pog.Connection) +} + +type PurchaseMessage { + PurchaseFinished( + attempt: Int, + result: Result( + factos_pog.Append, + factos_pog.Error(DomainError, DecodeError), + ), + ) +} + +pub fn main() -> Nil { + case run() { + Ok(summary) -> + io.println( + "ticket sale completed: " + <> int.to_string(summary.accepted) + <> " accepted, " + <> int.to_string(summary.sold_out) + <> " sold out, " + <> int.to_string(summary.recorded_events) + <> " recorded events", + ) + Error(_) -> io.println("ticket sale failed") + } +} + +pub fn run() -> Result(SaleSummary, factos_pog.Error(DomainError, DecodeError)) { + let TestGlobalData(connection) = global_data() + use _ <- result.try(reset_schema(connection)) + + let workers = process.new_subject() + let initial = + SaleSummary(attempts: 0, accepted: 0, sold_out: 0, recorded_events: 0) + + { + use _, attempt <- int.range(from: 1, to: concurrency + 1, with: Nil) + spawn_purchase(workers, connection, attempt) + } + + collect_purchases( + workers, + connection: connection, + remaining: purchase_attempts, + next_attempt: concurrency + 1, + summary: initial, + ) +} + +fn global_data() -> TestGlobalData { + global_value.create_with_unique_name("tickets_pog.global.data", fn() { + TestGlobalData(connection: start_connection()) + }) +} + +fn start_connection() -> pog.Connection { + let pool_name = process.new_name("tickets_pog") + let config = + pog.default_config(pool_name) + |> pog.host("127.0.0.1") + |> pog.port(5433) + |> pog.database("tickets_pog") + |> pog.user("postgres") + |> pog.password(Some("postgres")) + |> pog.ssl(pog.SslDisabled) + |> pog.pool_size(postgres_pool_size) + + let assert Ok(_) = pog.start(config) + process.sleep(100) + pog.named_connection(pool_name) +} + +fn reset_schema( + connection: pog.Connection, +) -> Result(Nil, factos_pog.Error(DomainError, DecodeError)) { + use _ <- result.try( + pog.query("drop table if exists factos_events") + |> pog.execute(on: connection) + |> result.map(fn(_) { Nil }) + |> result.map_error(factos_pog.StoreError), + ) + factos_pog.migrate(connection) +} + +fn spawn_purchase( + workers: process.Subject(PurchaseMessage), + connection: pog.Connection, + attempt: Int, +) -> Nil { + let _ = + process.spawn(fn() { + process.send( + workers, + PurchaseFinished(attempt, purchase(connection, attempt)), + ) + }) + Nil +} + +fn purchase( + connection: pog.Connection, + attempt: Int, +) -> Result(factos_pog.Append, factos_pog.Error(DomainError, DecodeError)) { + // Each buyer races through the same event-context query. PostgreSQL receives a + // large amount of concurrent work via the pool, while the backend's transaction + // lock preserves the capacity invariant for this arbitrary tag-based context. + factos_pog.dispatch_context( + connection, + stream: buyer_stream(attempt), + query: sale_query(), + decider: ticket_decider(), + codec: ticket_codec(), + command: BuyTicket(buyer_name(attempt)), + ) +} + +fn collect_purchases( + workers: process.Subject(PurchaseMessage), + connection connection: pog.Connection, + remaining remaining: Int, + next_attempt next_attempt: Int, + summary summary: SaleSummary, +) -> Result(SaleSummary, factos_pog.Error(DomainError, DecodeError)) { + case remaining { + 0 -> finalize_summary(connection, summary) + _ -> + case process.receive(workers, within: receive_timeout) { + Ok(PurchaseFinished(_, Ok(_))) -> { + case next_attempt <= purchase_attempts { + True -> spawn_purchase(workers, connection, next_attempt) + False -> Nil + } + collect_purchases( + workers, + connection: connection, + remaining: remaining - 1, + next_attempt: next_attempt + 1, + summary: SaleSummary( + attempts: summary.attempts + 1, + accepted: summary.accepted + 1, + sold_out: summary.sold_out, + recorded_events: summary.recorded_events, + ), + ) + } + Ok(PurchaseFinished(_, Error(factos_pog.DomainError(SoldOut(_))))) -> { + case next_attempt <= purchase_attempts { + True -> spawn_purchase(workers, connection, next_attempt) + False -> Nil + } + collect_purchases( + workers, + connection: connection, + remaining: remaining - 1, + next_attempt: next_attempt + 1, + summary: SaleSummary( + attempts: summary.attempts + 1, + accepted: summary.accepted, + sold_out: summary.sold_out + 1, + recorded_events: summary.recorded_events, + ), + ) + } + Ok(PurchaseFinished(_, Error(error))) -> Error(error) + Error(Nil) -> Error(factos_pog.DomainError(SoldOut(summary.accepted))) + } + } +} + +fn finalize_summary( + connection: pog.Connection, + summary: SaleSummary, +) -> Result(SaleSummary, factos_pog.Error(DomainError, DecodeError)) { + use context <- result.try(factos_pog.read_context( + connection, + query: sale_query(), + decider: ticket_decider(), + codec: ticket_codec(), + )) + + echo Ok(SaleSummary( + attempts: summary.attempts, + accepted: summary.accepted, + sold_out: summary.sold_out, + recorded_events: list.length(context.events), + )) +} + +fn sale_query() -> factos.Query { + factos.query([ + factos.query_item(types: [factos.event_type("TicketSold")], tags: [ + factos.tag("event:" <> event_id), + ]), + ]) +} + +fn ticket_decider() -> factos.Decider(Command, State, Event, DomainError) { + factos.decider( + initial: TicketWindow(capacity: ticket_capacity, sold: 0), + decide: decide, + evolve: evolve, + ) +} + +fn decide(state: State, command: Command) -> Result(List(Event), DomainError) { + let TicketWindow(capacity, sold) = state + case command { + BuyTicket(buyer) -> + case sold < capacity { + True -> Ok([TicketSold(buyer)]) + False -> Error(SoldOut(capacity)) + } + } +} + +fn evolve(state: State, event: Event) -> State { + let TicketWindow(capacity, sold) = state + case event { + TicketSold(_) -> TicketWindow(capacity: capacity, sold: sold + 1) + } +} + +fn ticket_codec() -> factos_pog.EventCodec(Event, DecodeError) { + factos_pog.EventCodec(encode: encode_event, decode: decode_event) +} + +fn encode_event(event: Event) -> factos_pog.Proposed(Event) { + case event { + TicketSold(buyer) -> + factos_pog.Proposed( + id: "ticket-sold-" <> buyer, + event: event, + type_: factos.event_type("TicketSold"), + tags: [factos.tag("event:" <> event_id)], + data: bit_array.from_string(buyer), + ) + } +} + +fn decode_event( + stored: factos_pog.StoredEvent, +) -> Result(factos.Decoded(Event), DecodeError) { + case factos.event_type_name(stored.type_) { + "TicketSold" -> { + use buyer <- result.try( + bit_array.to_string(stored.data) + |> result.replace_error(InvalidPayload("buyer was not utf-8")), + ) + Ok(factos.Decoded( + event: TicketSold(buyer), + type_: stored.type_, + tags: stored.tags, + )) + } + type_name -> Error(UnknownEventType(type_name)) + } +} + +fn buyer_name(attempt: Int) -> String { + "buyer-" <> int.to_string(attempt) +} + +fn buyer_stream(attempt: Int) -> String { + "ticket-buyer-" <> int.to_string(attempt) +} diff --git a/examples/tickets_pog/src/tickets_pog.gleam b/examples/tickets_pog/src/tickets_pog.gleam new file mode 100644 index 0000000..6250e4a --- /dev/null +++ b/examples/tickets_pog/src/tickets_pog.gleam @@ -0,0 +1,5 @@ +import ticket_sale + +pub fn main() -> Nil { + ticket_sale.main() +} diff --git a/examples/tickets_pog/test/tickets_pog_test.gleam b/examples/tickets_pog/test/tickets_pog_test.gleam new file mode 100644 index 0000000..0d2914a --- /dev/null +++ b/examples/tickets_pog/test/tickets_pog_test.gleam @@ -0,0 +1,15 @@ +import gleeunit +import ticket_sale + +pub fn main() -> Nil { + gleeunit.main() +} + +pub fn ticket_sale_preserves_capacity_under_high_concurrency_test() { + let assert Ok(ticket_sale.SaleSummary( + attempts: 300, + accepted: 100, + sold_out: 200, + recorded_events: 100, + )) = ticket_sale.run() +} diff --git a/gleam.toml b/gleam.toml index 4cec0a8..6d0a376 100644 --- a/gleam.toml +++ b/gleam.toml @@ -1,5 +1,10 @@ name = "factos" version = "1.0.0" +description = "Store-independent context-first Event Sourcing primitives for Gleam." +licences = ["Apache-2.0"] +links = [ + { title = "Simply Event Sourcing", href = "https://ricofritzsche.me/simply-event-sourcing/" }, +] # Fill out these fields if you intend to generate HTML documentation or publish # your project to the Hex package manager. diff --git a/src/factos.gleam b/src/factos.gleam index 3171453..5e2ae74 100644 --- a/src/factos.gleam +++ b/src/factos.gleam @@ -4,39 +4,101 @@ //// facts, command contexts, pure decision components, and pure views. Concrete //// storage concerns live in backend packages such as `factos_sqlight` and //// `factos_kurrentdb_erlang`. +//// +//// This package follows a context-first reading of Event Sourcing: accepted +//// facts are the authoritative state of the system, and the facts relevant to a +//// command are considered before new facts are accepted. Aggregates, stream-per- +//// object storage, CQRS, projections, and message brokers are implementation +//// choices rather than prerequisites. +//// +//// The central flow is: +//// +//// 1. Select a command context with `Query`. +//// 2. Fold the matching recorded events into a temporary decision state. +//// 3. Run a pure `Decider`. +//// 4. Append the produced facts only if the context is still stable. +//// +//// Backends implement the storage-specific parts of that flow. This module keeps +//// the shared types and pure computations small and portable. import gleam/list import gleam/option.{type Option, None, Some} import gleam/result -pub type EventType { +/// A store-visible event type name. +/// +/// Event types are part of the query contract. A backend may use them for +/// efficient context reads, and applications should keep names stable enough +/// for stored history to remain decodable. +pub opaque type EventType { EventType(String) } -pub type Tag { +/// A store-visible tag value. +/// +/// Tags expose selected payload information to the event store so commands can +/// query the facts relevant to a decision. For example, an event payload may +/// contain `username: "renata"`, while the stored event also carries the tag +/// `username:renata`. +pub opaque type Tag { Tag(String) } pub type Query { + /// Match every recorded event. AllEvents + + /// Match events using one or more query items. + /// + /// Query items are OR-combined. See `QueryItem` for the matching rules inside + /// each item. Query(items: List(QueryItem)) } pub type QueryItem { + /// One branch of a command-context query. + /// + /// Within an item, event types are OR-combined and tags are AND-combined. Empty + /// `types` means any event type matches. Empty `tags` means no tag constraint. + /// + /// A query item with `types: [UserRegistered, UsernameReserved]` and + /// `tags: [username:renata]` means: events of either type that also have the + /// `username:renata` tag. QueryItem(types: List(EventType), tags: List(Tag)) } pub type SequencePosition { + /// No global position was observed. NoPosition + + /// A backend-specific global sequence position. + /// + /// Positions are used by context append conditions to express "after this + /// observed point in history". They are not stream revisions. SequencePosition(Int) } pub type AppendCondition { + /// Append without an additional context condition. NoAppendCondition + + /// Append only if no event matching `query` appeared after `after`. + /// + /// This models Command Context Consistency. The decision was made from the + /// matching facts visible at `after`, so the append must fail if that relevant + /// context changed before the new facts are recorded. FailIfEventsMatch(query: Query, after: SequencePosition) } pub type Decider(command, state, event, domain_error) { + /// A pure command-side domain component. + /// + /// `initial` is the empty decision state. `evolve` folds accepted events into + /// state. `decide` applies a command to the folded state and either returns new + /// events or a domain error. + /// + /// A decider has no dependency on storage, transactions, codecs, projections, + /// subscriptions, or transports. Decider( initial: state, decide: fn(state, command) -> Result(List(event), domain_error), @@ -45,19 +107,37 @@ pub type Decider(command, state, event, domain_error) { } pub type View(state, event) { + /// A pure projection fold. + /// + /// Views derive read-side state from events. They are intentionally only the + /// computation; persistence, delivery, rebuilds, and subscription management are + /// outside the core library. View(initial: state, evolve: fn(state, event) -> state) } pub type Revision { + /// A stream has no events. NoEvents + + /// The last known revision of a stream. CurrentRevision(Int) } pub type Decoded(event) { + /// A domain event decoded from backend storage. + /// + /// Backends use codecs supplied by the application. The decoded value includes + /// the domain event plus the event type and tags that should participate in + /// query matching. Decoded(event: event, type_: EventType, tags: List(Tag)) } pub type Recorded(event) { + /// A stored event with backend metadata. + /// + /// `revision` is the per-stream revision. `position` is the global sequence + /// position used for context consistency. `type_` and `tags` are store-visible + /// query metadata. Recorded( id: String, stream: String, @@ -70,6 +150,11 @@ pub type Recorded(event) { } pub type Context(event, state) { + /// A command context read from history. + /// + /// The context contains the query that selected the relevant facts, the folded + /// decision state, the matching recorded events, the highest observed position, + /// and the append condition needed to protect the decision. Context( query: Query, state: state, @@ -80,6 +165,10 @@ pub type Context(event, state) { } pub type LoadedStream(event, state) { + /// A stream read from history. + /// + /// This is the stream-consistency counterpart to `Context`. It contains the + /// folded state for one stream and its current revision. LoadedStream( stream: String, state: state, @@ -88,24 +177,32 @@ pub type LoadedStream(event, state) { ) } +/// Wrap an event type name. pub fn event_type(name: String) -> EventType { EventType(name) } +/// Unwrap an event type name. pub fn event_type_name(event_type: EventType) -> String { let EventType(name) = event_type name } +/// Wrap a tag value. pub fn tag(value: String) -> Tag { Tag(value) } +/// Unwrap a tag value. pub fn tag_value(tag: Tag) -> String { let Tag(value) = tag value } +/// Build a query from query items. +/// +/// An empty list becomes `AllEvents`; otherwise the query contains the supplied +/// items. Query items are OR-combined by `matches_query`. pub fn query(items: List(QueryItem)) -> Query { case items { [] -> AllEvents @@ -113,6 +210,10 @@ pub fn query(items: List(QueryItem)) -> Query { } } +/// Build one command-context query branch. +/// +/// Event types are OR-combined. Tags are AND-combined. Empty lists act as wildcards +/// for that part of the item. pub fn query_item( types types: List(EventType), tags tags: List(Tag), @@ -120,6 +221,11 @@ pub fn query_item( QueryItem(types:, tags:) } +/// Build a pure command-side decider. +/// +/// The supplied functions remain owned by the application domain. Factos only +/// stores them together so backends and tests can run the same read-decide-append +/// flow consistently. pub fn decider( initial initial: state, decide decide: fn(state, command) -> Result(List(event), domain_error), @@ -128,6 +234,10 @@ pub fn decider( Decider(initial:, decide:, evolve:) } +/// Build a pure projection view. +/// +/// A view folds events into read-side state. It does not prescribe where that +/// state is stored or how events are delivered. pub fn view( initial initial: state, evolve evolve: fn(state, event) -> state, @@ -135,6 +245,10 @@ pub fn view( View(initial:, evolve:) } +/// Fold events with a decider and decide which new events a command produces. +/// +/// This is useful for unit tests and for in-memory command handling. It does not +/// perform any append or consistency check. pub fn compute_events( decider decider: Decider(command, state, event, domain_error), events events: List(event), @@ -144,6 +258,10 @@ pub fn compute_events( decide(fold_events(initial, events, evolve), command) } +/// Decide from an optional current state and return the state after produced events. +/// +/// If `current` is `None`, the decider's initial state is used. The function first +/// runs the decider, then folds the produced events into the decision state. pub fn compute_state( decider decider: Decider(command, state, event, domain_error), current current: Option(state), @@ -159,6 +277,10 @@ pub fn compute_state( Ok(fold_events(state, events, evolve)) } +/// Fold recorded events into state using a domain evolution function. +/// +/// Backends use this after decoding stored events. Only the domain event payload is +/// passed to `evolve`; storage metadata is ignored for state computation. pub fn evolve_recorded( initial initial: state, events events: List(Recorded(event)), @@ -168,6 +290,7 @@ pub fn evolve_recorded( evolve(state, recorded.event) } +/// Project events from a view's initial state. pub fn project( view view: View(state, event), events events: List(event), @@ -176,6 +299,7 @@ pub fn project( fold_events(initial, events, evolve) } +/// Project events starting from an already materialized view state. pub fn project_from( view view: View(state, event), state state: state, @@ -185,6 +309,10 @@ pub fn project_from( fold_events(state, events, evolve) } +/// Merge two views that consume the same event type. +/// +/// The resulting view keeps both states in a tuple and evolves both for every +/// event. This is a convenience for composing small pure projections. pub fn merge_views( first first: View(first_state, event), second second: View(second_state, event), @@ -197,6 +325,10 @@ pub fn merge_views( #(first_evolve(first_state, event), second_evolve(second_state, event)) } +/// Run a command against a previously read context. +/// +/// The returned tuple preserves the original context alongside the newly produced +/// events so a backend can append them with `context.append_condition`. pub fn decide_context( context: Context(event, state), command: command, @@ -207,6 +339,7 @@ pub fn decide_context( Ok(#(context, events)) } +/// Test whether a recorded event belongs to a query-defined context. pub fn matches_query(recorded: Recorded(event), query: Query) -> Bool { case query { AllEvents -> True @@ -214,6 +347,10 @@ pub fn matches_query(recorded: Recorded(event), query: Query) -> Bool { } } +/// Return the later of two sequence positions. +/// +/// `NoPosition` acts as absence of an observed position. If both positions are +/// concrete, the larger integer is returned. pub fn highest_position( left: SequencePosition, right: SequencePosition, -- 2.51.2