diff --git a/.wrangler/cache/pages.json b/.wrangler/cache/pages.json deleted file mode 100644 index 8881842..0000000 --- a/.wrangler/cache/pages.json +++ /dev/null @@ -1,4 +0,0 @@ -{ - "account_id": "9f52ed810ad09a589c0045ad0e552873", - "project_name": "factos" -} \ No newline at end of file diff --git a/.wrangler/cache/wrangler-account.json b/.wrangler/cache/wrangler-account.json deleted file mode 100644 index a36bd32..0000000 --- a/.wrangler/cache/wrangler-account.json +++ /dev/null @@ -1,6 +0,0 @@ -{ - "account": { - "id": "9f52ed810ad09a589c0045ad0e552873", - "name": "Renata.amutio@gmail.com's Account" - } -} \ No newline at end of file diff --git a/README.md b/README.md index e5c8d26..837c402 100644 --- a/README.md +++ b/README.md @@ -12,9 +12,6 @@ The library helps with the repetitive part of event-sourced applications: 5. return the committed records so your application can update views or trigger effects. -The main backend today is `factos_pog`, which stores events in PostgreSQL with -[`pog`](https://hex.pm/packages/pog). - Factos is not a large framework. Your application still defines the commands, events, state, errors, codecs, read models, and side effects. Factos gives those pieces a standard shape and gives backends a standard way to run the @@ -22,7 +19,7 @@ read-decide-append flow safely. ## What gets stored? -Backends store events. In `factos_pog`, those events are rows in PostgreSQL. +Backends store events. Materialized views are not stored by Factos itself. A `factos.View` is an in-memory fold over events. If you want a durable read model, your application @@ -101,54 +98,6 @@ factos.compute_events( ) ``` -## What does `factos_pog` provide? - -`factos_pog` is the PostgreSQL backend. It provides: - -- database migrations for an append-only event log; -- an application codec boundary for event bytes; -- context reads by event type and tag; -- stream reads by stream name; -- dispatch functions that run the read-decide-append flow in PostgreSQL; -- committed `factos.Recorded(event)` values after successful appends. - -The primary function is `dispatch_with_query`: - -```gleam -let assert Ok(dispatch) = - factos_pog.dispatch_with_query( - connection, - stream: buyer_stream(attempt), - query: sale_query(), - decider: ticket_decider(), - codec: ticket_codec(), - command: BuyTicket(buyer_name(attempt)), - ) -``` - -That call: - -1. starts a PostgreSQL transaction; -2. locks the event table; -3. reads events matching `sale_query()`; -4. decodes those rows with `ticket_codec()`; -5. folds them with the decider's `evolve` function; -6. calls the decider's `decide` function with `BuyTicket(...)`; -7. checks that no matching event appeared since the context was read; -8. inserts the new events if the decision succeeded; -9. returns `Dispatch(event)`. - -`Dispatch(event)` contains append metadata and the events committed by this -specific dispatch: - -```gleam -pub type Dispatch(event) { - Dispatch(append: Append, events: List(factos.Recorded(event))) -} -``` - -That is what your application can feed into reactors or projection updates. - ## Events, commands, and command sourcing Factos stores events: facts that were accepted by the application. A backend row @@ -251,61 +200,13 @@ Factos does not send the email, publish the webhook, or mark the effect as done. It keeps that work explicit so your application can choose the durability and retry strategy. -## Failure modes - -`factos_pog` separates common failure classes: - -- `DomainError(error)`: your decider rejected the command. -- `StoreError(error)`: PostgreSQL or `pog` failed. -- `AppendConditionFailed(condition)`: the context or stream changed before append. -- `DecodeError(error)`: stored bytes could not be decoded by your codec. - -This distinction matters operationally. A sold-out ticket is not a database -failure. A decode failure means stored history and current codec no longer agree. -An append-condition failure usually means the command should be retried from a -fresh context or rejected with newer information. - -## Scaling model - -The current `factos_pog` context dispatch path prioritizes correctness over write -throughput. It locks the event table while running `dispatch_with_query`, so -concurrent writers queue behind each other even if their contexts do not overlap. - -That simple lock is what makes arbitrary event-type/tag append conditions -correct in this backend. - -Use `dispatch` when one stream revision is the correct consistency boundary. Use -`dispatch_with_query` when the rule spans facts selected by event type and tag. -Future PostgreSQL backends can use more granular locking, but they must preserve -the same append-condition guarantee. - -Read scaling and projection scaling are application concerns. You can recompute -views by replaying events, maintain your own materialized tables, or build -subscription workers on top of backend reads. - -## Example - -Run the PostgreSQL ticket-sale example: - -```sh -cd examples/tickets_pog -docker compose up -d -gleam run -``` - -It starts many concurrent buyers for one event. Only 100 tickets can be accepted. -The backend stores the accepted `TicketSold` events, protects the tag-based -capacity context, and returns committed records for the reactor. - ## Repository packages This repository contains: 1. `factos`: core primitives and pure computations. -2. `factos_pog`: PostgreSQL backend using `pog`. -3. `factos_sqlight`: SQLite backend using `sqlight`. -4. `factos_kurrentdb_erlang`: KurrentDB backend for Erlang. -5. `factos_cf`: Cloudflare D1 backend. +2. `factos_kurrentdb_erlang`: KurrentDB backend for Erlang. +3. `factos_cf`: Cloudflare D1 backend. The core concepts are shared. Storage behaviour and scaling tradeoffs are backend specific. diff --git a/backends/factos_sqlight/gleam.toml b/backends/factos_sqlight/gleam.toml deleted file mode 100644 index e9daad1..0000000 --- a/backends/factos_sqlight/gleam.toml +++ /dev/null @@ -1,18 +0,0 @@ -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 = { git = "https://github.com/renatillas/factos", ref = "main" } -gleam_stdlib = ">= 1.0.0 and < 2.0.0" -sqlight = ">= 1.1.0 and < 2.0.0" - -[dev_dependencies] -gleeunit = ">= 1.0.0 and < 2.0.0" -gleam_erlang = ">= 1.3.0 and < 2.0.0" -simplifile = ">= 2.5.0 and < 3.0.0" diff --git a/backends/factos_sqlight/manifest.toml b/backends/factos_sqlight/manifest.toml deleted file mode 100644 index 254d8d9..0000000 --- a/backends/factos_sqlight/manifest.toml +++ /dev/null @@ -1,26 +0,0 @@ -# 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 = "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.5.0", build_tools = ["gleam"], requirements = ["filepath", "gleam_stdlib"], otp_app = "simplifile", source = "hex", outer_checksum = "6C72DCCDF25C38A5931740B30E823969F33106831FD1637719B5EDBCA30027A4" }, - { name = "sqlight", version = "1.1.0", build_tools = ["gleam"], requirements = ["esqlite", "gleam_stdlib"], otp_app = "sqlight", source = "hex", outer_checksum = "ECA1A4B45C35EB9EFCEEB7FAAC7BF5D8B2C777A7C1FC8A9C12CB67D54CED42E7" }, -] - -[requirements] -factos = { path = "../.." } -gleam_erlang = { version = ">= 1.3.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.5.0 and < 3.0.0" } -sqlight = { version = ">= 1.1.0 and < 2.0.0" } diff --git a/backends/factos_sqlight/priv/migrations.sql b/backends/factos_sqlight/priv/migrations.sql deleted file mode 100644 index 77c96ec..0000000 --- a/backends/factos_sqlight/priv/migrations.sql +++ /dev/null @@ -1,18 +0,0 @@ -create table if not exists factos_events ( - position integer primary key autoincrement, - id text not null, - stream text not null, - revision integer not null, - type text not null, - version integer not null, - tags text not null, - metadata text not null default '', - data blob not null, - unique(stream, revision) -); - -create index if not exists factos_events_stream_revision - on factos_events(stream, revision); - -create index if not exists factos_events_position - on factos_events(position); diff --git a/backends/factos_sqlight/src/factos/factos_sqlight.gleam b/backends/factos_sqlight/src/factos/factos_sqlight.gleam deleted file mode 100644 index c5ccdfd..0000000 --- a/backends/factos_sqlight/src/factos/factos_sqlight.gleam +++ /dev/null @@ -1,684 +0,0 @@ -//// 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 -import gleam/list -import gleam/result -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, - type_: factos.EventType, - version: Int, - tags: List(factos.Tag), - metadata: factos.Metadata, - data: BitArray, - ) -} - -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, - stream: String, - revision: Int, - type_: factos.EventType, - version: Int, - tags: List(factos.Tag), - metadata: factos.Metadata, - data: BitArray, - ) -} - -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), - ) -} - -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 Dispatch(event) { - /// Result of a successful dispatch. - /// - /// `append` has the stream revision and final global position. `events` are the - /// committed events recorded by this dispatch, suitable for pure Factos - /// reactors or backend-specific durable effect adapters. - Dispatch(append: Append, events: List(factos.Recorded(event))) -} - -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) -} - -/// Deprecated: read `priv/migrations.sql` from the `factos_sqlight` application -/// with `gleam/erlang/application.priv_directory` and run it with your -/// application's migration tool instead. -/// 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. -@deprecated("Use the SQL file at factos_sqlight/priv/migrations.sql instead.") -pub fn migrate(connection: sqlight.Connection) -> Result(Nil, Error(_, _)) { - sqlight.exec( - " - create table if not exists factos_events ( - position integer primary key autoincrement, - id text not null, - stream text not null, - revision integer not null, - type text not null, - version integer not null, - tags text not null, - metadata text not null default '', - data blob not null, - unique(stream, revision) - ); - create index if not exists factos_events_stream_revision - on factos_events(stream, revision); - create index if not exists factos_events_position - on factos_events(position); - ", - 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, - 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. -/// -/// 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, - query query: factos.Query, - decider decider: factos.Decider(command, state, event, domain_error), - codec codec: EventCodec(event, decode_error), - command command: command, -) -> Result(Dispatch(event), Error(domain_error, decode_error)) { - use _ <- result.try( - sqlight.exec("begin immediate", on: connection) - |> result.map_error(StoreError), - ) - - let result = { - use context <- result.try(read_context( - 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( - connection, - stream_name, - events, - codec, - context.append_condition, - ) - } - - 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, - 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. -/// -/// 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, - decider decider: factos.Decider(command, state, event, domain_error), - codec codec: EventCodec(event, decode_error), - command command: command, -) -> Result(Dispatch(event), Error(domain_error, decode_error)) { - use _ <- result.try( - sqlight.exec("begin immediate", on: connection) - |> result.map_error(StoreError), - ) - - let result = { - use loaded <- result.try(load_stream( - 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( - connection, - stream_name, - events, - codec, - loaded.revision, - ) - } - - finish_transaction(connection, result) -} - -fn finish_transaction( - connection: sqlight.Connection, - result: Result(Dispatch(event), Error(domain_error, decode_error)), -) -> Result(Dispatch(event), Error(domain_error, decode_error)) { - case result { - Ok(dispatch) -> - case sqlight.exec("commit", on: connection) { - Ok(Nil) -> Ok(dispatch) - Error(error) -> Error(StoreError(error)) - } - Error(error) -> { - let _ = sqlight.exec("rollback", on: connection) - Error(error) - } - } -} - -fn append_with_condition( - connection: sqlight.Connection, - stream_name: String, - events: List(event), - codec: EventCodec(event, decode_error), - condition: factos.AppendCondition, -) -> Result(Dispatch(event), 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: sqlight.Connection, - stream_name: String, - events: List(event), - codec: EventCodec(event, decode_error), -) -> Result(Dispatch(event), 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: sqlight.Connection, - stream_name: String, - events: List(event), - codec: EventCodec(event, decode_error), - expected: factos.Revision, -) -> Result(Dispatch(event), Error(domain_error, decode_error)) { - case events { - [] -> { - let append = - Append( - current_revision: revision_to_int(expected), - position: factos.NoPosition, - ) - Ok(Dispatch(append:, events: [])) - } - [_, ..] -> { - 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: sqlight.Connection, - stream_name: String, - events: List(event), - codec: EventCodec(event, decode_error), - revision: Int, - position: factos.SequencePosition, - recorded_events: List(factos.Recorded(event)), -) -> Result(Dispatch(event), Error(domain_error, decode_error)) { - case events { - [] -> { - let append = Append(current_revision: revision - 1, position: position) - Ok(Dispatch(append:, events: list.reverse(recorded_events))) - } - [event, ..rest] -> { - let EventCodec(encode, _) = codec - let Proposed(id, _, type_, version, tags, metadata, data) = encode(event) - use positions <- result.try( - sqlight.query( - " - insert into factos_events (id, stream, revision, type, version, tags, metadata, data) - values (?, ?, ?, ?, ?, ?, ?, ?) - returning position - ", - on: connection, - with: [ - sqlight.text(id), - sqlight.text(stream_name), - sqlight.int(revision), - sqlight.text(factos.event_type_name(type_)), - sqlight.int(version), - sqlight.text(tags_to_text(tags)), - sqlight.text(metadata_to_text(metadata)), - sqlight.blob(data), - ], - expecting: int_field_decoder(), - ) - |> result.map_error(StoreError), - ) - let position = case positions { - [position, ..] -> factos.SequencePosition(position) - [] -> position - } - let recorded = - factos.Recorded( - id: id, - stream: stream_name, - revision: revision, - position: position, - type_: type_, - version: version, - tags: tags, - metadata: metadata, - event: event, - ) - insert_events( - connection, - stream_name, - rest, - codec, - revision + 1, - position, - [recorded, ..recorded_events], - ) - } - } -} - -fn read_matching_events( - connection: sqlight.Connection, - query: factos.Query, - codec: EventCodec(event, decode_error), -) -> Result(List(factos.Recorded(event)), Error(domain_error, decode_error)) { - use rows <- result.try( - sqlight.query( - "select position, id, stream, revision, type, version, tags, metadata, data from factos_events order by position", - on: connection, - with: [], - expecting: stored_event_decoder(), - ) - |> result.map_error(StoreError), - ) - decode_rows(rows, codec) - |> result.map(list.filter(_, factos.matches_query(_, query))) -} - -fn read_stream_events( - connection: sqlight.Connection, - stream_name: String, - codec: EventCodec(event, decode_error), -) -> Result(List(factos.Recorded(event)), Error(domain_error, decode_error)) { - use rows <- result.try( - sqlight.query( - "select position, id, stream, revision, type, version, tags, metadata, data from factos_events where stream = ? order by revision", - on: connection, - with: [sqlight.text(stream_name)], - expecting: stored_event_decoder(), - ) - |> 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_, version, tags, metadata) = decoded - let StoredEvent(position, id, stream, revision, _, _, _, _, _) = row - - Ok(factos.Recorded( - id: id, - stream: stream, - revision: revision, - position: factos.SequencePosition(position), - type_: type_, - version: version, - tags: tags, - metadata: metadata, - 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 version <- decode.field(5, decode.int) - use tags <- decode.field(6, decode.string) - use metadata <- decode.field(7, decode.string) - use data <- decode.field(8, decode.bit_array) - decode.success(StoredEvent( - position: position, - id: id, - stream: stream, - revision: revision, - type_: factos.event_type(type_name), - version: version, - tags: tags_from_text(tags), - metadata: metadata_from_text(metadata), - data: data, - )) -} - -fn current_revision( - connection: sqlight.Connection, - stream_name: String, -) -> Result(Int, sqlight.Error) { - 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(), - )) - - case rows { - [revision, ..] -> revision - [] -> -1 - } -} - -fn has_matching_events_after( - connection: sqlight.Connection, - query: factos.Query, - after: factos.SequencePosition, -) -> Result(Bool, sqlight.Error) { - let after_position = case after { - factos.NoPosition -> -1 - factos.SequencePosition(position) -> position - } - sqlight.query( - "select type, tags from factos_events where position > ?", - on: connection, - with: [sqlight.int(after_position)], - expecting: query_match_decoder(), - ) - |> result.map( - list.any(_, fn(pair) { - 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_, - version: 1, - tags: tags, - metadata: factos.empty_metadata(), - 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) - } -} - -fn metadata_to_text(metadata: factos.Metadata) -> String { - metadata - |> factos.metadata_entries - |> list.map(fn(entry) { entry.0 <> "=" <> entry.1 }) - |> string.join(with: "\n") -} - -fn metadata_from_text(metadata: String) -> factos.Metadata { - case string.is_empty(metadata) { - True -> factos.empty_metadata() - False -> - metadata - |> string.split(on: "\n") - |> list.filter_map(fn(entry) { - case string.split(entry, on: "=") { - [key, value] -> Ok(#(key, value)) - _ -> Error(Nil) - } - }) - |> factos.metadata - } -} diff --git a/backends/factos_sqlight/test/factos_sqlight_test.gleam b/backends/factos_sqlight/test/factos_sqlight_test.gleam deleted file mode 100644 index 372365e..0000000 --- a/backends/factos_sqlight/test/factos_sqlight_test.gleam +++ /dev/null @@ -1,418 +0,0 @@ -import factos -import factos/factos_sqlight -import gleam/bit_array -import gleam/erlang/application -import gleam/int -import gleam/list -import gleam/result -import gleeunit -import simplifile -import sqlight - -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) -} - -pub fn dispatch_stream_persists_events_test() { - use connection <- sqlight.with_connection(":memory:") - execute_migration_file(connection) - - let assert Ok(dispatch) = - factos_sqlight.dispatch_stream( - connection, - stream: "user-renata", - decider: decider(), - codec: codec(), - command: RegisterUser("renata"), - ) - - let assert factos_sqlight.Append( - current_revision: 0, - position: factos.SequencePosition(_), - ) = dispatch.append - let assert [recorded] = dispatch.events - assert_user_recorded( - recorded, - stream: "user-renata", - revision: 0, - position: dispatch.append.position, - username: "renata", - ) - let reactor = factos.reactor(react: fn(recorded) { [recorded.event] }) - assert factos.react_all(reactor: reactor, events: dispatch.events) - == [ - UserRegistered("renata"), - ] - - let assert Ok(loaded) = - factos_sqlight.load_stream( - connection, - stream: "user-renata", - decider: decider(), - codec: codec(), - ) - - assert loaded.state == Taken - assert loaded.revision == factos.CurrentRevision(0) -} - -pub fn dispatch_stream_handles_many_events_test() { - use connection <- sqlight.with_connection(":memory:") - execute_migration_file(connection) - - let assert Ok(dispatch) = dispatch_counter_stream_many(connection, 250) - let assert factos_sqlight.Append( - current_revision: 249, - position: factos.SequencePosition(_), - ) = dispatch.append - let assert [recorded] = dispatch.events - assert_counter_recorded( - recorded, - stream: "counter-load", - revision: 249, - position: dispatch.append.position, - value: 250, - ) - let reactor = factos.reactor(react: fn(recorded) { [recorded.event] }) - assert factos.react_all(reactor: reactor, events: dispatch.events) - == [ - Incremented(250), - ] - - let assert Ok(loaded) = - factos_sqlight.load_stream( - connection, - stream: "counter-load", - decider: counter_decider(), - codec: counter_codec(), - ) - - assert loaded.state == CounterState(250) - assert loaded.revision == factos.CurrentRevision(249) - assert list.length(loaded.events) == 250 -} - -pub fn dispatch_context_handles_many_streams_test() { - use connection <- sqlight.with_connection(":memory:") - execute_migration_file(connection) - - let query = - factos.query([ - factos.query_item(types: [factos.event_type("Incremented")], tags: [ - factos.tag("counter:load"), - ]), - ]) - - let assert Ok(dispatch) = - dispatch_counter_context_many(connection, query, 100) - let assert factos_sqlight.Append( - current_revision: 0, - position: factos.SequencePosition(_), - ) = dispatch.append - let assert [recorded] = dispatch.events - assert_counter_recorded( - recorded, - stream: "counter-context-1", - revision: 0, - position: dispatch.append.position, - value: 100, - ) - let reactor = factos.reactor(react: fn(recorded) { [recorded.event] }) - assert factos.react_all(reactor: reactor, events: dispatch.events) - == [ - Incremented(100), - ] - - let assert Ok(context) = - factos_sqlight.read_context( - connection, - query: query, - decider: counter_decider(), - codec: counter_codec(), - ) - - assert context.state == CounterState(100) - assert list.length(context.events) == 100 - assert context.position != factos.NoPosition -} - -fn execute_migration_file(connection: sqlight.Connection) -> Nil { - let assert Ok(priv_directory) = application.priv_directory("factos_sqlight") - let assert Ok(sql) = simplifile.read(priv_directory <> "/migrations.sql") - let assert Ok(Nil) = sqlight.exec(sql, on: connection) - Nil -} - -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_sqlight.EventCodec(Event, DecodeError) { - factos_sqlight.EventCodec(encode:, decode:) -} - -fn encode(event: Event) -> factos_sqlight.Proposed(Event) { - factos_sqlight.Proposed( - id: "event-" <> event.username, - event: event, - type_: factos.event_type("UserRegistered"), - version: 1, - tags: [factos.tag("username:" <> event.username)], - metadata: factos.empty_metadata(), - data: bit_array.from_string(event.username), - ) -} - -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.replace_error(InvalidData), - ) - Ok(factos.Decoded( - event: UserRegistered(username), - type_: stored.type_, - version: stored.version, - tags: stored.tags, - metadata: stored.metadata, - )) - } - _ -> Error(UnknownEvent) - } -} - -fn assert_user_recorded( - recorded: factos.Recorded(Event), - stream stream_name: String, - revision revision: Int, - position position: factos.SequencePosition, - username username: String, -) -> Nil { - assert recorded.id == "event-" <> username - assert recorded.stream == stream_name - assert recorded.revision == revision - assert recorded.position == position - assert recorded.type_ == factos.event_type("UserRegistered") - assert recorded.version == 1 - assert recorded.tags == [factos.tag("username:" <> username)] - assert recorded.metadata == factos.empty_metadata() - assert recorded.event == UserRegistered(username) -} - -fn dispatch_counter_stream_many( - connection: sqlight.Connection, - remaining: Int, -) -> Result( - factos_sqlight.Dispatch(CounterEvent), - factos_sqlight.Error(Nil, DecodeError), -) { - case remaining { - 0 -> - factos_sqlight.dispatch_stream( - connection, - stream: "counter-load", - decider: counter_decider(), - codec: counter_codec(), - command: Increment, - ) - _ -> { - let result = - factos_sqlight.dispatch_stream( - connection, - stream: "counter-load", - decider: counter_decider(), - codec: counter_codec(), - command: Increment, - ) - case remaining, result { - 1, _ -> result - _, Ok(_) -> dispatch_counter_stream_many(connection, remaining - 1) - _, Error(error) -> Error(error) - } - } - } -} - -fn dispatch_counter_context_many( - connection: sqlight.Connection, - query: factos.Query, - remaining: Int, -) -> Result( - factos_sqlight.Dispatch(CounterEvent), - factos_sqlight.Error(Nil, DecodeError), -) { - case remaining { - 0 -> - factos_sqlight.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_sqlight.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_sqlight.EventCodec(CounterEvent, DecodeError) { - factos_sqlight.EventCodec( - encode: encode_counter_event, - decode: decode_counter_event, - ) -} - -fn encode_counter_event( - event: CounterEvent, -) -> factos_sqlight.Proposed(CounterEvent) { - case event { - Incremented(value) -> - factos_sqlight.Proposed( - id: "counter-event-" <> int.to_string(value), - event: event, - type_: factos.event_type("Incremented"), - version: 1, - tags: [factos.tag("counter:load")], - metadata: factos.empty_metadata(), - data: bit_array.from_string(int.to_string(value)), - ) - } -} - -fn decode_counter_event( - stored: factos_sqlight.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_, - version: stored.version, - tags: stored.tags, - metadata: stored.metadata, - )) - } - _ -> Error(UnknownEvent) - } -} - -fn assert_counter_recorded( - recorded: factos.Recorded(CounterEvent), - stream stream_name: String, - revision revision: Int, - position position: factos.SequencePosition, - value value: Int, -) -> Nil { - assert recorded.id == "counter-event-" <> int.to_string(value) - assert recorded.stream == stream_name - assert recorded.revision == revision - assert recorded.position == position - assert recorded.type_ == factos.event_type("Incremented") - assert recorded.version == 1 - assert recorded.tags == [factos.tag("counter:load")] - assert recorded.metadata == factos.empty_metadata() - assert recorded.event == Incremented(value) -} diff --git a/docs/postgresql-backend.md b/docs/postgresql-backend.md deleted file mode 100644 index a40f703..0000000 --- a/docs/postgresql-backend.md +++ /dev/null @@ -1,264 +0,0 @@ -# PostgreSQL Backend - -`factos_pog` is the PostgreSQL backend for Factos. It uses -[`pog`](https://hex.pm/packages/pog) and implements the full read-decide-append -flow for context-first Event Sourcing. - -Use it when PostgreSQL is your event store and command consistency should be -protected by event type and tag queries. - -## Responsibilities - -`factos_pog` owns storage mechanics: - -- creating the event tables; -- reading contexts by event type and tag; -- reading streams by stream name; -- running PostgreSQL transactions; -- checking append conditions; -- assigning per-stream revisions and global positions; -- returning committed `factos.Recorded(event)` values. - -Your application still owns the domain: - -- command, event, state, and error types; -- deciders; -- event encoding and decoding; -- query tags; -- effect execution and retries. - -It does not maintain materialized views. `factos.View` values are in-memory -folds, and applications decide where durable read models live. - -It does not execute side effects. Successful dispatch returns committed records, -and applications decide how reactors/effect delivery should run. - -## Schema - -The backend ships its schema as `priv/migrations.sql`. Run that SQL once before -dispatching commands. In an Erlang-target application, locate the file through -the package `priv` directory: - -```gleam -import gleam/erlang/application - -let assert Ok(priv_directory) = application.priv_directory("factos_pog") -let migration_path = priv_directory <> "/migrations.sql" -``` - -`factos_pog.migrate` still exists for v1 compatibility, but it is deprecated. - -The backend creates two tables. - -`factos_events` is the append-only log: - -- `position`: global append order; -- `id`: application event id; -- `stream`: stream name; -- `revision`: per-stream revision; -- `type`: event type name; -- `version`: event version; -- `tags`: newline-encoded tags; -- `metadata`: newline-encoded metadata; -- `data`: opaque bytes. - -`factos_event_tags` mirrors tags by event position. This gives PostgreSQL an -indexed shape for tag queries without decoding event payload bytes. - -## Codecs - -PostgreSQL does not understand your domain event payload. Your application -provides a codec: - -```gleam -fn ticket_codec() -> factos_pog.EventCodec(Event) { - factos_pog.codec(encode: encode_event, decode: decode_event) -} -``` - -The encoder turns a domain event into `Proposed(event)`: - -```gleam -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"), - version: 1, - tags: [factos.tag("event:gleamconf-2026")], - metadata: factos.empty_metadata(), - data: bit_array.from_string(buyer), - ) - } -} -``` - -The decoder turns a stored row into `factos.Decoded(event)`: - -```gleam -fn decode_event( - stored: factos_pog.StoredEvent, -) -> Result(factos.Decoded(Event), factos_pog.DecodeError) { - case factos.event_type_name(stored.type_) { - "TicketSold" -> { - use buyer <- result.try( - bit_array.to_string(stored.data) - |> result.replace_error(factos_pog.InvalidData), - ) - Ok(factos.Decoded( - event: TicketSold(buyer), - type_: stored.type_, - version: stored.version, - tags: stored.tags, - metadata: stored.metadata, - )) - } - _ -> Error(factos_pog.UnknownEvent) - } -} -``` - -Decode errors stop read and dispatch flows with `factos_pog.DecodeError`. Library -code does not panic for malformed event data. - -## Context dispatch - -`dispatch_with_query` is the main context-first API: - -```gleam -factos_pog.dispatch_with_query( - connection, - stream: buyer_stream(attempt), - query: sale_query(), - decider: ticket_decider(), - codec: ticket_codec(), - command: BuyTicket(buyer_name(attempt)), -) -``` - -It runs inside a PostgreSQL transaction and performs this sequence: - -1. lock `factos_events` in exclusive mode; -2. select rows matching the query; -3. decode rows with the application codec; -4. fold those events into decision state; -5. run the decider; -6. check `FailIfEventsMatch(query, after: position)`; -7. insert the produced events; -8. mirror tags into `factos_event_tags`; -9. return `Dispatch(event)`. - -The return value contains both append metadata and committed records: - -```gleam -pub type Dispatch(event) { - Dispatch(append: Append, events: List(factos.Recorded(event))) -} -``` - -Those records are the safe input for reactors because the transaction has already -accepted them. - -## Why the table lock exists - -PostgreSQL has row locks, advisory locks, serializable transactions, and unique -constraints, but it does not have a built-in primitive for: - -> append these rows only if no row matching this arbitrary event-type/tag query -> appeared after position N. - -`factos_pog` uses `lock table factos_events in exclusive mode` to make this -correct for every `factos.Query`. This is conservative. Concurrent writers queue -behind each other, even when their contexts do not overlap. - -The tradeoff is deliberate for this backend: simple correctness before -throughput. A future PostgreSQL backend could use advisory locks or query-specific -lock keys, but only if it preserves the same context-stability guarantee. - -## Stream dispatch - -`dispatch` is available when one stream revision is the intended consistency -boundary: - -```gleam -factos_pog.dispatch( - connection, - stream: "ticket-sale-renata", - decider: ticket_decider(), - codec: ticket_codec(), - command: BuyTicket("renata"), -) -``` - -It loads the stream, folds state, runs the decider, and appends only if the -stream revision still matches the revision that was loaded. - -Use stream dispatch for stream-shaped rules. Use context dispatch for rules that -need facts selected by event type and tag. - -## Reads - -`read_context` loads the facts selected by a query and returns a -`factos.Context` with folded state and an append condition. - -`load_stream` loads one stream and returns a `factos.LoadedStream` with folded -state, decoded recorded events, and the current stream revision. - -These functions are useful for tests, diagnostics, projections, and custom -application flows. - -## Reacting to committed events - -A successful dispatch returns the committed records for that dispatch: - -```gleam -let assert Ok(dispatch) = - factos_pog.dispatch_with_query( - connection, - stream: "ticket-sale-renata", - query: sale_query(), - decider: ticket_decider(), - codec: ticket_codec(), - command: BuyTicket("renata"), - ) - -let effects = factos.react_all(ticket_reactor(), dispatch.events) -``` - -`factos_pog` does not execute effects. It exposes the accepted records so the -application can make a deliberate choice: - -- run effects immediately; -- store effects durably in another table; -- send effects to a queue; -- retry failures; -- ignore reactions during replay. - -## Errors - -`factos_pog.Error(domain_error)` has four cases: - -- `DomainError(domain_error)`: the decider rejected the command; -- `StoreError(pog.QueryError)`: PostgreSQL or `pog` failed; -- `AppendConditionFailed(factos.AppendCondition)`: the context or stream changed; -- `DecodeError(factos_pog.DecodeError)`: stored bytes could not be decoded. - -This keeps business rejection separate from storage, concurrency, and decode -failures. - -## Example - -Run the ticket sale example: - -```sh -cd examples/tickets_pog -docker compose up -d -gleam run -``` - -The example starts many concurrent buyers for the same event. Only 100 tickets -can be accepted. The backend serializes writes through the PostgreSQL transaction -lock, protects the `TicketSold` + `event:gleamconf-2026` context, and returns the -committed records so the example can run its reactor. diff --git a/examples/orders_sqlight/gleam.toml b/examples/orders_sqlight/gleam.toml deleted file mode 100644 index eda8d63..0000000 --- a/examples/orders_sqlight/gleam.toml +++ /dev/null @@ -1,13 +0,0 @@ -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 deleted file mode 100644 index 649b039..0000000 --- a/examples/orders_sqlight/manifest.toml +++ /dev/null @@ -1,28 +0,0 @@ -# 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 deleted file mode 100644 index 8865775..0000000 --- a/examples/orders_sqlight/src/order_workflow.gleam +++ /dev/null @@ -1,874 +0,0 @@ -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.Dispatch(Event), - 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.Dispatch(Event), - 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), - version: 1, - tags: [factos.tag("restaurant"), ..tags], - metadata: factos.empty_metadata(), - 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_, - version: stored.version, - tags: stored.tags, - metadata: stored.metadata, - )) -} - -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 deleted file mode 100644 index 09c68b8..0000000 --- a/examples/orders_sqlight/src/orders_sqlight.gleam +++ /dev/null @@ -1,5 +0,0 @@ -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 deleted file mode 100644 index be92879..0000000 --- a/examples/orders_sqlight/test/orders_sqlight_test.gleam +++ /dev/null @@ -1,16 +0,0 @@ -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 deleted file mode 100644 index a6570a6..0000000 --- a/examples/tickets_pog/compose.yml +++ /dev/null @@ -1,14 +0,0 @@ -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 deleted file mode 100644 index 85a5bc5..0000000 --- a/examples/tickets_pog/gleam.toml +++ /dev/null @@ -1,13 +0,0 @@ -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 deleted file mode 100644 index 964e5cc..0000000 --- a/examples/tickets_pog/manifest.toml +++ /dev/null @@ -1,33 +0,0 @@ -# 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/tickets_pog.gleam b/examples/tickets_pog/src/tickets_pog.gleam deleted file mode 100644 index af97e65..0000000 --- a/examples/tickets_pog/src/tickets_pog.gleam +++ /dev/null @@ -1,353 +0,0 @@ -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) -} - -type Effect { - AnnounceTicketSale(buyer: String, position: factos.SequencePosition) -} - -pub type State { - TicketWindow(capacity: Int, sold: Int) -} - -pub type DomainError { - SoldOut(capacity: Int) -} - -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.Dispatch(Event), factos_pog.Error(DomainError)), - ) -} - -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)) { - 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)) { - use _ <- result.try( - pog.query("drop table if exists factos_event_tags") - |> pog.execute(on: connection) - |> result.map(fn(_) { Nil }) - |> result.map_error(factos_pog.StoreError), - ) - 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 { - process.spawn(fn() { - process.send( - workers, - PurchaseFinished(attempt, purchase(connection, attempt)), - ) - }) - Nil -} - -fn purchase( - connection: pog.Connection, - attempt: Int, -) -> Result(factos_pog.Dispatch(Event), factos_pog.Error(DomainError)) { - // 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_with_query( - 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)) { - case remaining { - 0 -> finalize_summary(connection, summary) - _ -> - case process.receive(workers, within: receive_timeout) { - Ok(PurchaseFinished(_, Ok(dispatch))) -> { - run_effects(factos.react_all(ticket_reactor(), dispatch.events)) - 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)) { - use context <- result.try(factos_pog.read_context( - connection, - query: sale_query(), - decider: ticket_decider(), - codec: ticket_codec(), - )) - - 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_reactor() -> factos.Reactor(Event, Effect) { - factos.reactor(fn(recorded) { - case recorded.event { - TicketSold(buyer) -> [ - AnnounceTicketSale(buyer: buyer, position: recorded.position), - ] - } - }) -} - -fn run_effects(effects: List(Effect)) -> Nil { - use effect <- list.each(effects) - case effect { - AnnounceTicketSale(buyer:, position:) -> - io.println( - "reactor: ticket sold to " - <> buyer - <> " at position " - <> position_to_string(position), - ) - } -} - -fn position_to_string(position: factos.SequencePosition) -> String { - case position { - factos.NoPosition -> "none" - factos.SequencePosition(position) -> int.to_string(position) - } -} - -fn ticket_codec() -> factos_pog.EventCodec(Event) { - factos_pog.codec(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"), - version: 1, - tags: [factos.tag("event:" <> event_id)], - metadata: factos.empty_metadata(), - data: bit_array.from_string(buyer), - ) - } -} - -fn decode_event( - stored: factos_pog.StoredEvent, -) -> Result(factos.Decoded(Event), factos_pog.DecodeError) { - case factos.event_type_name(stored.type_) { - "TicketSold" -> { - use buyer <- result.try( - bit_array.to_string(stored.data) - |> result.replace_error(factos_pog.InvalidData), - ) - Ok(factos.Decoded( - event: TicketSold(buyer), - type_: stored.type_, - version: stored.version, - tags: stored.tags, - metadata: stored.metadata, - )) - } - _ -> Error(factos_pog.UnknownEvent) - } -} - -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/test/tickets_pog_test.gleam b/examples/tickets_pog/test/tickets_pog_test.gleam deleted file mode 100644 index 64879f1..0000000 --- a/examples/tickets_pog/test/tickets_pog_test.gleam +++ /dev/null @@ -1,15 +0,0 @@ -import gleeunit -import tickets_pog - -pub fn main() -> Nil { - gleeunit.main() -} - -pub fn ticket_sale_preserves_capacity_under_high_concurrency_test() { - let assert Ok(tickets_pog.SaleSummary( - attempts: 300, - accepted: 100, - sold_out: 200, - recorded_events: 100, - )) = tickets_pog.run() -} diff --git a/gleam.toml b/gleam.toml index f82a28d..2d0fa5b 100644 --- a/gleam.toml +++ b/gleam.toml @@ -6,8 +6,8 @@ links = [ { title = "Simply Event Sourcing", href = "https://ricofritzsche.me/simply-event-sourcing/" }, ] [repository] -type = "github" -user = "renatillas" +type = "tangled" +user = "renatillas.dev" repo = "factos" [[documentation.pages]] @@ -15,11 +15,6 @@ title = "Core Model" path = "core-model.html" source = "./docs/core-model.md" -[[documentation.pages]] -title = "PostgreSQL Backend" -path = "postgresql-backend.html" -source = "./docs/postgresql-backend.md" - [[documentation.pages]] title = "Event Logs and Command Dispatch" path = "event-sourcing.html" diff --git a/src/factos.gleam b/src/factos.gleam index 2a14476..e691acc 100644 --- a/src/factos.gleam +++ b/src/factos.gleam @@ -107,6 +107,15 @@ pub type Decider(command, state, event, domain_error) { /// /// A decider has no dependency on storage, transactions, codecs, projections, /// subscriptions, or transports. + /// + /// WARNING: deciders used by retrying backends must be pure. + /// + /// A backend may run `decide` and `evolve` more than once for the same logical + /// command when it retries a transaction conflict. These functions must be + /// deterministic and side-effect free: do not perform IO, mutate external + /// state, allocate ids from an external system, publish messages, or otherwise + /// affect the host system. Return domain events only; external effects belong + /// in durable outbox records after commit. Decider( initial: state, decide: fn(state, command) -> Result(List(event), domain_error), @@ -130,6 +139,14 @@ pub type Reactor(event, effect) { /// effect values. They do not run IO. Applications or backend adapters decide /// whether to execute effects immediately, persist them durably, retry them, or /// ignore them during replay. + /// + /// WARNING: reactors used by retrying backends must be pure. + /// + /// A backend may run `react` more than once for the same committed event while + /// retrying a transaction conflict. Reactors must only derive effect values; + /// they must not execute IO, publish messages, call external services, mutate + /// state, or allocate externally-visible ids. Persist durable effect values and + /// execute them after commit. Reactor(react: fn(Recorded(event)) -> List(effect)) } @@ -319,6 +336,11 @@ pub fn query_item( /// 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. +/// +/// WARNING: deciders used by retrying backends must be pure. `decide` and +/// `evolve` may be called more than once for the same logical command during a +/// transaction retry. They must be deterministic and side-effect free; external +/// effects belong in durable outbox records after commit. pub fn decider( initial initial: state, decide decide: fn(state, command) -> Result(List(event), domain_error), @@ -343,6 +365,10 @@ pub fn view( /// Reactors are the side-effect planning counterpart to views: they consume /// committed recorded events and return application-owned effect values without /// executing IO. +/// +/// WARNING: reactors used by retrying backends must be pure. `react` may be +/// called more than once for the same committed event during a transaction +/// retry. It must only derive effect values; executing IO belongs after commit. pub fn reactor( react react: fn(Recorded(event)) -> List(effect), ) -> Reactor(event, effect) {