diff --git a/backends/factos_sqlight/test/factos_sqlight_test.gleam b/backends/factos_sqlight/test/factos_sqlight_test.gleam index 5b83b5e..7276dce 100644 --- a/backends/factos_sqlight/test/factos_sqlight_test.gleam +++ b/backends/factos_sqlight/test/factos_sqlight_test.gleam @@ -3,6 +3,7 @@ import factos/factos_sqlight import gleam/dynamic/decode import gleam/erlang/application import gleam/erlang/process +import gleam/int import gleam/json import gleam/list import gleam/result @@ -42,6 +43,18 @@ type FireMessage { ) } +type CounterCommand { + Increment +} + +type CounterEvent { + Incremented(value: Int) +} + +type CounterState { + CounterState(total: Int) +} + pub fn bootstrap_and_dbmate_migrations_create_streamless_schema_test() -> Nil { use bootstrap_connection <- sqlight.with_connection(":memory:") let assert Ok(Nil) = factos_sqlight.migrate(bootstrap_connection) @@ -349,6 +362,211 @@ pub fn fire_and_forget_continues_after_callback_error_in_append_order_test() -> Nil } +pub fn dispatch_builder_with_one_retry_attempt_persists_events_test() -> Nil { + use connection <- sqlight.with_connection(":memory:") + execute_migration_file(connection) + + let assert Ok(dispatch) = + factos.new_dispatch( + connection:, + decision_context: username_context("renata"), + decider: decider(), + encode: encode_event, + decode: decode_event, + ) + |> factos.with_retry_attempts(attempts: 1) + |> factos_sqlight.dispatch(RegisterUser(username: "renata"), event_id: fn() { + "single-retry" + }) + + let assert factos.SequencePosition(_) = dispatch.position + let assert [recorded] = dispatch.events + assert_user_recorded( + recorded, + position: dispatch.position, + username: "renata", + ) + + let assert Ok(context) = + factos_sqlight.read( + connection, + username_context("renata"), + decider(), + decode_event, + ) + assert context.state == State(usernames: ["renata"]) + assert context.events == [recorded] + Nil +} + +pub fn dispatch_builder_with_query_filters_before_decoding_unknown_events_test() -> Nil { + use connection <- sqlight.with_connection(":memory:") + execute_migration_file(connection) + insert_unknown_event(connection) + + let assert Ok(dispatch) = + dispatch_user(connection, "renata", event_id: "filtered-dispatch") + + let assert factos.SequencePosition(_) = dispatch.position + let assert [recorded] = dispatch.events + assert_user_recorded( + recorded, + position: dispatch.position, + username: "renata", + ) + + let assert Ok(context) = + factos_sqlight.read( + connection, + username_context("renata"), + decider(), + decode_event, + ) + let assert [event] = context.events + assert event.event.payload == UserRegistered(username: "renata") + assert context.state == State(usernames: ["renata"]) + assert context.position != factos.NoPosition + Nil +} + +pub fn dispatch_builder_with_query_handles_many_events_test() -> Nil { + use connection <- sqlight.with_connection(":memory:") + execute_migration_file(connection) + + let query = + factos.Matching(items: [ + factos.item(types: [factos.event_type("Incremented")], tags: [ + factos.tag("counter:load"), + ]), + ]) + + let assert Ok(dispatch) = dispatch_counter_context_many(connection, query, 25) + let assert factos.SequencePosition(_) = dispatch.position + let assert [recorded] = dispatch.events + assert_counter_recorded( + recorded, + position: dispatch.position, + value: 25, + type_: factos.event_type("Incremented"), + ) + + let assert Ok(context) = + factos_sqlight.read( + connection, + query, + counter_decider(), + decode_counter_event, + ) + assert context.state == CounterState(25) + assert list.length(context.events) == 25 + Nil +} + +pub fn context_semantics_conformance_test() -> Nil { + use connection <- sqlight.with_connection(":memory:") + execute_migration_file(connection) + + let decision_context = empty_query() + let assert Ok(renata_dispatch) = + factos.new_dispatch( + connection:, + decision_context:, + decider: accepting_decider(), + encode: encode_event, + decode: decode_event, + ) + |> factos_sqlight.dispatch(RegisterUser(username: "renata"), event_id: fn() { + "conformance-renata" + }) + let assert Ok(lucy_dispatch) = + factos.new_dispatch( + connection:, + decision_context:, + decider: accepting_decider(), + encode: encode_event, + decode: decode_event, + ) + |> factos_sqlight.dispatch(RegisterUser(username: "lucy"), event_id: fn() { + "conformance-lucy" + }) + let assert Ok(marc_dispatch) = + factos.new_dispatch( + connection:, + decision_context: factos.AllEvents, + decider: accepting_decider(), + encode: encode_event, + decode: decode_event, + ) + |> factos_sqlight.dispatch(RegisterUser(username: "marc"), event_id: fn() { + "conformance-marc" + }) + let renata_position = renata_dispatch.position + let lucy_position = lucy_dispatch.position + let marc_position = marc_dispatch.position + + let assert Ok(empty_context) = + factos_sqlight.read(connection, decision_context, decider(), decode_event) + assert empty_context.state == State(usernames: []) + assert empty_context.events == [] + assert empty_context.position == factos.NoPosition + assert empty_context.append_condition + == factos.FailIfEventsMatch(decision_context:, after: factos.NoPosition) + + let compound_query = username_conformance_query() + let assert Ok(compound_context) = + factos_sqlight.read(connection, compound_query, decider(), decode_event) + let assert [lucy] = compound_context.events + assert_user_recorded(lucy, position: lucy_position, username: "lucy") + assert compound_context.state == State(usernames: ["lucy"]) + assert compound_context.position == lucy_position + assert compound_context.append_condition + == factos.FailIfEventsMatch( + decision_context: compound_query, + after: lucy_position, + ) + + let assert Ok(all_context) = + factos_sqlight.read(connection, factos.AllEvents, decider(), decode_event) + let assert [renata, lucy, marc] = all_context.events + assert_user_recorded(renata, position: renata_position, username: "renata") + assert_user_recorded(lucy, position: lucy_position, username: "lucy") + assert_user_recorded(marc, position: marc_position, username: "marc") + assert all_context.state == State(usernames: ["renata", "lucy", "marc"]) + assert all_context.position == marc_position + assert all_context.append_condition + == factos.FailIfEventsMatch( + decision_context: factos.AllEvents, + after: marc_position, + ) + Nil +} + +pub fn no_context_dispatch_skips_existing_events_test() -> Nil { + use connection <- sqlight.with_connection(":memory:") + execute_migration_file(connection) + + let assert Ok(first_dispatch) = + dispatch_user(connection, "renata", event_id: "no-context-first") + let assert Ok(second_dispatch) = + factos.new_dispatch( + connection:, + decision_context: factos.NoContext, + decider: decider(), + encode: encode_event, + decode: decode_event, + ) + |> factos_sqlight.dispatch(RegisterUser(username: "renata"), event_id: fn() { + "no-context-second" + }) + + let assert [first] = first_dispatch.events + let assert [second] = second_dispatch.events + assert first.event.payload == UserRegistered(username: "renata") + assert second.event.payload == UserRegistered(username: "renata") + assert first.position != second.position + assert all_recorded_events(connection) == [first, second] +} + fn decider() -> factos.Decider(Command, State, Event, DomainError) { factos.decider( initial: State(usernames: []), @@ -622,3 +840,168 @@ fn string_pair_decoder() -> decode.Decoder(#(String, String)) { use second <- decode.field(1, decode.string) decode.success(#(first, second)) } + +fn empty_query() -> factos.DecisionContext { + factos.Matching(items: []) +} + +fn username_conformance_query() -> factos.DecisionContext { + factos.Matching(items: [ + factos.item(types: [factos.event_type("UserRegistered")], tags: [ + factos.tag("username:renata"), + factos.tag("username:lucy"), + ]), + factos.item( + types: [ + factos.event_type("UnknownEventType"), + factos.event_type("UserRegistered"), + ], + tags: [factos.tag("username:lucy")], + ), + ]) +} + +fn assert_user_recorded( + recorded: factos.Recorded(Event), + position position: factos.SequencePosition, + username username: String, +) -> Nil { + assert recorded.position == position + assert recorded.event.descriptor.type_ == factos.event_type("UserRegistered") + assert recorded.event.descriptor.version == 1 + assert recorded.event.descriptor.tags == [factos.tag("username:" <> username)] + assert recorded.event.descriptor.metadata == factos.empty_metadata() + assert recorded.event.payload == UserRegistered(username:) +} + +fn insert_unknown_event(connection: sqlight.Connection) -> Nil { + let tag = "username:intruder" + let assert Ok(_) = + sqlight.query( + "insert into factos_events (id, type, version, tags, metadata, data) + values (?, ?, ?, ?, ?, ?) + returning position", + on: connection, + with: [ + sqlight.text("unknown-event"), + sqlight.text("UnknownEventType"), + sqlight.int(1), + sqlight.text("\n" <> tag <> "\n"), + sqlight.text(""), + sqlight.text("\"unknown\""), + ], + expecting: int_field_decoder(), + ) + Nil +} + +fn all_recorded_events( + connection: sqlight.Connection, +) -> List(factos.Recorded(Event)) { + let assert Ok(events) = + factos_sqlight.read_after( + connection, + factos.AllEvents, + factos.NoPosition, + 100, + decode_event, + ) + events +} + +fn dispatch_counter_context_many( + connection: sqlight.Connection, + decision_context: factos.DecisionContext, + remaining: Int, +) -> Result(factos.Dispatch(CounterEvent), factos_sqlight.Error(Nil, Nil)) { + let result = + factos.new_dispatch( + connection:, + decision_context:, + decider: counter_decider(), + encode: encode_counter_event, + decode: decode_counter_event, + ) + |> factos_sqlight.dispatch(Increment, event_id: fn() { + "counter-" <> int.to_string(remaining) + }) + case remaining, result { + 1, _ -> result + _, Ok(_) -> + dispatch_counter_context_many(connection, decision_context, remaining - 1) + _, Error(error) -> Error(error) + } +} + +fn counter_decider() -> factos.Decider( + CounterCommand, + CounterState, + CounterEvent, + Nil, +) { + factos.decider( + initial: CounterState(0), + decide: counter_decide, + evolve: counter_evolve, + ) +} + +fn counter_decide( + state: CounterState, + command: CounterCommand, +) -> Result(List(CounterEvent), Nil) { + let CounterState(total) = state + case command { + Increment -> Ok([Incremented(total + 1)]) + } +} + +fn counter_evolve(state: CounterState, event: CounterEvent) -> CounterState { + let CounterState(total) = state + case event { + Incremented(_) -> CounterState(total + 1) + } +} + +fn encode_counter_event(event: CounterEvent) -> factos.Event(json.Json) { + case event { + Incremented(value) -> + factos.new_event( + type_: factos.event_type("Incremented"), + version: 1, + data: json.int(value), + ) + |> factos.with_tags(tags: [factos.tag("counter:load")]) + } +} + +fn decode_counter_event( + stored: factos.Recorded(String), +) -> Result(CounterEvent, factos.Recorded(String)) { + case + factos.event_type_to_string(stored.event.descriptor.type_), + stored.event.descriptor.version + { + "Incremented", 1 -> + json.parse( + stored.event.payload, + using: decode.int |> decode.map(Incremented), + ) + |> result.replace_error(stored) + _, _ -> Error(stored) + } +} + +fn assert_counter_recorded( + recorded: factos.Recorded(CounterEvent), + position position: factos.SequencePosition, + value value: Int, + type_ type_: factos.EventType, +) -> Nil { + assert recorded.position == position + assert recorded.event.descriptor.type_ == type_ + assert recorded.event.descriptor.version == 1 + assert recorded.event.descriptor.tags == [factos.tag("counter:load")] + assert recorded.event.descriptor.metadata == factos.empty_metadata() + assert recorded.event.payload == Incremented(value) +}