diff --git a/backends/factos_cf/src/factos/factos_cf.gleam b/backends/factos_cf/src/factos/factos_cf.gleam index 26d02f6..879a1e8 100644 --- a/backends/factos_cf/src/factos/factos_cf.gleam +++ b/backends/factos_cf/src/factos/factos_cf.gleam @@ -9,12 +9,6 @@ //// conditional `insert ... select ... returning` statement for appends. That keeps //// the append condition and writes in the same SQLite statement, which is the //// atomic boundary D1 exposes through prepared statements. -//// -//// Side effects (extracted from the codec's `side_effects` list) run -//// **after** a successful append, receiving the domain events that were -//// just persisted. Dispatch functions wait for side effects to finish before -//// resolving so callers can safely depend on side-effect continuations such as -//// outbox inserts and queue notifications. import cf/d1 import factos @@ -88,13 +82,12 @@ pub fn new(database: d1.Database) -> Client { Client(database) } -/// Application-owned codec and side-effect configuration for this D1 backend. +/// Application-owned codec for this D1 backend. /// pub opaque type EventCodec(event, state) { EventCodec( encode: fn(event) -> Proposed(event), decode: fn(StoredEvent) -> Result(factos.Decoded(event), EventDecodeError), - side_effects: List(fn(List(event)) -> Promise(Nil)), ) } @@ -106,24 +99,12 @@ pub opaque type EventCodec(event, state) { /// /// `decode` turns a stored row back into a domain event. Decode failures are /// returned as `DecodeError` and stop load/read flows rather than panicking. -/// -/// `side_effects` are async hooks that run after a successful append. Each hook -/// receives the exact domain events that were just persisted. Dispatch functions -/// await all hooks before resolving, which is important for outbox workflows: -/// a side effect can insert `factos/outbox.Entry` rows and enqueue the returned -/// outbox IDs before the command response completes. -/// -/// Side effects are not part of the append transaction. Events are persisted -/// first; side effects run after the append succeeds. If a side effect can fail -/// and must be retried, prefer writing an outbox row and processing it from a -/// queue or scheduled worker. pub fn codec( encode encode: fn(event) -> Proposed(event), decode decode: fn(StoredEvent) -> Result(factos.Decoded(event), EventDecodeError), - side_effects side_effects: List(fn(List(event)) -> Promise(Nil)), ) -> EventCodec(event, state) { - EventCodec(encode:, decode:, side_effects:) + EventCodec(encode:, decode:) } /// Result of a successful append. @@ -137,6 +118,15 @@ pub type Append { 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) { /// The decider rejected the command with a domain error. DomainError(domain_error) @@ -172,16 +162,8 @@ type QuerySql { /// Create or update the D1 schema required by this backend. /// /// This function is idempotent and safe to run during application startup or in -/// tests. It creates: -/// -/// - `factos_events`, the append-only event store. -/// - indexes used by stream reads, context queries, and projection cursors. -/// - `event_outbox`, the optional outbox table used by `factos/outbox`. -/// -/// The outbox table is included here because side-effect hooks commonly need to -/// insert outbox rows immediately after successful appends. Keeping both tables -/// in one migration function prevents consumers from accidentally deploying the -/// event store without the side-effect infrastructure. +/// tests. It creates `factos_events`, the append-only event store, and indexes +/// used by stream reads, context queries, and projection cursors. pub fn migrate(client: Client) -> Promise(Result(Nil, Error(_))) { let Client(database:) = client use _ <- promise.try_await(execute_migration( @@ -208,27 +190,12 @@ pub fn migrate(client: Client) -> Promise(Result(Nil, Error(_))) { on factos_events(stream, revision) ", )) - use _ <- promise.try_await(execute_migration( + execute_migration( database, " create index if not exists factos_events_position on factos_events(position) ", - )) - execute_migration( - database, - " - create table if not exists event_outbox ( - id integer primary key autoincrement, - stream text not null, - event_type text not null, - payload text not null, - status text not null default 'pending', - error text, - created_at integer not null default (unixepoch()), - processed_at integer - ) - ", ) } @@ -283,11 +250,7 @@ pub fn read_context( /// /// This function: /// -/// - reads all events matching `query`, -/// - folds them into state, -/// - asks the decider to produce new events, -/// - appends those events only if no matching events were added meanwhile, -/// - runs and awaits codec side effects after a successful append. +/// - appends those events only if no matching events were added meanwhile. /// /// Use this for commands whose validity depends on facts outside a single /// stream. If the context changed between the read and the append, the function @@ -299,7 +262,7 @@ pub fn dispatch_with_query( decider decider: factos.Decider(command, state, event, domain_error), codec codec: EventCodec(event, state), command command: command, -) -> Promise(Result(Append, Error(domain_error))) { +) -> Promise(Result(Dispatch(event), Error(domain_error))) { use context <- promise.try_await(read_context( client, query:, @@ -313,24 +276,18 @@ pub fn dispatch_with_query( |> promise.resolve(), ) - use append <- promise.try_await(append_with_condition( + append_with_condition( client.database, stream_name, events, codec, context.append_condition, - )) - - use _nil_list <- promise.await(run_side_effects(codec, events)) - - promise.resolve(Ok(append)) + ) } /// Load and fold one stream. /// -/// Reads every event from `stream_name`, decodes the rows with `codec`, and -/// folds them through `decider.evolve`. This does not append events and does not -/// run side effects. +/// folds them through `decider.evolve`. This does not append events. pub fn load_stream( client: Client, stream stream_name: String, @@ -360,15 +317,15 @@ pub fn load_stream( /// appends the produced events with an expected-revision condition. If another /// write has advanced the stream, the append fails with `AppendConditionFailed`. /// -/// After a successful append, all codec side effects are awaited before the -/// returned promise resolves. +/// The returned dispatch includes the committed recorded events so callers can +/// run pure Factos reactors or persist backend-specific durable effects. pub fn dispatch( client: Client, stream stream_name: String, decider decider: factos.Decider(command, state, event, domain_error), codec codec: EventCodec(event, state), command command: command, -) -> Promise(Result(Append, Error(domain_error))) { +) -> Promise(Result(Dispatch(event), Error(domain_error))) { let Client(database:) = client use loaded <- promise.try_await(load_stream( client, @@ -379,31 +336,17 @@ pub fn dispatch( let factos.Decider(_, decide, _) = decider case decide(loaded.state, command) { Error(error) -> promise.resolve(Error(DomainError(error))) - Ok(events) -> { - use result <- promise.await(append_stream_events( + Ok(events) -> + append_stream_events( database, stream_name, events, codec, loaded.revision, - )) - case result { - Ok(_) -> { - use _ <- promise.await(run_side_effects(codec, events)) - promise.resolve(result) - } - Error(_) -> promise.resolve(result) - } - } + ) } } -fn run_side_effects( - codec: EventCodec(event, state), - events: List(event), -) -> promise.Promise(List(Nil)) { - promise.await_list(list.map(codec.side_effects, fn(f) { f(events) })) -} fn append_with_condition( database: d1.Database, @@ -411,7 +354,7 @@ fn append_with_condition( events: List(event), codec: EventCodec(event, state), condition: factos.AppendCondition, -) -> Promise(Result(Append, Error(domain_error))) { +) -> Promise(Result(Dispatch(event), Error(domain_error))) { case condition { factos.NoAppendCondition -> append_events(database, stream_name, events, codec, CurrentStream) @@ -432,7 +375,7 @@ fn append_stream_events( events: List(event), codec: EventCodec(event, state), expected: factos.Revision, -) -> Promise(Result(Append, Error(domain_error))) { +) -> Promise(Result(Dispatch(event), Error(domain_error))) { append_events(database, stream_name, events, codec, ExpectedStream(expected)) } @@ -442,31 +385,38 @@ fn append_events( events: List(event), codec: EventCodec(event, state), mode: AppendMode, -) -> Promise(Result(Append, Error(domain_error))) { +) -> Promise(Result(Dispatch(event), Error(domain_error))) { case events { [] -> current_revision(database, stream_name) |> promise.map(fn(result) { result |> result.map(fn(revision) { - Append(current_revision: revision, position: factos.NoPosition) + let append = Append( + current_revision: revision, + position: factos.NoPosition, + ) + Dispatch(append:, events: []) }) }) [_, ..] -> { - let #(sql, values) = append_sql(stream_name, events, codec, mode) + let EventCodec(encode, _) = codec + let proposed_events = list.map(events, encode) + let #(sql, values) = append_sql(stream_name, proposed_events, mode) let event_statement = d1.prepare(database, sql) |> d1.bind(values) d1.batch(database, [event_statement]) - |> decode_batch_result(events, mode) + |> decode_batch_result(stream_name, proposed_events, mode) } } } fn decode_batch_result( batch_result: Promise(Result(array.Array(d1.RunResult), String)), - events: List(event), + stream_name: String, + events: List(Proposed(event)), mode: AppendMode, -) -> Promise(Result(Append, Error(domain_error))) { +) -> Promise(Result(Dispatch(event), Error(domain_error))) { use result <- promise.map(batch_result) use run_results <- result.try(result |> result.map_error(StoreError)) let run_result_list = array.to_list(run_results) @@ -483,15 +433,48 @@ fn decode_batch_result( case list.length(appended) == list.length(events) { True -> { let #(position, revision) = last_append_row(appended) - Ok(Append( + let append = Append( current_revision: revision, position: factos.SequencePosition(position), + ) + Ok(Dispatch( + append: append, + events: recorded_append_rows(stream_name, events, appended), )) } False -> Error(AppendConditionFailed(append_condition_for(mode))) } } +fn recorded_append_rows( + stream_name: String, + events: List(Proposed(event)), + rows: List(#(Int, Int)), +) -> List(factos.Recorded(event)) { + case events, rows { + [], _ -> [] + _, [] -> [] + [event, ..events], [row, ..rows] -> { + let Proposed(id, domain_event, type_, version, tags, metadata, _) = event + let #(position, revision) = row + [ + factos.Recorded( + id: id, + stream: stream_name, + revision: revision, + position: factos.SequencePosition(position), + type_: type_, + version: version, + tags: tags, + metadata: metadata, + event: domain_event, + ), + ..recorded_append_rows(stream_name, events, rows) + ] + } + } +} + fn read_matching_events( database: d1.Database, query: factos.Query, @@ -638,7 +621,7 @@ fn decode_row( codec: EventCodec(event, state), ) -> Result(factos.Recorded(event), Error(domain_error)) { use stored <- result.try(decode_stored_event(row)) - let EventCodec(_, decode_event, _) = codec + let EventCodec(_, decode_event) = codec use decoded <- result.try( decode_event(stored) |> result.map_error(EventDecodeError), ) @@ -730,14 +713,13 @@ fn decode_string_field( fn append_sql( stream_name: String, - events: List(event), - codec: EventCodec(event, state), + events: List(Proposed(event)), mode: AppendMode, ) -> #(String, List(String)) { let rows = events |> list.index_map(fn(event, index) { - append_select_sql(stream_name, event, codec, mode, index) + append_select_sql(stream_name, event, mode, index) }) let sql = @@ -751,13 +733,11 @@ fn append_sql( fn append_select_sql( stream_name: String, - event: event, - codec: EventCodec(event, state), + event: Proposed(event), mode: AppendMode, index: Int, ) -> #(String, List(String)) { - let EventCodec(encode, _, _) = codec - let Proposed(id, _, type_, version, tags, metadata, data) = encode(event) + let Proposed(id, _, type_, version, tags, metadata, data) = event let base_values = [ id, stream_name, diff --git a/backends/factos_cf/src/factos/factos_cf/outbox.gleam b/backends/factos_cf/src/factos/factos_cf/outbox.gleam deleted file mode 100644 index 2b07556..0000000 --- a/backends/factos_cf/src/factos/factos_cf/outbox.gleam +++ /dev/null @@ -1,321 +0,0 @@ -//// Generic event outbox helpers for Cloudflare Workers D1. -//// -//// Provides insertion, read, status-update, and decoding helpers for rows in -//// the `event_outbox` table created by `factos_cf_workers.migrate`. -//// -//// The intended workflow is: -//// -//// 1. A Factos side effect converts domain events into `Entry` values. -//// 2. The side effect calls `insert`, receiving the inserted outbox IDs. -//// 3. The application sends those IDs to a Cloudflare Queue. -//// 4. A queue consumer calls `read_pending_outbox` by ID. -//// 5. After processing, the consumer calls `mark_outbox_sent` or -//// `mark_outbox_failed`. -//// -//// The module intentionally stores opaque payload strings. Applications decide -//// how to encode and decode payloads for each `event_type`. - -import cf/d1 -import gleam/dynamic.{type Dynamic} -import gleam/dynamic/decode -import gleam/int -import gleam/javascript/array -import gleam/javascript/promise.{type Promise} -import gleam/list -import gleam/option.{type Option} -import gleam/pair -import gleam/result -import gleam/time/timestamp - -/// A full row from the `event_outbox` table. -/// -/// `id` is the database-generated identifier used to enqueue and later read a -/// pending side effect. -/// -/// `stream` is application-owned context. Many consumers use the aggregate or -/// stream ID here so related side effects can be traced. -/// -/// `event_type` identifies what kind of side effect should be processed. Queue -/// consumers commonly dispatch on this value. -/// -/// `payload` is an opaque application-owned string, typically JSON. -/// -/// `status` is expected to be `pending`, `sent`, or `failed` by the helpers in -/// this module. The database default is `pending`. -/// -/// `error` contains a failure message for failed rows. It is decoded as an empty -/// string when the database value is null. -/// -/// `created_at` and `processed_at` are Unix timestamps. `processed_at` is -/// decoded as `0` when the database value is null. -pub type OutboxRecord { - OutboxRecord( - id: Int, - stream: String, - event_type: String, - payload: String, - status: String, - error: String, - created_at: Int, - processed_at: Int, - ) -} - -/// Errors returned by outbox reads and status updates. -/// -/// `StoreError` wraps D1 execution failures. `RowDecodeError` means D1 returned -/// a row shape that did not match the columns requested by this module. -pub type Error { - StoreError(String) - RowDecodeError(List(decode.DecodeError)) -} - -/// A pending side-effect row to insert into the outbox. -/// -/// The type is opaque so callers construct valid entries through `entry` rather -/// than depending on the table representation. This keeps the insert API small: -/// status, timestamps, errors, and IDs are owned by the database workflow. -pub opaque type Entry { - Entry(stream: String, event_type: String, payload: String) -} - -/// Construct an outbox entry for later insertion. -/// -/// `stream` should be the domain stream or aggregate identifier related to the -/// side effect. -/// -/// `event_type` should be stable and specific enough for queue consumers to -/// dispatch safely, such as `purchase.fulfillment_email_requested`. -/// -/// `payload` is intentionally a string. Use JSON or another application-owned -/// format and decode it in the queue consumer after `read_pending_outbox`. -pub fn entry( - stream stream: String, - event_type event_type: String, - payload payload: String, -) { - Entry(stream:, event_type:, payload:) -} - -/// Insert outbox entries and return their generated IDs. -/// -/// The inserts are batched through D1, preserving all-or-nothing behaviour for -/// the supplied entries. The SQL uses `returning id`, and the returned IDs are -/// intended to be sent to a queue for asynchronous processing. -/// -/// Passing an empty list is valid and returns `Ok([])` without calling -/// `d1.batch`. This matters because Cloudflare D1 rejects empty batches with -/// `No SQL statements detected`. -/// -/// Failures are collapsed to `Error(Nil)` because callers generally only need -/// to know whether IDs were produced. Use D1 logs for lower-level diagnostics. -pub fn insert( - database: d1.Database, - entries: List(Entry), -) -> Promise(Result(List(Int), Nil)) { - case entries { - [] -> promise.resolve(Ok([])) - _ -> { - let results = - { - use entry <- list.map(entries) - - d1.prepare( - database, - "insert into event_outbox (stream, event_type, payload) values (?, ?, ?) returning id;", - ) - |> d1.bind([entry.stream, entry.event_type, entry.payload]) - } - |> d1.batch(database, _) - - use results <- promise.await(results) - - case results { - Ok(results) -> { - results - |> array.to_list - |> list.map(fn(result) { - case result { - d1.RunResult(success: True, results:, ..) -> - decode_ids_from_rows(results) - d1.RunResult(success: False, ..) -> [] - } - }) - |> list.flatten - |> Ok - |> promise.resolve() - } - Error(_) -> promise.resolve(Error(Nil)) - } - } - } -} - -/// Fetch a single pending outbox row by ID. -/// -/// Returns `Ok(option.None)` if the row does not exist or is no longer pending. -/// Queue consumers should treat that as already processed or not actionable. -/// -/// Only pending rows are returned so retrying a queue message after successful -/// processing does not re-run the side effect. -pub fn read_pending_outbox( - database: d1.Database, - id: Int, -) -> Promise(Result(Option(OutboxRecord), Error)) { - d1.prepare( - database, - "select id, stream, event_type, payload, status, error, created_at, processed_at from event_outbox where id = ? and status = 'pending' limit 1", - ) - |> d1.bind([int.to_string(id)]) - |> d1.raw - |> promise.map(fn(result) { - use rows <- result.try(result |> result.map_error(StoreError)) - case rows |> array.to_list { - [] -> Ok(option.None) - [row, ..] -> row |> decode_row |> result.map(option.Some) - } - }) -} - -/// Mark an outbox row as successfully processed. -/// -/// `processed_at` should be a Unix timestamp chosen by the application. The row -/// status is changed to `sent`, and future `read_pending_outbox` calls for the -/// same ID will return `option.None`. -pub fn mark_outbox_sent( - database: d1.Database, - id: Int, - processed_at: timestamp.Timestamp, -) -> Promise(Result(Nil, Error)) { - d1.prepare( - database, - "update event_outbox set status = 'sent', processed_at = ? where id = ?", - ) - |> d1.bind([ - processed_at - |> timestamp.to_unix_seconds_and_nanoseconds - |> pair.first - |> int.to_string, - int.to_string(id), - ]) - |> d1.run - |> promise.map(fn(result) { - result - |> result.map(constant_nil) - |> result.map_error(StoreError) - }) -} - -/// Mark an outbox row as failed with an error message. -/// -/// `error_message` is stored for operational debugging. The helper also records -/// `processed_at`, allowing consumers to distinguish unprocessed pending rows -/// from failed rows that were attempted. -pub fn mark_outbox_failed( - database: d1.Database, - id: Int, - error_message: String, - processed_at: timestamp.Timestamp, -) -> Promise(Result(Nil, Error)) { - use result <- promise.map( - d1.prepare( - database, - "update event_outbox set status = 'failed', error = ?, processed_at = ? where id = ?", - ) - |> d1.bind([ - error_message, - processed_at - |> timestamp.to_unix_seconds_and_nanoseconds - |> pair.first - |> int.to_string, - int.to_string(id), - ]) - |> d1.run, - ) - result - |> result.map(constant_nil) - |> result.map_error(StoreError) -} - -fn decode_row(row: array.Array(Dynamic)) -> Result(OutboxRecord, Error) { - use id <- result.try(decode_int_field(row, 0)) - use stream <- result.try(decode_string_field(row, 1)) - use event_type <- result.try(decode_string_field(row, 2)) - use payload <- result.try(decode_string_field(row, 3)) - use status <- result.try(decode_string_field(row, 4)) - use error <- result.try(decode_nullable_string_field(row, 5)) - use created_at <- result.try(decode_int_field(row, 6)) - use processed_at <- result.try(decode_nullable_int_field(row, 7)) - - Ok(OutboxRecord( - id:, - stream:, - event_type:, - payload:, - status:, - error:, - created_at:, - processed_at:, - )) -} - -fn decode_int_field( - row: array.Array(Dynamic), - index: Int, -) -> Result(Int, Error) { - use value <- result.try( - array.get(row, index) - |> result.replace_error(RowDecodeError([])), - ) - decode.run(value, decode.int) - |> result.map_error(RowDecodeError) -} - -fn decode_string_field( - row: array.Array(Dynamic), - index: Int, -) -> Result(String, Error) { - use value <- result.try( - array.get(row, index) - |> result.replace_error(RowDecodeError([])), - ) - decode.run(value, decode.string) - |> result.map_error(RowDecodeError) -} - -fn decode_nullable_string_field( - row: array.Array(Dynamic), - index: Int, -) -> Result(String, Error) { - use value <- result.try( - array.get(row, index) - |> result.replace_error(RowDecodeError([])), - ) - decode.run(value, decode.optional(decode.string)) - |> result.map_error(RowDecodeError) - |> result.map(fn(opt) { option.unwrap(opt, "") }) -} - -fn decode_nullable_int_field( - row: array.Array(Dynamic), - index: Int, -) -> Result(Int, Error) { - use value <- result.try( - array.get(row, index) - |> result.replace_error(RowDecodeError([])), - ) - decode.run(value, decode.optional(decode.int)) - |> result.map_error(RowDecodeError) - |> result.map(fn(opt) { option.unwrap(opt, 0) }) -} - -fn decode_ids_from_rows(rows: array.Array(dynamic.Dynamic)) -> List(Int) { - use row <- list.filter_map(array.to_list(rows)) - - let decoder = decode.field("id", decode.int, decode.success) - decode.run(row, decoder) |> result.map_error(constant_nil) -} - -fn constant_nil(_: a) -> Nil { - Nil -} diff --git a/backends/factos_cf/test/factos_cf_test.gleam b/backends/factos_cf/test/factos_cf_test.gleam index 2ac6948..7b91047 100644 --- a/backends/factos_cf/test/factos_cf_test.gleam +++ b/backends/factos_cf/test/factos_cf_test.gleam @@ -3,10 +3,6 @@ import cf/miniflare import cf/miniflare/bindings import factos import factos/factos_cf -import factos/factos_cf/outbox -import gleam/dynamic.{type Dynamic} -import gleam/dynamic/decode -import gleam/int import gleam/javascript/promise.{type Promise} import gleam/list import gleam/option @@ -69,15 +65,14 @@ pub fn dispatch_stream_appends_and_loads_events_test() -> Promise(Nil) { command: Reserve("renata"), )) case append_result { - Ok(factos_cf.Append( - current_revision: 0, - position: factos.SequencePosition(_), - )) -> Nil - Ok(append) -> { - let message = - "stream append returned unexpected metadata: " - <> append_to_string(append) - panic as message + Ok(dispatch) -> { + let assert factos_cf.Append( + current_revision: 0, + position: factos.SequencePosition(_), + ) = dispatch.append + let assert [recorded] = dispatch.events + assert recorded.event == Reserved("renata") + Nil } Error(error) -> { let message = "stream append failed: " <> error_to_string(error) @@ -148,15 +143,14 @@ pub fn dispatch_context_rejects_changed_context_test() -> Promise(Nil) { command: Reserve("context-renata"), )) case first_result { - Ok(factos_cf.Append( - current_revision: 0, - position: factos.SequencePosition(_), - )) -> Nil - Ok(append) -> { - let message = - "context append returned unexpected metadata: " - <> append_to_string(append) - panic as message + Ok(dispatch) -> { + let assert factos_cf.Append( + current_revision: 0, + position: factos.SequencePosition(_), + ) = dispatch.append + let assert [recorded] = dispatch.events + assert recorded.event == Reserved("context-renata") + Nil } Error(error) -> { let message = "context append failed: " <> error_to_string(error) @@ -200,93 +194,6 @@ pub fn dispatch_context_rejects_changed_context_test() -> Promise(Nil) { promise.resolve(Nil) } -pub fn dispatch_awaits_side_effects_test() -> Promise(Nil) { - use test_database <- promise.await(new_test_database()) - let client = test_client(test_database.database) - use Nil <- promise.await(migrate_or_panic(client)) - use Nil <- promise.await(clear_events_or_panic(test_database.database)) - use Nil <- promise.await(create_side_effect_marker_table( - test_database.database, - )) - - use append_result <- promise.await(factos_cf.dispatch( - client, - stream: "reservation-side-effect", - decider: reservation_decider(), - codec: codec_with_side_effect(test_database.database), - command: Reserve("side-effect"), - )) - case append_result { - Ok(_) -> Nil - Error(error) -> { - let message = "stream append failed: " <> error_to_string(error) - panic as message - } - } - - use marker_written <- promise.await(side_effect_marker_written( - test_database.database, - )) - case marker_written { - True -> Nil - False -> panic as "dispatch returned before side effect finished" - } - - use Nil <- promise.await(dispose(test_database)) - promise.resolve(Nil) -} - -pub fn dispatch_with_state_awaits_side_effects_test() -> Promise(Nil) { - use test_database <- promise.await(new_test_database()) - let client = test_client(test_database.database) - use Nil <- promise.await(migrate_or_panic(client)) - use Nil <- promise.await(clear_events_or_panic(test_database.database)) - use Nil <- promise.await(create_side_effect_marker_table( - test_database.database, - )) - - use append_result <- promise.await(factos_cf.dispatch( - client, - stream: "reservation-side-effect-state", - decider: reservation_decider(), - codec: codec_with_side_effect(test_database.database), - command: Reserve("side-effect-state"), - )) - case append_result { - Ok(_) -> Nil - Error(error) -> { - let message = "stream append failed: " <> error_to_string(error) - panic as message - } - } - - use marker_written <- promise.await(side_effect_marker_written( - test_database.database, - )) - case marker_written { - True -> Nil - False -> panic as "dispatch_with_state returned before side effect finished" - } - - use Nil <- promise.await(dispose(test_database)) - promise.resolve(Nil) -} - -pub fn outbox_insert_with_no_entries_returns_empty_ids_test() -> Promise(Nil) { - use test_database <- promise.await(new_test_database()) - let client = test_client(test_database.database) - use Nil <- promise.await(migrate_or_panic(client)) - - use ids <- promise.await(outbox.insert(test_database.database, [])) - case ids { - Ok([]) -> Nil - Ok(_) -> panic as "empty outbox insert returned ids" - Error(_) -> panic as "empty outbox insert failed" - } - - use Nil <- promise.await(dispose(test_database)) - promise.resolve(Nil) -} fn new_test_database() -> Promise(TestDatabase) { let worker = @@ -313,27 +220,6 @@ fn dispose(test_database: TestDatabase) -> Promise(Nil) { miniflare.dispose(test_database.miniflare) } -fn migrate_or_panic(client: factos_cf.Client) -> Promise(Nil) { - use migrate_result <- promise.await(factos_cf.migrate(client)) - case migrate_result { - Ok(Nil) -> promise.resolve(Nil) - Error(error) -> { - let message = "migration failed: " <> error_to_string(error) - panic as message - } - } -} - -fn clear_events_or_panic(database: d1.Database) -> Promise(Nil) { - use clear_result <- promise.await(clear_events(database)) - case clear_result { - Ok(Nil) -> promise.resolve(Nil) - Error(error) -> { - let message = "clear events failed: " <> error - panic as message - } - } -} fn clear_events(database: d1.Database) -> Promise(Result(Nil, String)) { d1.prepare(database, "delete from factos_events") @@ -341,47 +227,6 @@ fn clear_events(database: d1.Database) -> Promise(Result(Nil, String)) { |> promise.map(fn(run_result) { run_result |> result.map(fn(_) { Nil }) }) } -fn create_side_effect_marker_table(database: d1.Database) -> Promise(Nil) { - use result <- promise.await(d1.exec( - database, - "create table side_effect_marker (id integer primary key autoincrement);", - )) - case result { - Ok(_) -> promise.resolve(Nil) - Error(error) -> panic as error - } -} - -fn insert_side_effect_marker(database: d1.Database) -> Promise(Nil) { - use Nil <- promise.await(promise.wait(50)) - use result <- promise.await( - d1.prepare(database, "insert into side_effect_marker default values") - |> d1.run, - ) - case result { - Ok(_) -> promise.resolve(Nil) - Error(error) -> panic as error - } -} - -fn side_effect_marker_written(database: d1.Database) -> Promise(Bool) { - use result <- promise.await( - d1.prepare(database, "select count(*) as count from side_effect_marker") - |> d1.first, - ) - case result { - Ok(row) -> promise.resolve(decode_count(row) > 0) - Error(error) -> panic as error - } -} - -fn decode_count(row: Dynamic) -> Int { - let decoder = decode.field("count", decode.int, decode.success) - case decode.run(row, decoder) { - Ok(count) -> count - Error(_) -> 0 - } -} fn error_to_string(error: factos_cf.Error(DomainError)) -> String { case error { @@ -393,20 +238,6 @@ fn error_to_string(error: factos_cf.Error(DomainError)) -> String { } } -fn append_to_string(append: factos_cf.Append) -> String { - let factos_cf.Append(current_revision:, position:) = append - "current_revision=" - <> int.to_string(current_revision) - <> ", 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 reservation_decider() -> factos.Decider( Command, @@ -437,16 +268,9 @@ fn evolve(state: List(String), event: Event) -> List(String) { } fn codec() -> factos_cf.EventCodec(Event, List(String)) { - factos_cf.codec(encode:, decode:, side_effects: []) + factos_cf.codec(encode:, decode:) } -fn codec_with_side_effect( - database: d1.Database, -) -> factos_cf.EventCodec(Event, List(String)) { - factos_cf.codec(encode: encode, decode: decode, side_effects: [ - fn(_) { insert_side_effect_marker(database) }, - ]) -} fn encode(event: Event) -> factos_cf.Proposed(Event) { case event { diff --git a/backends/factos_kurrentdb_erlang/src/factos/factos_kurrentdb_erlang.gleam b/backends/factos_kurrentdb_erlang/src/factos/factos_kurrentdb_erlang.gleam index 3e636af..e43bdb9 100644 --- a/backends/factos_kurrentdb_erlang/src/factos/factos_kurrentdb_erlang.gleam +++ b/backends/factos_kurrentdb_erlang/src/factos/factos_kurrentdb_erlang.gleam @@ -53,6 +53,14 @@ pub type EventCodec(event, decode_error) { ) } +pub type Dispatch(event) { + /// Result of a successful dispatch. + /// + /// `append` is the KurrentDB append response. `events` are the committed events + /// observed after the append, suitable for pure Factos reactors. + Dispatch(append: append_to_stream.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) @@ -112,7 +120,7 @@ pub fn dispatch_context( codec codec: EventCodec(event, decode_error), command command: command, timeout timeout: Int, -) -> Result(append_to_stream.Append, Error(domain_error, decode_error)) { +) -> Result(Dispatch(event), Error(domain_error, decode_error)) { use context <- result.try(read_context( connection, query:, @@ -126,14 +134,15 @@ pub fn dispatch_context( ) let #(context, events) = pair - append_with_condition( + use append <- result.try(append_with_condition( connection, stream_name, events, codec, context.append_condition, timeout, - ) + )) + Ok(Dispatch(append:, events: [])) } /// Load and fold one KurrentDB stream. @@ -168,7 +177,7 @@ pub fn dispatch_stream( codec codec: EventCodec(event, decode_error), command command: command, timeout timeout: Int, -) -> Result(append_to_stream.Append, Error(domain_error, decode_error)) { +) -> Result(Dispatch(event), Error(domain_error, decode_error)) { let factos.Decider(initial, decide, evolve) = decider use loaded <- result.try(load_stream_events( @@ -185,14 +194,26 @@ pub fn dispatch_stream( |> result.map_error(DomainError), ) - append_stream_events( + use append <- result.try(append_stream_events( connection, stream_name, events, codec, loaded.revision, timeout, - ) + )) + use loaded_after_append <- result.try(load_stream_events( + connection, + stream_name, + initial, + evolve, + codec, + timeout, + )) + Ok(Dispatch( + append: append, + events: appended_events_after(loaded_after_append.events, loaded.revision), + )) } fn append_with_condition( @@ -275,6 +296,17 @@ fn load_stream_events( ) } +fn appended_events_after( + events: List(factos.Recorded(event)), + revision: factos.Revision, +) -> List(factos.Recorded(event)) { + case revision { + factos.NoEvents -> events + factos.CurrentRevision(revision) -> + list.filter(events, fn(event) { event.revision > revision }) + } +} + fn append_stream_events( connection: kurrentdb_erlang.Connection, stream_name: String, diff --git a/backends/factos_kurrentdb_erlang/test/factos_kurrentdb_erlang_test.gleam b/backends/factos_kurrentdb_erlang/test/factos_kurrentdb_erlang_test.gleam index 26cd640..90a2e16 100644 --- a/backends/factos_kurrentdb_erlang/test/factos_kurrentdb_erlang_test.gleam +++ b/backends/factos_kurrentdb_erlang/test/factos_kurrentdb_erlang_test.gleam @@ -47,8 +47,22 @@ pub fn dispatch_stream_handles_many_events_integration_test() { let stream_name = unique_name("counter-stream") let event_type = unique_name("FactosCounterIncremented") - let assert Ok(append_to_stream.Append(current_revision: 99, position: _)) = + let assert Ok(dispatch) = dispatch_counter_stream_many(stream_name, event_type, 100) + let assert append_to_stream.Append(current_revision: 99, position: _) = + dispatch.append + let assert [recorded] = dispatch.events + assert_counter_recorded( + recorded, + stream: stream_name, + revision: 99, + value: 100, + type_: factos.event_type(event_type), + ) + let reactor = factos.reactor(react: fn(recorded) { [recorded.event] }) + assert factos.react_all(reactor: reactor, events: dispatch.events) == [ + Incremented(100), + ] let assert Ok(loaded) = factos_kurrentdb_erlang.load_stream( @@ -73,8 +87,22 @@ pub fn read_context_handles_many_streams_integration_test() { ]), ]) - let assert Ok(append_to_stream.Append(current_revision: 0, position: _)) = + let assert Ok(dispatch) = dispatch_counter_context_streams_many(event_type, 50) + let assert append_to_stream.Append(current_revision: 0, position: _) = + dispatch.append + let assert [recorded] = dispatch.events + assert_counter_recorded( + recorded, + stream: event_type <> "-counter-context-1", + revision: 0, + value: 1, + type_: factos.event_type(event_type), + ) + let reactor = factos.reactor(react: fn(recorded) { [recorded.event] }) + assert factos.react_all(reactor: reactor, events: dispatch.events) == [ + Incremented(1), + ] let assert Ok(context) = factos_kurrentdb_erlang.read_context( @@ -112,7 +140,7 @@ fn dispatch_counter_stream_many( event_type: String, remaining: Int, ) -> Result( - append_to_stream.Append, + factos_kurrentdb_erlang.Dispatch(CounterEvent), factos_kurrentdb_erlang.Error(Nil, DecodeError), ) { let result = @@ -137,13 +165,13 @@ fn dispatch_counter_context_streams_many( event_type: String, remaining: Int, ) -> Result( - append_to_stream.Append, + factos_kurrentdb_erlang.Dispatch(CounterEvent), factos_kurrentdb_erlang.Error(Nil, DecodeError), ) { let result = factos_kurrentdb_erlang.dispatch_stream( connection(), - stream: unique_name("counter-context"), + stream: event_type <> "-counter-context-" <> int.to_string(remaining), decider: counter_decider(), codec: counter_codec(event_type), command: Increment, @@ -243,6 +271,23 @@ fn decode_counter_event( } } +fn assert_counter_recorded( + recorded: factos.Recorded(CounterEvent), + stream stream_name: String, + revision revision: Int, + value value: Int, + type_ type_: factos.EventType, +) -> Nil { + assert recorded.stream == stream_name + assert recorded.revision == revision + assert recorded.position != factos.NoPosition + assert recorded.type_ == type_ + assert recorded.version == 1 + assert recorded.tags == [factos.tag("counter:load")] + assert recorded.metadata == factos.empty_metadata() + assert recorded.event == Incremented(value) +} + fn unique_name(prefix: String) -> String { prefix <> "-" <> uuid.to_string(uuid.v4()) } diff --git a/backends/factos_pog/README.md b/backends/factos_pog/README.md index 12f0140..e670c61 100644 --- a/backends/factos_pog/README.md +++ b/backends/factos_pog/README.md @@ -13,7 +13,7 @@ instead of decoding unrelated rows in the application. PostgreSQL does not need to understand payloads, but any payload value needed by future context queries must be exposed as a tag when the event is written. -`dispatch_context` runs inside a PostgreSQL transaction and locks the event +`dispatch_with_query` runs inside a PostgreSQL transaction and locks the event table before reading, deciding, checking, and appending. This is intentionally conservative. It makes arbitrary `FailIfEventsMatch(query, after)` checks correct without trying to infer lock keys from dynamic query metadata. A @@ -21,13 +21,13 @@ higher-throughput backend could replace the table lock with advisory locks or more granular query-specific locks, but only if it preserves the same context-stability guarantee. -`dispatch_stream` is also available for applications where one stream revision really is the intended consistency boundary. It is an implementation strategy, not the definition of Event Sourcing. - -`factos/factos_pog/outbox` provides the PostgreSQL counterpart to the Cloudflare outbox helpers. `migrate` creates the `event_outbox` table so codec side effects can persist retryable work after a successful append. +`dispatch` is also available for applications where one stream revision really +is the intended consistency boundary. It is an implementation strategy, not the +definition of Event Sourcing. ## Usage -Start a `pog` pool in your application supervision tree, run `migrate`, build a codec with `factos_pog.codec`, then call `dispatch_with_query`/`dispatch_context` or `dispatch`/`dispatch_stream` with your domain decider and command. +Start a `pog` pool in your application supervision tree, run `migrate`, build a codec with `factos_pog.codec`, then call `dispatch_with_query` or `dispatch` with your domain decider and command. ```gleam let connection = pog.named_connection(pool_name) diff --git a/backends/factos_pog/compose.yml b/backends/factos_pog/compose.yml index 504ee59..031d187 100644 --- a/backends/factos_pog/compose.yml +++ b/backends/factos_pog/compose.yml @@ -6,7 +6,7 @@ services: POSTGRES_USER: postgres POSTGRES_PASSWORD: postgres ports: - - "5432:5432" + - "55432:5432" healthcheck: test: ["CMD-SHELL", "pg_isready -U postgres -d factos_pog"] interval: 1s diff --git a/backends/factos_pog/src/factos/factos_pog.gleam b/backends/factos_pog/src/factos/factos_pog.gleam index 00a5ce0..2b581fd 100644 --- a/backends/factos_pog/src/factos/factos_pog.gleam +++ b/backends/factos_pog/src/factos/factos_pog.gleam @@ -63,12 +63,9 @@ pub type EventCodec(event) { /// `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`. - /// `side_effects` run after a successful append and receive the exact domain - /// events that were just persisted. EventCodec( encode: fn(event) -> Proposed(event), decode: fn(StoredEvent) -> Result(factos.Decoded(event), DecodeError), - side_effects: List(fn(List(event)) -> Nil), ) } @@ -81,6 +78,15 @@ pub type Append { 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) { /// The decider rejected the command with a domain error. DomainError(domain_error) @@ -112,14 +118,11 @@ type QueryParameter { /// `encode` turns a domain event into a `Proposed` event ready for persistence. /// `decode` turns a stored row back into a domain event. Decode failures are /// returned as `DecodeError` and stop load/read flows rather than panicking. -/// `side_effects` run after a successful append. They are outside the database -/// transaction; use an outbox table if the effect must be retried. pub fn codec( encode encode: fn(event) -> Proposed(event), decode decode: fn(StoredEvent) -> Result(factos.Decoded(event), DecodeError), - side_effects side_effects: List(fn(List(event)) -> Nil), ) -> EventCodec(event) { - EventCodec(encode:, decode:, side_effects:) + EventCodec(encode:, decode:) } /// The schema is an append-only `factos_events` table with a global identity @@ -197,28 +200,6 @@ pub fn migrate(connection: pog.Connection) -> Result(Nil, Error(_)) { on conflict do nothing ", )) - use _ <- result.try(execute_migration( - connection, - " - create table if not exists event_outbox ( - id bigint generated always as identity primary key, - stream text not null, - event_type text not null, - payload text not null, - status text not null default 'pending', - error text, - created_at bigint not null default extract(epoch from now())::bigint, - processed_at bigint - ) - ", - )) - use _ <- result.try(execute_migration( - connection, - " - create index if not exists event_outbox_status_id - on event_outbox(status, id) - ", - )) Ok(Nil) } @@ -280,7 +261,7 @@ pub fn dispatch_with_query( decider decider: factos.Decider(command, state, event, domain_error), codec codec: EventCodec(event), command command: command, -) -> Result(Append, Error(domain_error)) { +) -> Result(Dispatch(event), Error(domain_error)) { use transaction_connection <- run_locked_transaction(connection) use context <- result.try(read_context( transaction_connection, @@ -294,15 +275,13 @@ pub fn dispatch_with_query( ) let #(context, events) = pair - use append <- result.try(append_with_condition( + append_with_condition( transaction_connection, stream_name, events, codec, context.append_condition, - )) - run_side_effects(codec, events) - Ok(append) + ) } /// Load and fold one stream. @@ -335,14 +314,14 @@ pub fn load_stream( /// /// Use this when one stream is intentionally the consistency boundary. It remains /// useful, but it is not required by Event Sourcing. For command-specific rules, -/// prefer `dispatch_context` so the protected boundary follows the decision. +/// prefer `dispatch_with_query` so the protected boundary follows the decision. pub fn dispatch( connection: pog.Connection, stream stream_name: String, decider decider: factos.Decider(command, state, event, domain_error), codec codec: EventCodec(event), command command: command, -) -> Result(Append, Error(domain_error)) { +) -> Result(Dispatch(event), Error(domain_error)) { use transaction_connection <- run_locked_transaction(connection) use loaded <- result.try(load_stream( transaction_connection, @@ -356,27 +335,19 @@ pub fn dispatch( |> result.map_error(DomainError), ) - use append <- result.try(append_stream_events( + append_stream_events( transaction_connection, stream_name, events, codec, loaded.revision, - )) - run_side_effects(codec, events) - Ok(append) -} - -fn run_side_effects(codec: EventCodec(event), events: List(event)) -> Nil { - let EventCodec(_, _, side_effects) = codec - use side_effect <- list.each(side_effects) - side_effect(events) + ) } fn run_locked_transaction( connection: pog.Connection, - work: fn(pog.Connection) -> Result(Append, Error(domain_error)), -) -> Result(Append, Error(domain_error)) { + work: fn(pog.Connection) -> Result(Dispatch(event), Error(domain_error)), +) -> Result(Dispatch(event), Error(domain_error)) { case { use transaction_connection <- pog.transaction(connection) @@ -389,7 +360,7 @@ fn run_locked_transaction( work(transaction_connection) } { - Ok(append) -> Ok(append) + Ok(dispatch) -> Ok(dispatch) Error(pog.TransactionQueryError(error)) -> Error(StoreError(error)) Error(pog.TransactionRolledBack(error)) -> Error(error) } @@ -401,7 +372,7 @@ fn append_with_condition( events: List(event), codec: EventCodec(event), condition: factos.AppendCondition, -) -> Result(Append, Error(domain_error)) { +) -> Result(Dispatch(event), Error(domain_error)) { case condition { factos.NoAppendCondition -> append_current_stream(connection, stream_name, events, codec) @@ -420,7 +391,7 @@ fn append_current_stream( stream_name: String, events: List(event), codec: EventCodec(event), -) -> Result(Append, Error(domain_error)) { +) -> Result(Dispatch(event), Error(domain_error)) { use revision <- result.try( current_revision(connection, stream_name) |> result.map_error(StoreError), @@ -440,13 +411,15 @@ fn append_stream_events( events: List(event), codec: EventCodec(event), expected: factos.Revision, -) -> Result(Append, Error(domain_error)) { +) -> Result(Dispatch(event), Error(domain_error)) { case events { - [] -> - Ok(Append( + [] -> { + 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) @@ -462,6 +435,7 @@ fn append_stream_events( codec, current + 1, factos.NoPosition, + [], ) } } @@ -475,11 +449,15 @@ fn insert_events( codec: EventCodec(event), revision: Int, position: factos.SequencePosition, -) -> Result(Append, Error(domain_error)) { + recorded_events: List(factos.Recorded(event)), +) -> Result(Dispatch(event), Error(domain_error)) { case events { - [] -> Ok(Append(current_revision: revision - 1, position: position)) + [] -> { + let append = Append(current_revision: revision - 1, position: position) + Ok(Dispatch(append:, events: list.reverse(recorded_events))) + } [event, ..rest] -> { - let EventCodec(encode, _, _) = codec + let EventCodec(encode, _) = codec let Proposed(id, _, type_, version, tags, metadata, data) = encode(event) use returned <- result.try( pog.query( @@ -505,6 +483,17 @@ fn insert_events( [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, + ) use _ <- result.try(insert_event_tags(connection, position, tags)) insert_events( connection, @@ -513,6 +502,7 @@ fn insert_events( codec, revision + 1, position, + [recorded, ..recorded_events], ) } } @@ -578,7 +568,7 @@ fn decode_row( row: StoredEvent, codec: EventCodec(event), ) -> Result(factos.Recorded(event), Error(domain_error)) { - let EventCodec(_, decode_event, _) = codec + 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 diff --git a/backends/factos_pog/src/factos/factos_pog/outbox.gleam b/backends/factos_pog/src/factos/factos_pog/outbox.gleam deleted file mode 100644 index 01712e1..0000000 --- a/backends/factos_pog/src/factos/factos_pog/outbox.gleam +++ /dev/null @@ -1,188 +0,0 @@ -//// Generic event outbox helpers for PostgreSQL. -//// -//// Provides insertion, read, status-update, and decoding helpers for rows in -//// the `event_outbox` table created by `factos_pog.migrate`. - -import gleam/dynamic/decode -import gleam/int -import gleam/list -import gleam/option -import gleam/result -import gleam/time/timestamp -import pog - -/// A full row from the `event_outbox` table. -pub type OutboxRecord { - OutboxRecord( - id: Int, - stream: String, - event_type: String, - payload: String, - status: String, - error: String, - created_at: Int, - processed_at: Int, - ) -} - -/// Errors returned by outbox reads and status updates. -pub type Error { - StoreError(pog.QueryError) - RowDecodeError(List(decode.DecodeError)) -} - -/// A pending side-effect row to insert into the outbox. -pub opaque type Entry { - Entry(stream: String, event_type: String, payload: String) -} - -/// Construct an outbox entry for later insertion. -pub fn entry( - stream stream: String, - event_type event_type: String, - payload payload: String, -) -> Entry { - Entry(stream:, event_type:, payload:) -} - -/// Insert outbox entries and return their generated IDs. -pub fn insert( - connection: pog.Connection, - entries: List(Entry), -) -> Result(List(Int), Error) { - insert_entries(connection, entries, []) -} - -/// Fetch a single pending outbox row by ID. -pub fn read_pending_outbox( - connection: pog.Connection, - id: Int, -) -> Result(option.Option(OutboxRecord), Error) { - use returned <- result.try( - pog.query( - "select id, stream, event_type, payload, status, error, created_at, processed_at - from event_outbox - where id = $1 and status = 'pending' - limit 1", - ) - |> pog.parameter(pog.int(id)) - |> pog.returning(outbox_record_decoder()) - |> pog.execute(on: connection) - |> result.map_error(StoreError), - ) - - case returned.rows { - [] -> Ok(option.None) - [row, ..] -> Ok(option.Some(row)) - } -} - -/// Mark an outbox row as successfully processed. -pub fn mark_outbox_sent( - connection: pog.Connection, - id: Int, - processed_at: timestamp.Timestamp, -) -> Result(Nil, Error) { - mark_outbox_processed(connection, id, "sent", "", processed_at) -} - -/// Mark an outbox row as failed with an error message. -pub fn mark_outbox_failed( - connection: pog.Connection, - id: Int, - error_message: String, - processed_at: timestamp.Timestamp, -) -> Result(Nil, Error) { - mark_outbox_processed(connection, id, "failed", error_message, processed_at) -} - -fn insert_entries( - connection: pog.Connection, - entries: List(Entry), - ids: List(Int), -) -> Result(List(Int), Error) { - case entries { - [] -> Ok(list.reverse(ids)) - [entry, ..rest] -> { - let Entry(stream, event_type, payload) = entry - use returned <- result.try( - pog.query( - "insert into event_outbox (stream, event_type, payload) - values ($1, $2, $3) - returning id", - ) - |> pog.parameter(pog.text(stream)) - |> pog.parameter(pog.text(event_type)) - |> pog.parameter(pog.text(payload)) - |> pog.returning(int_field_decoder()) - |> pog.execute(on: connection) - |> result.map_error(StoreError), - ) - let ids = case returned.rows { - [] -> ids - [id, ..] -> [id, ..ids] - } - insert_entries(connection, rest, ids) - } - } -} - -fn mark_outbox_processed( - connection: pog.Connection, - id: Int, - status: String, - error_message: String, - processed_at: timestamp.Timestamp, -) -> Result(Nil, Error) { - let processed_at_seconds = - processed_at - |> timestamp.to_unix_seconds_and_nanoseconds - |> fn(pair) { pair.0 } - - pog.query( - "update event_outbox - set status = $1, error = $2, processed_at = $3 - where id = $4", - ) - |> pog.parameter(pog.text(status)) - |> pog.parameter(pog.text(error_message)) - |> pog.parameter(pog.int(processed_at_seconds)) - |> pog.parameter(pog.int(id)) - |> pog.execute(on: connection) - |> result.map(fn(_) { Nil }) - |> result.map_error(StoreError) -} - -fn outbox_record_decoder() -> decode.Decoder(OutboxRecord) { - use id <- decode.field(0, decode.int) - use stream <- decode.field(1, decode.string) - use event_type <- decode.field(2, decode.string) - use payload <- decode.field(3, decode.string) - use status <- decode.field(4, decode.string) - use error <- decode.field(5, decode.optional(decode.string)) - use created_at <- decode.field(6, decode.int) - use processed_at <- decode.field(7, decode.optional(decode.int)) - decode.success(OutboxRecord( - id:, - stream:, - event_type:, - payload:, - status:, - error: option.unwrap(error, ""), - created_at:, - processed_at: option.unwrap(processed_at, 0), - )) -} - -fn int_field_decoder() -> decode.Decoder(Int) { - use value <- decode.field(0, decode.int) - decode.success(value) -} - -pub fn error_to_string(error: Error) -> String { - case error { - StoreError(_) -> "store error" - RowDecodeError(errors) -> - "database row decode error: " <> int.to_string(list.length(errors)) - } -} diff --git a/backends/factos_pog/test/factos_pog_test.gleam b/backends/factos_pog/test/factos_pog_test.gleam index 887b0d3..1e0d1a1 100644 --- a/backends/factos_pog/test/factos_pog_test.gleam +++ b/backends/factos_pog/test/factos_pog_test.gleam @@ -1,6 +1,5 @@ import factos import factos/factos_pog -import factos/factos_pog/outbox import gleam/bit_array import gleam/erlang/process import gleam/int @@ -17,7 +16,6 @@ pub fn main() -> Nil { type Command { RegisterUser(username: String) - RegisterUsers(first: String, second: String) } type Event { @@ -49,15 +47,12 @@ type TestGlobalData { TestGlobalData(connection: pog.Connection) } -type SideEffectMessage { - SideEffect(events: List(Event)) -} pub fn dispatch_alias_persists_events_test() { let TestGlobalData(connection) = global_data() reset_schema(connection) - let assert Ok(factos_pog.Append(current_revision: 0, position: _)) = + let assert Ok(dispatch) = factos_pog.dispatch( connection, stream: "user-renata", @@ -66,6 +61,23 @@ pub fn dispatch_alias_persists_events_test() { command: RegisterUser("renata"), ) + let assert factos_pog.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_pog.load_stream( connection, @@ -85,10 +97,7 @@ pub fn dispatch_with_query_filters_before_decoding_unknown_events_test() { let query = username_query("renata") - let assert Ok(factos_pog.Append( - current_revision: 0, - position: factos.SequencePosition(_), - )) = + let assert Ok(dispatch) = factos_pog.dispatch_with_query( connection, stream: "user-renata", @@ -98,6 +107,23 @@ pub fn dispatch_with_query_filters_before_decoding_unknown_events_test() { command: RegisterUser("renata"), ) + let assert factos_pog.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(context) = factos_pog.read_context( connection, @@ -112,7 +138,7 @@ pub fn dispatch_with_query_filters_before_decoding_unknown_events_test() { assert context.position != factos.NoPosition } -pub fn dispatch_context_handles_many_streams_test() { +pub fn dispatch_with_query_handles_many_streams_test() { let TestGlobalData(connection) = global_data() reset_schema(connection) @@ -123,10 +149,24 @@ pub fn dispatch_context_handles_many_streams_test() { ]), ]) - let assert Ok(factos_pog.Append( + let assert Ok(dispatch) = dispatch_counter_context_many(connection, query, 25) + let assert factos_pog.Append( current_revision: 0, position: factos.SequencePosition(_), - )) = dispatch_counter_context_many(connection, query, 25) + ) = dispatch.append + let assert [recorded] = dispatch.events + assert_counter_recorded( + recorded, + stream: "counter-context-1", + revision: 0, + position: dispatch.append.position, + value: 25, + type_: factos.event_type("Incremented"), + ) + let reactor = factos.reactor(react: fn(recorded) { [recorded.event] }) + assert factos.react_all(reactor: reactor, events: dispatch.events) == [ + Incremented(25), + ] let assert Ok(context) = factos_pog.read_context( @@ -140,68 +180,6 @@ pub fn dispatch_context_handles_many_streams_test() { assert list.length(context.events) == 25 } -pub fn dispatch_alias_runs_side_effect_after_successful_append_test() { - let TestGlobalData(connection) = global_data() - reset_schema(connection) - let side_effect_messages = process.new_subject() - - let assert Ok(append) = - factos_pog.dispatch( - connection, - stream: "user-side-effect", - decider: decider(), - codec: codec_with_side_effect(fn(events) { - process.send(side_effect_messages, SideEffect(events)) - }), - command: RegisterUsers("side-effect-a", "side-effect-b"), - ) - - let assert Ok(loaded) = - factos_pog.load_stream( - connection, - stream: "user-side-effect", - decider: decider(), - codec: codec(), - ) - let assert [first, second] = loaded.events - let assert Ok(SideEffect(side_effect_events)) = - process.receive(side_effect_messages, within: 0) - - assert append.current_revision == 1 - assert append.position == second.position - assert side_effect_events == [first.event, second.event] -} - -pub fn outbox_insert_returns_readable_pending_records_test() { - let TestGlobalData(connection) = global_data() - reset_schema(connection) - - let assert Ok([first_id, second_id]) = - outbox.insert(connection, [ - outbox.entry( - stream: "outbox-stream-1", - event_type: "UserRegistered", - payload: "renata", - ), - outbox.entry( - stream: "outbox-stream-2", - event_type: "UserRegistered", - payload: "ada", - ), - ]) - - let assert Ok(Some(first)) = outbox.read_pending_outbox(connection, first_id) - let assert Ok(Some(second)) = - outbox.read_pending_outbox(connection, second_id) - - assert first.id == first_id - assert first.stream == "outbox-stream-1" - assert first.event_type == "UserRegistered" - assert first.payload == "renata" - assert second.id == second_id - assert second.stream == "outbox-stream-2" - assert second.payload == "ada" -} fn global_data() -> TestGlobalData { global_value.create_with_unique_name("factos_pog_test.global.data", fn() { @@ -214,7 +192,7 @@ fn start_test_connection() -> pog.Connection { let config = pog.default_config(pool_name) |> pog.host("127.0.0.1") - |> pog.port(5432) + |> pog.port(55432) |> pog.database("factos_pog") |> pog.user("postgres") |> pog.password(Some("postgres")) @@ -226,9 +204,6 @@ fn start_test_connection() -> pog.Connection { } fn reset_schema(connection: pog.Connection) -> Nil { - let assert Ok(_) = - pog.query("drop table if exists event_outbox") - |> pog.execute(on: connection) let assert Ok(_) = pog.query("drop table if exists factos_event_tags") |> pog.execute(on: connection) @@ -274,6 +249,24 @@ fn username_query(username: String) -> factos.Query { ]) } +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 decider() -> factos.Decider(Command, State, Event, DomainError) { factos.decider(initial: Available, decide:, evolve:) } @@ -281,13 +274,7 @@ fn decider() -> factos.Decider(Command, State, Event, DomainError) { fn decide(state: State, command: Command) -> Result(List(Event), DomainError) { case state, command { Available, RegisterUser(username) -> Ok([UserRegistered(username)]) - Available, RegisterUsers(first, second) -> - Ok([ - UserRegistered(first), - UserRegistered(second), - ]) Taken, RegisterUser(_) -> Error(AlreadyTaken) - Taken, RegisterUsers(_, _) -> Error(AlreadyTaken) } } @@ -296,13 +283,7 @@ fn evolve(_state: State, _event: Event) -> State { } fn codec() -> factos_pog.EventCodec(Event) { - factos_pog.codec(encode:, decode:, side_effects: []) -} - -fn codec_with_side_effect( - side_effect: fn(List(Event)) -> Nil, -) -> factos_pog.EventCodec(Event) { - factos_pog.codec(encode:, decode:, side_effects: [side_effect]) + factos_pog.codec(encode:, decode:) } fn encode(event: Event) -> factos_pog.Proposed(Event) { @@ -342,7 +323,7 @@ fn dispatch_counter_context_many( connection: pog.Connection, query: factos.Query, remaining: Int, -) -> Result(factos_pog.Append, factos_pog.Error(Nil)) { +) -> Result(factos_pog.Dispatch(CounterEvent), factos_pog.Error(Nil)) { case remaining { 0 -> factos_pog.dispatch_with_query( @@ -408,7 +389,6 @@ fn counter_codec() -> factos_pog.EventCodec(CounterEvent) { factos_pog.codec( encode: encode_counter_event, decode: decode_counter_event, - side_effects: [], ) } @@ -453,3 +433,22 @@ fn decode_counter_event( _ -> Error(factos_pog.UnknownEvent) } } + +fn assert_counter_recorded( + recorded: factos.Recorded(CounterEvent), + stream stream_name: String, + revision revision: Int, + position position: factos.SequencePosition, + value value: Int, + type_ type_: factos.EventType, +) -> 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_ == type_ + 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/backends/factos_sqlight/src/factos/factos_sqlight.gleam b/backends/factos_sqlight/src/factos/factos_sqlight.gleam index b2c5f61..e28d3e6 100644 --- a/backends/factos_sqlight/src/factos/factos_sqlight.gleam +++ b/backends/factos_sqlight/src/factos/factos_sqlight.gleam @@ -81,6 +81,15 @@ pub type Append { 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) @@ -172,7 +181,7 @@ pub fn dispatch_context( decider decider: factos.Decider(command, state, event, domain_error), codec codec: EventCodec(event, decode_error), command command: command, -) -> Result(Append, Error(domain_error, decode_error)) { +) -> Result(Dispatch(event), Error(domain_error, decode_error)) { use _ <- result.try( sqlight.exec("begin immediate", on: connection) |> result.map_error(StoreError), @@ -243,7 +252,7 @@ pub fn dispatch_stream( decider decider: factos.Decider(command, state, event, domain_error), codec codec: EventCodec(event, decode_error), command command: command, -) -> Result(Append, Error(domain_error, decode_error)) { +) -> Result(Dispatch(event), Error(domain_error, decode_error)) { use _ <- result.try( sqlight.exec("begin immediate", on: connection) |> result.map_error(StoreError), @@ -276,12 +285,12 @@ pub fn dispatch_stream( fn finish_transaction( connection: sqlight.Connection, - result: Result(Append, Error(domain_error, decode_error)), -) -> Result(Append, Error(domain_error, decode_error)) { + result: Result(Dispatch(event), Error(domain_error, decode_error)), +) -> Result(Dispatch(event), Error(domain_error, decode_error)) { case result { - Ok(append) -> + Ok(dispatch) -> case sqlight.exec("commit", on: connection) { - Ok(Nil) -> Ok(append) + Ok(Nil) -> Ok(dispatch) Error(error) -> Error(StoreError(error)) } Error(error) -> { @@ -297,7 +306,7 @@ fn append_with_condition( events: List(event), codec: EventCodec(event, decode_error), condition: factos.AppendCondition, -) -> Result(Append, Error(domain_error, decode_error)) { +) -> Result(Dispatch(event), Error(domain_error, decode_error)) { case condition { factos.NoAppendCondition -> append_current_stream(connection, stream_name, events, codec) @@ -316,7 +325,7 @@ fn append_current_stream( stream_name: String, events: List(event), codec: EventCodec(event, decode_error), -) -> Result(Append, Error(domain_error, decode_error)) { +) -> Result(Dispatch(event), Error(domain_error, decode_error)) { use revision <- result.try( current_revision(connection, stream_name) |> result.map_error(StoreError), @@ -336,13 +345,15 @@ fn append_stream_events( events: List(event), codec: EventCodec(event, decode_error), expected: factos.Revision, -) -> Result(Append, Error(domain_error, decode_error)) { +) -> Result(Dispatch(event), Error(domain_error, decode_error)) { case events { - [] -> - Ok(Append( + [] -> { + 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) @@ -358,6 +369,7 @@ fn append_stream_events( codec, current + 1, factos.NoPosition, + [], ) } } @@ -371,9 +383,13 @@ fn insert_events( codec: EventCodec(event, decode_error), revision: Int, position: factos.SequencePosition, -) -> Result(Append, Error(domain_error, decode_error)) { + recorded_events: List(factos.Recorded(event)), +) -> Result(Dispatch(event), Error(domain_error, decode_error)) { case events { - [] -> Ok(Append(current_revision: revision - 1, position: position)) + [] -> { + 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) @@ -403,6 +419,17 @@ fn insert_events( [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, @@ -410,6 +437,7 @@ fn insert_events( codec, revision + 1, position, + [recorded, ..recorded_events], ) } } diff --git a/backends/factos_sqlight/test/factos_sqlight_test.gleam b/backends/factos_sqlight/test/factos_sqlight_test.gleam index 2caa161..7344078 100644 --- a/backends/factos_sqlight/test/factos_sqlight_test.gleam +++ b/backends/factos_sqlight/test/factos_sqlight_test.gleam @@ -49,7 +49,7 @@ pub fn dispatch_stream_persists_events_test() { use connection <- sqlight.with_connection(":memory:") let assert Ok(Nil) = factos_sqlight.migrate(connection) - let assert Ok(factos_sqlight.Append(current_revision: 0, position: _)) = + let assert Ok(dispatch) = factos_sqlight.dispatch_stream( connection, stream: "user-renata", @@ -58,6 +58,23 @@ pub fn dispatch_stream_persists_events_test() { 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, @@ -74,10 +91,23 @@ pub fn dispatch_stream_handles_many_events_test() { use connection <- sqlight.with_connection(":memory:") let assert Ok(Nil) = factos_sqlight.migrate(connection) - let assert Ok(factos_sqlight.Append( + let assert Ok(dispatch) = dispatch_counter_stream_many(connection, 250) + let assert factos_sqlight.Append( current_revision: 249, position: factos.SequencePosition(_), - )) = dispatch_counter_stream_many(connection, 250) + ) = 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( @@ -103,10 +133,23 @@ pub fn dispatch_context_handles_many_streams_test() { ]), ]) - let assert Ok(factos_sqlight.Append( + let assert Ok(dispatch) = dispatch_counter_context_many(connection, query, 100) + let assert factos_sqlight.Append( current_revision: 0, position: factos.SequencePosition(_), - )) = dispatch_counter_context_many(connection, query, 100) + ) = 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( @@ -173,10 +216,28 @@ fn decode( } } +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.Append, factos_sqlight.Error(Nil, DecodeError)) { +) -> Result(factos_sqlight.Dispatch(CounterEvent), factos_sqlight.Error(Nil, DecodeError)) { case remaining { 0 -> factos_sqlight.dispatch_stream( @@ -208,7 +269,7 @@ fn dispatch_counter_context_many( connection: sqlight.Connection, query: factos.Query, remaining: Int, -) -> Result(factos_sqlight.Append, factos_sqlight.Error(Nil, DecodeError)) { +) -> Result(factos_sqlight.Dispatch(CounterEvent), factos_sqlight.Error(Nil, DecodeError)) { case remaining { 0 -> factos_sqlight.dispatch_context( @@ -318,3 +379,21 @@ fn decode_counter_event( _ -> 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/examples/orders_sqlight/src/order_workflow.gleam b/examples/orders_sqlight/src/order_workflow.gleam index 010776b..8865775 100644 --- a/examples/orders_sqlight/src/order_workflow.gleam +++ b/examples/orders_sqlight/src/order_workflow.gleam @@ -360,7 +360,7 @@ fn dispatch_with_retry( command: Command, attempts attempts: Int, ) -> Result( - factos_sqlight.Append, + factos_sqlight.Dispatch(Event), factos_sqlight.Error(DomainError, DecodeError), ) { let result = dispatch(connection, order_id, command) @@ -411,7 +411,7 @@ fn dispatch( order_id: String, command: Command, ) -> Result( - factos_sqlight.Append, + factos_sqlight.Dispatch(Event), factos_sqlight.Error(DomainError, DecodeError), ) { factos_sqlight.dispatch_stream( diff --git a/examples/tickets_pog/src/ticket_sale.gleam b/examples/tickets_pog/src/ticket_sale.gleam deleted file mode 100644 index 02fdd3f..0000000 --- a/examples/tickets_pog/src/ticket_sale.gleam +++ /dev/null @@ -1,320 +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) -} - -pub type State { - TicketWindow(capacity: Int, sold: Int) -} - -pub type DomainError { - SoldOut(capacity: Int) -} - -pub type DecodeError { - UnknownEventType(String) - InvalidPayload(String) -} - -pub type SaleSummary { - SaleSummary(attempts: Int, accepted: Int, sold_out: Int, recorded_events: Int) -} - -type TestGlobalData { - TestGlobalData(connection: pog.Connection) -} - -type PurchaseMessage { - PurchaseFinished( - attempt: Int, - result: Result( - factos_pog.Append, - factos_pog.Error(DomainError, DecodeError), - ), - ) -} - -pub fn main() -> Nil { - case run() { - Ok(summary) -> - io.println( - "ticket sale completed: " - <> int.to_string(summary.accepted) - <> " accepted, " - <> int.to_string(summary.sold_out) - <> " sold out, " - <> int.to_string(summary.recorded_events) - <> " recorded events", - ) - Error(_) -> io.println("ticket sale failed") - } -} - -pub fn run() -> Result(SaleSummary, factos_pog.Error(DomainError, DecodeError)) { - let TestGlobalData(connection) = global_data() - use _ <- result.try(reset_schema(connection)) - - let workers = process.new_subject() - let initial = - SaleSummary(attempts: 0, accepted: 0, sold_out: 0, recorded_events: 0) - - { - use _, attempt <- int.range(from: 1, to: concurrency + 1, with: Nil) - spawn_purchase(workers, connection, attempt) - } - - collect_purchases( - workers, - connection: connection, - remaining: purchase_attempts, - next_attempt: concurrency + 1, - summary: initial, - ) -} - -fn global_data() -> TestGlobalData { - global_value.create_with_unique_name("tickets_pog.global.data", fn() { - TestGlobalData(connection: start_connection()) - }) -} - -fn start_connection() -> pog.Connection { - let pool_name = process.new_name("tickets_pog") - let config = - pog.default_config(pool_name) - |> pog.host("127.0.0.1") - |> pog.port(5433) - |> pog.database("tickets_pog") - |> pog.user("postgres") - |> pog.password(Some("postgres")) - |> pog.ssl(pog.SslDisabled) - |> pog.pool_size(postgres_pool_size) - - let assert Ok(_) = pog.start(config) - process.sleep(100) - pog.named_connection(pool_name) -} - -fn reset_schema( - connection: pog.Connection, -) -> Result(Nil, factos_pog.Error(DomainError, DecodeError)) { - use _ <- result.try( - pog.query("drop table if exists factos_events") - |> pog.execute(on: connection) - |> result.map(fn(_) { Nil }) - |> result.map_error(factos_pog.StoreError), - ) - factos_pog.migrate(connection) -} - -fn spawn_purchase( - workers: process.Subject(PurchaseMessage), - connection: pog.Connection, - attempt: Int, -) -> Nil { - process.spawn(fn() { - process.send( - workers, - PurchaseFinished(attempt, purchase(connection, attempt)), - ) - }) - Nil -} - -fn purchase( - connection: pog.Connection, - attempt: Int, -) -> Result(factos_pog.Append, factos_pog.Error(DomainError, DecodeError)) { - // Each buyer races through the same event-context query. PostgreSQL receives a - // large amount of concurrent work via the pool, while the backend's transaction - // lock preserves the capacity invariant for this arbitrary tag-based context. - factos_pog.dispatch_context( - connection, - stream: buyer_stream(attempt), - query: sale_query(), - decider: ticket_decider(), - codec: ticket_codec(), - command: BuyTicket(buyer_name(attempt)), - ) -} - -fn collect_purchases( - workers: process.Subject(PurchaseMessage), - connection connection: pog.Connection, - remaining remaining: Int, - next_attempt next_attempt: Int, - summary summary: SaleSummary, -) -> Result(SaleSummary, factos_pog.Error(DomainError, DecodeError)) { - case remaining { - 0 -> finalize_summary(connection, summary) - _ -> - case process.receive(workers, within: receive_timeout) { - Ok(PurchaseFinished(_, Ok(_))) -> { - case next_attempt <= purchase_attempts { - True -> spawn_purchase(workers, connection, next_attempt) - False -> Nil - } - collect_purchases( - workers, - connection: connection, - remaining: remaining - 1, - next_attempt: next_attempt + 1, - summary: SaleSummary( - attempts: summary.attempts + 1, - accepted: summary.accepted + 1, - sold_out: summary.sold_out, - recorded_events: summary.recorded_events, - ), - ) - } - Ok(PurchaseFinished(_, Error(factos_pog.DomainError(SoldOut(_))))) -> { - case next_attempt <= purchase_attempts { - True -> spawn_purchase(workers, connection, next_attempt) - False -> Nil - } - collect_purchases( - workers, - connection: connection, - remaining: remaining - 1, - next_attempt: next_attempt + 1, - summary: SaleSummary( - attempts: summary.attempts + 1, - accepted: summary.accepted, - sold_out: summary.sold_out + 1, - recorded_events: summary.recorded_events, - ), - ) - } - Ok(PurchaseFinished(_, Error(error))) -> Error(error) - Error(Nil) -> Error(factos_pog.DomainError(SoldOut(summary.accepted))) - } - } -} - -fn finalize_summary( - connection: pog.Connection, - summary: SaleSummary, -) -> Result(SaleSummary, factos_pog.Error(DomainError, DecodeError)) { - use context <- result.try(factos_pog.read_context( - connection, - query: sale_query(), - decider: ticket_decider(), - codec: ticket_codec(), - )) - - Ok(SaleSummary( - attempts: summary.attempts, - accepted: summary.accepted, - sold_out: summary.sold_out, - recorded_events: list.length(context.events), - )) -} - -fn sale_query() -> factos.Query { - factos.query([ - factos.query_item(types: [factos.event_type("TicketSold")], tags: [ - factos.tag("event:" <> event_id), - ]), - ]) -} - -fn ticket_decider() -> factos.Decider(Command, State, Event, DomainError) { - factos.decider( - initial: TicketWindow(capacity: ticket_capacity, sold: 0), - decide: decide, - evolve: evolve, - ) -} - -fn decide(state: State, command: Command) -> Result(List(Event), DomainError) { - let TicketWindow(capacity, sold) = state - case command { - BuyTicket(buyer) -> - case sold < capacity { - True -> Ok([TicketSold(buyer)]) - False -> Error(SoldOut(capacity)) - } - } -} - -fn evolve(state: State, event: Event) -> State { - let TicketWindow(capacity, sold) = state - case event { - TicketSold(_) -> TicketWindow(capacity: capacity, sold: sold + 1) - } -} - -fn ticket_codec() -> factos_pog.EventCodec(Event, DecodeError) { - factos_pog.codec(encode: encode_event, decode: decode_event, side_effects: []) -} - -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), DecodeError) { - case factos.event_type_name(stored.type_) { - "TicketSold" -> { - use buyer <- result.try( - bit_array.to_string(stored.data) - |> result.replace_error(InvalidPayload("buyer was not utf-8")), - ) - Ok(factos.Decoded( - event: TicketSold(buyer), - type_: stored.type_, - version: stored.version, - tags: stored.tags, - metadata: stored.metadata, - )) - } - type_name -> Error(UnknownEventType(type_name)) - } -} - -fn buyer_name(attempt: Int) -> String { - "buyer-" <> int.to_string(attempt) -} - -fn buyer_stream(attempt: Int) -> String { - "ticket-buyer-" <> int.to_string(attempt) -} diff --git a/examples/tickets_pog/src/tickets_pog.gleam b/examples/tickets_pog/src/tickets_pog.gleam index 6250e4a..af97e65 100644 --- a/examples/tickets_pog/src/tickets_pog.gleam +++ b/examples/tickets_pog/src/tickets_pog.gleam @@ -1,5 +1,353 @@ -import ticket_sale +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 { - ticket_sale.main() + 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 index 0d2914a..64879f1 100644 --- a/examples/tickets_pog/test/tickets_pog_test.gleam +++ b/examples/tickets_pog/test/tickets_pog_test.gleam @@ -1,15 +1,15 @@ import gleeunit -import ticket_sale +import tickets_pog pub fn main() -> Nil { gleeunit.main() } pub fn ticket_sale_preserves_capacity_under_high_concurrency_test() { - let assert Ok(ticket_sale.SaleSummary( + let assert Ok(tickets_pog.SaleSummary( attempts: 300, accepted: 100, sold_out: 200, recorded_events: 100, - )) = ticket_sale.run() + )) = tickets_pog.run() } diff --git a/gleam.toml b/gleam.toml index e07abcf..48cd5cd 100644 --- a/gleam.toml +++ b/gleam.toml @@ -1,7 +1,7 @@ name = "factos" version = "1.0.0" description = "Store-independent context-first Event Sourcing primitives for Gleam." -licences = ["Apache-2.0"] +licences = ["MIT"] links = [ { title = "Simply Event Sourcing", href = "https://ricofritzsche.me/simply-event-sourcing/" }, ] diff --git a/src/factos.gleam b/src/factos.gleam index c69c462..71a6828 100644 --- a/src/factos.gleam +++ b/src/factos.gleam @@ -123,6 +123,16 @@ pub type View(state, event) { View(initial: state, evolve: fn(state, event) -> state) } +pub type Reactor(event, effect) { + /// A pure event reaction. + /// + /// Reactors inspect committed recorded events and produce application-owned + /// 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. + Reactor(react: fn(Recorded(event)) -> List(effect)) +} + pub type Revision { /// A stream has no events. NoEvents @@ -277,6 +287,17 @@ pub fn view( View(initial:, evolve:) } +/// Build a pure event reactor. +/// +/// Reactors are the side-effect planning counterpart to views: they consume +/// committed recorded events and return application-owned effect values without +/// executing IO. +pub fn reactor( + react react: fn(Recorded(event)) -> List(effect), +) -> Reactor(event, effect) { + Reactor(react:) +} + /// Fold events with a decider and decide which new events a command produces. /// /// This is useful for unit tests and for in-memory command handling. It does not @@ -341,6 +362,36 @@ pub fn project_from( fold_events(state, events, evolve) } +/// Produce effect values for one committed recorded event. +pub fn react( + reactor reactor: Reactor(event, effect), + event event: Recorded(event), +) -> List(effect) { + let Reactor(react) = reactor + react(event) +} + +/// Produce effect values for committed recorded events, preserving event order. +pub fn react_all( + reactor reactor: Reactor(event, effect), + events events: List(Recorded(event)), +) -> List(effect) { + use event <- list.flat_map(events) + react(reactor, event) +} + +/// Merge two reactors that consume the same event type. +/// +/// The resulting reactor runs both reactors for each event and concatenates their +/// produced effects in argument order. +pub fn merge_reactors( + first first: Reactor(event, effect), + second second: Reactor(event, effect), +) -> Reactor(event, effect) { + use event <- Reactor + list.append(react(first, event), react(second, event)) +} + /// Merge two views that consume the same event type. /// /// The resulting view keeps both states in a tuple and evolves both for every diff --git a/test/factos_test.gleam b/test/factos_test.gleam index f8cfcbe..35f625e 100644 --- a/test/factos_test.gleam +++ b/test/factos_test.gleam @@ -199,6 +199,82 @@ pub fn merge_views_projects_same_events_into_tuple_state_test() { == #(1, 2) } +pub fn react_maps_one_recorded_event_into_effect_values_test() -> Nil { + let event = + recorded(UserRegistered("renata"), [factos.tag("username:renata")], revision: 5) + let reactor = + factos.reactor(react: fn(recorded) { + case recorded.event { + UserRegistered(username) -> [ + recorded.id <> ":" <> int.to_string(recorded.revision) <> ":" <> username, + ] + UsernameReserved(_) -> [] + DisplayNameChanged(_) -> [] + } + }) + + assert factos.react(reactor: reactor, event: event) + == ["event-5:5:renata"] +} + +pub fn react_all_flattens_reactions_in_recorded_event_order_test() -> Nil { + let reactor = + factos.reactor(react: fn(recorded) { + case recorded.event { + UserRegistered(username) -> [ + "registered:" <> username, + "welcome:" <> username, + ] + UsernameReserved(_) -> [] + DisplayNameChanged(name) -> ["display:" <> name] + } + }) + + assert factos.react_all(reactor: reactor, events: [ + recorded(UserRegistered("renata"), [], revision: 0), + recorded(UsernameReserved("renata"), [], revision: 1), + recorded(DisplayNameChanged("Rena"), [], revision: 2), + recorded(UserRegistered("lucy"), [], revision: 3), + ]) + == [ + "registered:renata", + "welcome:renata", + "display:Rena", + "registered:lucy", + "welcome:lucy", + ] +} + +pub fn merge_reactors_combines_outputs_for_the_same_recorded_event_test() -> Nil { + let event = + recorded(UserRegistered("renata"), [factos.tag("username:renata")], revision: 7) + let audit = + factos.reactor(react: fn(recorded) { + case recorded.event { + UserRegistered(username) -> [ + "audit:" <> recorded.id <> ":" <> username, + ] + UsernameReserved(_) -> [] + DisplayNameChanged(_) -> [] + } + }) + let notification = + factos.reactor(react: fn(recorded) { + case recorded.event { + UserRegistered(username) -> [ + "notify:" <> int.to_string(recorded.revision) <> ":" <> username, + ] + UsernameReserved(_) -> [] + DisplayNameChanged(_) -> [] + } + }) + + let merged = factos.merge_reactors(audit, notification) + + assert factos.react(reactor: merged, event: event) + == ["audit:event-7:renata", "notify:7:renata"] +} + fn evolve(state: State, event: Event) -> State { case state, event { UsernameAvailable, UsernameReserved(_) -> UsernameTaken