diff --git a/backends/factos_sqlight/gleam.toml b/backends/factos_sqlight/gleam.toml index 61e7c79..820994c 100644 --- a/backends/factos_sqlight/gleam.toml +++ b/backends/factos_sqlight/gleam.toml @@ -26,6 +26,7 @@ gleam_stdlib = ">= 1.0.0 and < 2.0.0" sqlight = ">= 1.1.0 and < 2.0.0" gleam_erlang = ">= 1.3.0 and < 2.0.0" exception = ">= 2.1.1 and < 3.0.0" +gleam_json = ">= 3.1.0 and < 4.0.0" [dev_dependencies] gleeunit = ">= 1.0.0 and < 2.0.0" diff --git a/backends/factos_sqlight/manifest.toml b/backends/factos_sqlight/manifest.toml index 459add6..c6e3613 100644 --- a/backends/factos_sqlight/manifest.toml +++ b/backends/factos_sqlight/manifest.toml @@ -12,6 +12,7 @@ packages = [ { name = "factos", version = "2.0.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], source = "local", path = "../.." }, { name = "filepath", version = "1.1.2", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "filepath", source = "hex", outer_checksum = "B06A9AF0BF10E51401D64B98E4B627F1D2E48C154967DA7AF4D0914780A6D40A" }, { name = "gleam_erlang", version = "1.3.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_erlang", source = "hex", outer_checksum = "1124AD3AA21143E5AF0FC5CF3D9529F6DB8CA03E43A55711B60B6B7B3874375C" }, + { name = "gleam_json", version = "3.1.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_json", source = "hex", outer_checksum = "44FDAA8847BE8FC48CA7A1C089706BD54BADCC4C45B237A992EDDF9F2CDB2836" }, { name = "gleam_stdlib", version = "1.0.3", build_tools = ["gleam"], requirements = [], otp_app = "gleam_stdlib", source = "hex", outer_checksum = "1F543AFBA5D33DA493E6087F4E4C4F20D899411343512686C98A8ABB2963CF22" }, { name = "gleeunit", version = "1.11.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleeunit", source = "hex", outer_checksum = "EC31ABA74256AEA531EDF8169931D775BBB384FED0A8A1BDC4DD9354E3E21826" }, { name = "simplifile", version = "2.5.0", build_tools = ["gleam"], requirements = ["filepath", "gleam_stdlib"], otp_app = "simplifile", source = "hex", outer_checksum = "6C72DCCDF25C38A5931740B30E823969F33106831FD1637719B5EDBCA30027A4" }, @@ -22,6 +23,7 @@ packages = [ exception = { version = ">= 2.1.1 and < 3.0.0" } factos = { path = "../.." } gleam_erlang = { version = ">= 1.3.0 and < 2.0.0" } +gleam_json = { version = ">= 3.1.0 and < 4.0.0" } gleam_stdlib = { version = ">= 1.0.0 and < 2.0.0" } gleeunit = { version = ">= 1.0.0 and < 2.0.0" } simplifile = { version = ">= 2.5.0 and < 3.0.0" } diff --git a/backends/factos_sqlight/src/factos/factos_sqlight.gleam b/backends/factos_sqlight/src/factos/factos_sqlight.gleam index 42a8775..df9913c 100644 --- a/backends/factos_sqlight/src/factos/factos_sqlight.gleam +++ b/backends/factos_sqlight/src/factos/factos_sqlight.gleam @@ -10,13 +10,38 @@ import exception import factos +import gleam/dict import gleam/dynamic/decode import gleam/erlang/process +import gleam/json import gleam/list import gleam/result import gleam/string import sqlight +pub type Decoder(event) = + fn(factos.Recorded(String)) -> Result(event, factos.Recorded(String)) + +pub type Error(domain_error, subscription_error) = + factos.Error( + domain_error, + subscription_error, + sqlight.Error, + factos.Recorded(String), + ) + +type QuerySql { + QuerySql(sql: String, arguments: List(sqlight.Value)) +} + +type PreparedEvent(event) { + PreparedEvent( + id: String, + event: factos.Event(event), + encoded_payload: json.Json, + ) +} + /// Execute a shared dispatch builder against SQLite. /// /// Dispatch uses `BEGIN IMMEDIATE` to read the decision context, decide, append, @@ -34,22 +59,22 @@ pub fn dispatch( command, state, event, - BitArray, + json.Json, + String, domain_error, subscription_error, sqlight.Connection, + factos.Recorded(String), ), command: command, event_id event_id: fn() -> String, -) -> Result( - factos.Dispatch(event), - factos.Error(domain_error, subscription_error, sqlight.Error), -) { +) -> Result(factos.Dispatch(event), Error(domain_error, subscription_error)) { let factos.DispatchBuilder( connection:, decision_context:, decider:, - codec:, + encode:, + decode:, retry_attempts:, subscriptions:, ) = builder @@ -59,7 +84,8 @@ pub fn dispatch( connection, decision_context:, decider:, - codec:, + encode:, + decode:, command:, event_id:, retry_attempts:, @@ -79,17 +105,13 @@ pub fn dispatch( } } -type QuerySql { - QuerySql(sql: String, arguments: List(sqlight.Value)) -} - /// Create the fresh SQLite schema required by this backend. /// /// Applications with existing databases should copy the statements into their /// own immutable migration history rather than using this convenience at startup. pub fn migrate( connection: sqlight.Connection, -) -> Result(Nil, factos.Error(domain_error, subscription_error, sqlight.Error)) { +) -> Result(Nil, Error(domain_error, subscription_error)) { sqlight.exec(migration_sql, on: connection) |> result.map_error(factos.StoreError) } @@ -103,19 +125,19 @@ pub fn sqlight_error_to_string(error: sqlight.Error) -> String { @internal pub fn read( connection: sqlight.Connection, - decision_context decision_context: factos.DecisionContext, - decider decider: factos.Decider(command, state, event, domain_error), - codec codec: factos.EventCodec(event, BitArray), + decision_context: factos.DecisionContext, + decider: factos.Decider(command, state, event, domain_error), + decode: fn(factos.Recorded(String)) -> Result(event, factos.Recorded(String)), ) -> Result( factos.Context(event, state), - factos.Error(domain_error, subscription_error, sqlight.Error), + Error(domain_error, subscription_error), ) { let factos.Decider(initial:, evolve:, ..) = decider use events <- result.try(read_matching_events( connection, decision_context, - codec, + decode, )) let position = factos.highest_recorded_position(events) @@ -138,13 +160,13 @@ pub fn read( @internal pub fn read_after( connection: sqlight.Connection, - decision_context decision_context: factos.DecisionContext, - after after: factos.SequencePosition, - limit limit: Int, - codec codec: factos.EventCodec(event, BitArray), + decision_context: factos.DecisionContext, + after: factos.SequencePosition, + limit: Int, + decode: Decoder(event), ) -> Result( List(factos.Recorded(event)), - factos.Error(domain_error, subscription_error, sqlight.Error), + Error(domain_error, subscription_error), ) { case limit <= 0 { True -> Ok([]) @@ -157,9 +179,8 @@ pub fn read_after( order by position limit ?", on: connection, with: list.append(arguments, [ sqlight.int(limit), - ]), expecting: stored_event_decoder()) + ]), expecting: stored_row_decoder(decode)) |> result.map_error(factos.StoreError) - |> result.try(decode_stored_events(_, codec)) } } } @@ -168,88 +189,71 @@ fn dispatch_context( connection: sqlight.Connection, decision_context decision_context: factos.DecisionContext, decider decider: factos.Decider(command, state, event, domain_error), - codec codec: factos.EventCodec(event, BitArray), + encode encode: fn(event) -> factos.Event(json.Json), + decode decode: Decoder(event), command command: command, event_id event_id: fn() -> String, retry_attempts retry_attempts: Int, subscriptions subscriptions: List( factos.Subscription(event, subscription_error, sqlight.Connection), ), -) -> Result( - factos.Dispatch(event), - factos.Error(domain_error, subscription_error, sqlight.Error), -) { - run_immediate_transaction( +) -> Result(factos.Dispatch(event), Error(domain_error, subscription_error)) { + use transaction_connection <- run_serializable_transaction( connection, retry_attempts, - fn(transaction_connection) { - use context <- result.try(read( - transaction_connection, - decision_context, - decider, - codec, - )) - use pair <- result.try( - factos.decide_context(context, command, decider) - |> result.map_error(factos.DomainError), - ) - let #(context, events) = pair - use dispatch <- result.try(append_events( - transaction_connection, - events, - codec, - event_id, - context.append_condition, - )) - use _ <- result.try(run_strong_subscriptions( - transaction_connection, - subscriptions, - dispatch.events, - )) - Ok(dispatch) - }, ) + use context <- result.try(read( + transaction_connection, + decision_context, + decider, + decode, + )) + use pair <- result.try( + factos.decide_context(context, command, decider) + |> result.map_error(factos.DomainError), + ) + let #(context, events) = pair + use dispatch <- result.try(append_events( + transaction_connection, + events, + encode, + event_id, + context.append_condition, + )) + use _ <- result.try(run_strong_subscriptions( + transaction_connection, + subscriptions, + dispatch.events, + )) + Ok(dispatch) } -fn run_immediate_transaction( +fn run_serializable_transaction( connection: sqlight.Connection, retry_attempts: Int, work: fn(sqlight.Connection) -> - Result( - factos.Dispatch(event), - factos.Error(domain_error, subscription_error, sqlight.Error), - ), -) -> Result( - factos.Dispatch(event), - factos.Error(domain_error, subscription_error, sqlight.Error), -) { - run_immediate_transaction_attempt( + Result(factos.Dispatch(event), Error(domain_error, subscription_error)), +) -> Result(factos.Dispatch(event), Error(domain_error, subscription_error)) { + run_serializable_transaction_attempt( connection, work, attempts_remaining: retry_attempts, ) } -fn run_immediate_transaction_attempt( +fn run_serializable_transaction_attempt( connection: sqlight.Connection, work: fn(sqlight.Connection) -> - Result( - factos.Dispatch(event), - factos.Error(domain_error, subscription_error, sqlight.Error), - ), + Result(factos.Dispatch(event), Error(domain_error, subscription_error)), attempts_remaining attempts_remaining: Int, -) -> Result( - factos.Dispatch(event), - factos.Error(domain_error, subscription_error, sqlight.Error), -) { - let transaction_result = run_immediate_transaction_once(connection, work) - +) -> Result(factos.Dispatch(event), Error(domain_error, subscription_error)) { + let transaction_result = run_serializable_transaction_once(connection, work) case transaction_result { Ok(dispatch) -> Ok(dispatch) Error(error) -> case attempts_remaining > 1 && retryable_transaction_error(error) { True -> - run_immediate_transaction_attempt( + run_serializable_transaction_attempt( connection, work, attempts_remaining: attempts_remaining - 1, @@ -259,17 +263,11 @@ fn run_immediate_transaction_attempt( } } -fn run_immediate_transaction_once( +fn run_serializable_transaction_once( connection: sqlight.Connection, work: fn(sqlight.Connection) -> - Result( - factos.Dispatch(event), - factos.Error(domain_error, subscription_error, sqlight.Error), - ), -) -> Result( - factos.Dispatch(event), - factos.Error(domain_error, subscription_error, sqlight.Error), -) { + Result(factos.Dispatch(event), Error(domain_error, subscription_error)), +) -> Result(factos.Dispatch(event), Error(domain_error, subscription_error)) { use _ <- result.try( sqlight.exec("begin immediate", on: connection) |> result.map_error(factos.StoreError), @@ -281,13 +279,10 @@ fn run_immediate_transaction_once( fn finish_transaction( transaction_result: Result( factos.Dispatch(event), - factos.Error(domain_error, subscription_error, sqlight.Error), + Error(domain_error, subscription_error), ), connection: sqlight.Connection, -) -> Result( - factos.Dispatch(event), - factos.Error(domain_error, subscription_error, sqlight.Error), -) { +) -> Result(factos.Dispatch(event), Error(domain_error, subscription_error)) { case transaction_result { Ok(dispatch) -> case sqlight.exec("commit", on: connection) { @@ -305,7 +300,7 @@ fn finish_transaction( } fn retryable_transaction_error( - error: factos.Error(domain_error, subscription_error, sqlight.Error), + error: Error(domain_error, subscription_error), ) -> Bool { case error { factos.StoreError(sqlight.SqlightError(code:, ..)) -> @@ -334,21 +329,16 @@ fn run_strong_subscriptions( factos.Subscription(event, subscription_error, sqlight.Connection), ), events: List(factos.Recorded(event)), -) -> Result(Nil, factos.Error(domain_error, subscription_error, sqlight.Error)) { +) -> Result(Nil, Error(domain_error, subscription_error)) { case subscriptions { [] -> Ok(Nil) - [factos.Subscription(decision_context:, consistency:, handle:), ..remaining] -> + [factos.Subscription(consistency:, handle:), ..remaining] -> case consistency { factos.FireAndForget -> run_strong_subscriptions(connection, remaining, events) factos.StrongConsistency -> { use _ <- result.try( - run_strong_subscription_events( - connection, - decision_context, - handle, - events, - ) + run_strong_subscription_events(connection, handle, events) |> result.map_error(factos.SubscriptionError), ) run_strong_subscriptions(connection, remaining, events) @@ -359,32 +349,16 @@ fn run_strong_subscriptions( fn run_strong_subscription_events( connection: sqlight.Connection, - decision_context: factos.DecisionContext, handle: fn(sqlight.Connection, factos.Recorded(event)) -> Result(Nil, subscription_error), events: List(factos.Recorded(event)), ) -> Result(Nil, subscription_error) { case events { [] -> Ok(Nil) - [recorded, ..remaining] -> - case factos.matches_decision_context(recorded, decision_context) { - True -> { - use _ <- result.try(handle(connection, recorded)) - run_strong_subscription_events( - connection, - decision_context, - handle, - remaining, - ) - } - False -> - run_strong_subscription_events( - connection, - decision_context, - handle, - remaining, - ) - } + [recorded, ..remaining] -> { + use Nil <- result.try(handle(connection, recorded)) + run_strong_subscription_events(connection, handle, remaining) + } } } @@ -399,12 +373,6 @@ fn enqueue_fire_and_forget_subscriptions( case subscription.consistency { factos.StrongConsistency -> Nil factos.FireAndForget -> { - let events = - list.filter(events, factos.matches_decision_context( - _, - subscription.decision_context, - )) - case events { [] -> Nil events -> { @@ -438,92 +406,98 @@ fn run_fire_and_forget_events( fn append_events( connection: sqlight.Connection, events: List(event), - codec: factos.EventCodec(event, BitArray), + encode: fn(event) -> factos.Event(json.Json), event_id: fn() -> String, condition: factos.AppendCondition, -) -> Result( - factos.Dispatch(event), - factos.Error(domain_error, subscription_error, sqlight.Error), -) { +) -> Result(factos.Dispatch(event), Error(domain_error, subscription_error)) { case events { [] -> Ok(factos.Dispatch(position: factos.NoPosition, events: [])) - [_, ..] -> - case has_matching_events_after_condition(connection, condition) { - Error(error) -> Error(factos.StoreError(error)) - Ok(True) -> Error(factos.AppendConditionFailed(condition)) - Ok(False) -> - insert_events( - connection, - events, - codec, - event_id, - factos.NoPosition, - [], - ) - } + [_, ..] -> { + let prepared_events = prepare_events(events, encode, event_id) + insert_event_batch(connection, prepared_events, condition) + } } } -fn insert_events( - connection: sqlight.Connection, +fn prepare_events( events: List(event), - codec: factos.EventCodec(event, BitArray), + encode: fn(event) -> factos.Event(json.Json), event_id: fn() -> String, - position: factos.SequencePosition, - recorded_events: List(factos.Recorded(event)), -) -> Result( - factos.Dispatch(event), - factos.Error(domain_error, subscription_error, sqlight.Error), -) { +) -> List(PreparedEvent(event)) { case events { - [] -> Ok(factos.Dispatch(position:, events: list.reverse(recorded_events))) - [event, ..remaining] -> { - let factos.Event( - payload:, - descriptor: factos.EventDescriptor(type_:, version:, tags:, metadata:), - ) = codec.encode(event) - let id = event_id() - use positions <- result.try( - sqlight.query( - "insert into factos_events (id, type, version, tags, metadata, data) - values (?, ?, ?, ?, ?, ?) - returning position", - on: connection, - with: [ - sqlight.text(id), - sqlight.text(factos.event_type_name(type_)), - sqlight.int(version), - sqlight.text(tags_to_text(tags)), - sqlight.text(metadata_to_text(metadata)), - sqlight.blob(payload), - ], - expecting: int_field_decoder(), + [] -> [] + [payload, ..remaining] -> { + let encoded = encode(payload) + let prepared = + PreparedEvent( + id: event_id(), + event: factos.Event(payload:, descriptor: encoded.descriptor), + encoded_payload: encoded.payload, ) - |> result.map_error(factos.StoreError), - ) + [prepared, ..prepare_events(remaining, encode, event_id)] + } + } +} + +fn insert_event_batch( + connection: sqlight.Connection, + events: List(PreparedEvent(event)), + condition condition: factos.AppendCondition, +) -> Result(factos.Dispatch(event), Error(domain_error, subscription_error)) { + use conflicting <- result.try( + has_matching_events_after_condition(connection, condition) + |> result.map_error(factos.StoreError), + ) - case positions { - [] -> Error(factos.DecodeError(factos.InvalidData)) - [returned_position, ..] -> { - let position = factos.SequencePosition(returned_position) - let recorded = - factos.Recorded( - id:, - position:, - event:, - descriptor: factos.EventDescriptor( - type_:, - version:, - tags:, - metadata:, - ), + case conflicting { + True -> Error(factos.AppendConditionFailed(condition)) + False -> { + use recorded_events <- result.try( + list.try_map(events, fn(event) { + let factos.EventDescriptor(type_:, version:, tags:, metadata:) = + event.event.descriptor + use positions <- result.try( + sqlight.query( + "insert into factos_events (id, type, version, tags, metadata, data) + values (?, ?, ?, ?, ?, ?) + returning position", + on: connection, + with: [ + sqlight.text(event.id), + sqlight.text(factos.event_type_to_string(type_)), + sqlight.int(version), + sqlight.text(tags_to_text(tags)), + sqlight.text(metadata_to_text(metadata)), + sqlight.text(json.to_string(event.encoded_payload)), + ], + expecting: int_field_decoder(), ) - insert_events(connection, remaining, codec, event_id, position, [ - recorded, - ..recorded_events - ]) - } - } + |> result.map_error(factos.StoreError), + ) + + case positions { + [] -> + Error( + factos.StoreError(sqlight.SqlightError( + code: sqlight.GenericError, + message: "insert into factos_events returned no position", + offset: -1, + )), + ) + [position, ..] -> + Ok(factos.Recorded( + id: event.id, + position: factos.SequencePosition(position), + event: event.event, + )) + } + }), + ) + + Ok(factos.Dispatch( + position: factos.highest_recorded_position(recorded_events), + events: recorded_events, + )) } } } @@ -539,81 +513,72 @@ fn has_matching_events_after_condition( fn read_matching_events( connection: sqlight.Connection, decision_context: factos.DecisionContext, - codec: factos.EventCodec(event, BitArray), + decode: Decoder(event), ) -> Result( List(factos.Recorded(event)), - factos.Error(domain_error, subscription_error, sqlight.Error), + Error(domain_error, subscription_error), ) { let QuerySql(where_sql, arguments) = query_to_sql(decision_context) sqlight.query("select position, id, type, version, tags, metadata, data from factos_events " <> where_sql <> " - order by position", on: connection, with: arguments, expecting: stored_event_decoder()) + order by position", on: connection, with: arguments, expecting: stored_row_decoder( + decode, + )) |> result.map_error(factos.StoreError) - |> result.try(decode_stored_events(_, codec)) -} - -fn decode_stored_events( - rows: List(factos.Recorded(BitArray)), - codec: factos.EventCodec(event, BitArray), -) -> Result( - List(factos.Recorded(event)), - factos.Error(domain_error, subscription_error, sqlight.Error), -) { - decode_stored_events_loop(rows, codec, []) } -fn decode_stored_events_loop( - rows: List(factos.Recorded(BitArray)), - codec: factos.EventCodec(event, BitArray), - decoded: List(factos.Recorded(event)), -) -> Result( - List(factos.Recorded(event)), - factos.Error(domain_error, subscription_error, sqlight.Error), -) { - case rows { - [] -> Ok(list.reverse(decoded)) - [row, ..remaining] -> { - use recorded <- result.try(decode_stored_event(row, codec)) - decode_stored_events_loop(remaining, codec, [recorded, ..decoded]) - } +fn stored_row_decoder( + decode: Decoder(event), +) -> decode.Decoder(factos.Recorded(event)) { + use position <- decode.field(0, decode.int) + use id <- decode.field(1, decode.string) + use type_ <- decode.field(2, decode.string |> decode.map(factos.event_type)) + use version <- decode.field(3, decode.int) + use tags <- decode.field(4, tags_column_decoder()) + use metadata <- decode.field(5, metadata_column_decoder()) + use payload <- decode.field(6, decode.string) + let string_event = + factos.Recorded( + id:, + position: factos.SequencePosition(position), + event: factos.Event( + descriptor: factos.EventDescriptor(type_:, version:, tags:, metadata:), + payload:, + ), + ) + case decode(string_event) { + Ok(payload) -> + decode.success( + factos.Recorded( + ..string_event, + event: factos.Event(..string_event.event, payload:), + ), + ) + Error(error) -> decode.failure(coerce(error), expected: "Decodable payload") } } -fn decode_stored_event( - stored: factos.Recorded(BitArray), - codec: factos.EventCodec(event, BitArray), -) -> Result( - factos.Recorded(event), - factos.Error(domain_error, subscription_error, sqlight.Error), -) { - let factos.Recorded(position:, id:, descriptor:, ..) = stored - let factos.EventCodec(decode:, ..) = codec - use event <- result.try( - decode(stored) |> result.map_error(factos.DecodeError), - ) - - Ok(factos.Recorded(id:, position:, event:, descriptor:)) -} +@external(erlang, "gleam_stdlib", "identity") +fn coerce(a: a) -> b -fn stored_event_decoder() -> decode.Decoder(factos.Recorded(BitArray)) { - use position <- decode.field(0, decode.int) - use id <- decode.field(1, decode.string) - use type_name <- decode.field(2, decode.string) - use version <- decode.field(3, decode.int) - use tags <- decode.field(4, decode.string) - use metadata <- decode.field(5, decode.string) - use event <- decode.field(6, decode.bit_array) - decode.success(factos.Recorded( - position: factos.SequencePosition(position), - id:, - descriptor: factos.EventDescriptor( - type_: factos.event_type(type_name), - version:, - tags: tags_from_text(tags), - metadata: metadata_from_text(metadata), - ), - event:, - )) +fn metadata_column_decoder() -> decode.Decoder(factos.Metadata) { + use value <- decode.then(decode.string) + case json.parse(value, decode.dict(decode.string, decode.string)) { + Ok(entries) -> + entries + |> dict.to_list + |> factos.metadata + |> decode.success + Error(_) -> + case string.starts_with(string.trim(value), "{") { + True -> + decode.failure( + factos.empty_metadata(), + expected: "a JSON object containing string values", + ) + False -> decode.success(metadata_from_text(value)) + } + } } fn has_matching_events_after( @@ -722,7 +687,7 @@ fn types_to_sql(types: List(factos.EventType)) -> QuerySql { QuerySql( sql: "type in (" <> placeholders(list.length(types)) <> ")", arguments: list.map(types, fn(type_) { - sqlight.text(factos.event_type_name(type_)) + sqlight.text(factos.event_type_to_string(type_)) }), ) } @@ -779,6 +744,24 @@ fn tags_from_text(tags: String) -> List(factos.Tag) { } } +fn tags_column_decoder() -> decode.Decoder(List(factos.Tag)) { + use tags <- decode.then(decode.string) + + case + json.parse( + tags, + using: decode.string |> decode.map(factos.tag) |> decode.list, + ) + { + Ok(tags) -> decode.success(tags) + Error(_) -> + case string.starts_with(string.trim(tags), "[") { + True -> decode.failure([], "expected a json array of string tags") + False -> decode.success(tags_from_text(tags)) + } + } +} + fn metadata_to_text(metadata: factos.Metadata) -> String { metadata |> factos.metadata_entries @@ -802,7 +785,6 @@ fn metadata_from_text(metadata: String) -> factos.Metadata { } } -/// Render a backend error with application formatters for its generic errors. const migration_sql = " create table if not exists factos_events ( position integer primary key autoincrement, diff --git a/backends/factos_sqlight/test/factos_sqlight_test.gleam b/backends/factos_sqlight/test/factos_sqlight_test.gleam index 031e863..5b83b5e 100644 --- a/backends/factos_sqlight/test/factos_sqlight_test.gleam +++ b/backends/factos_sqlight/test/factos_sqlight_test.gleam @@ -1,9 +1,9 @@ import factos import factos/factos_sqlight -import gleam/bit_array import gleam/dynamic/decode import gleam/erlang/application import gleam/erlang/process +import gleam/json import gleam/list import gleam/result import gleam/string @@ -61,8 +61,8 @@ pub fn dispatch_uses_decision_context_and_application_event_ids_test() -> Nil { let assert [renata] = renata_dispatch.events assert renata_dispatch.position == renata.position assert renata.id == "event-renata" - assert renata.event == UserRegistered(username: "renata") - assert renata.descriptor.tags == [factos.tag("username:renata")] + assert renata.event.payload == UserRegistered(username: "renata") + assert renata.event.descriptor.tags == [factos.tag("username:renata")] let duplicate = dispatch_user(connection, "renata", event_id: "unused") assert duplicate @@ -77,9 +77,9 @@ pub fn dispatch_uses_decision_context_and_application_event_ids_test() -> Nil { let assert Ok(context) = factos_sqlight.read( connection, - decision_context: username_context("renata"), - decider: decider(), - codec: codec(), + username_context("renata"), + decider(), + decode_event, ) assert context.state == State(usernames: ["renata"]) assert context.events == [renata] @@ -118,30 +118,30 @@ pub fn read_after_filters_orders_and_rejects_unknown_events_test() -> Nil { let assert Ok([renata]) = factos_sqlight.read_after( connection, - decision_context: username_context("renata"), - after: factos.NoPosition, - limit: 10, - codec: codec(), + username_context("renata"), + factos.NoPosition, + 10, + decode_event, ) assert renata == first let assert Ok([maria]) = factos_sqlight.read_after( connection, - decision_context: factos.AllEvents, - after: first.position, - limit: 1, - codec: codec(), + factos.AllEvents, + first.position, + 1, + decode_event, ) assert maria == second let assert Ok([]) = factos_sqlight.read_after( connection, - decision_context: factos.AllEvents, - after: factos.NoPosition, - limit: 0, - codec: codec(), + factos.AllEvents, + factos.NoPosition, + 0, + decode_event, ) let assert Ok(_) = @@ -149,19 +149,20 @@ pub fn read_after_filters_orders_and_rejects_unknown_events_test() -> Nil { connection:, decision_context: factos.NoContext, decider: accepting_decider(), - codec: hidden_codec(), + encode: hidden_encode_event, + decode: decode_event, ) |> factos_sqlight.dispatch(RegisterUser(username: "hidden"), event_id: fn() { "hidden" }) - let assert Error(factos.DecodeError(factos.UnknownEvent)) = + let assert Error(factos.StoreError(_)) = factos_sqlight.read_after( connection, - decision_context: factos.AllEvents, - after: second.position, - limit: 10, - codec: codec(), + factos.AllEvents, + second.position, + 10, + decode_event, ) Nil } @@ -173,7 +174,6 @@ pub fn empty_dispatch_has_no_position_and_skips_subscriptions_test() -> Nil { let event_id_invocations = process.new_subject() let subscriptions = [ factos.new_subscription( - decision_context: factos.AllEvents, consistency: factos.StrongConsistency, handle: fn(_connection, _recorded) { process.send(invocations, "strong") @@ -187,7 +187,8 @@ pub fn empty_dispatch_has_no_position_and_skips_subscriptions_test() -> Nil { connection:, decision_context: factos.NoContext, decider: empty_decider(), - codec: codec(), + encode: encode_event, + decode: decode_event, ) |> factos.with_subscriptions(subscriptions:) |> factos_sqlight.dispatch(DoNothing, event_id: fn() { @@ -207,7 +208,6 @@ pub fn strong_subscription_commits_with_dispatch_test() -> Nil { create_projection_table(connection) let subscription = factos.new_subscription( - decision_context: factos.AllEvents, consistency: factos.StrongConsistency, handle: insert_projection, ) @@ -232,13 +232,11 @@ pub fn strong_subscription_failure_rolls_back_everything_test() -> Nil { let fire_deliveries = process.new_subject() let insert = factos.new_subscription( - decision_context: factos.AllEvents, consistency: factos.StrongConsistency, handle: insert_projection, ) let fail_after_observing_insert = factos.new_subscription( - decision_context: factos.AllEvents, consistency: factos.StrongConsistency, handle: fn(transaction_connection, recorded) { case projection_rows(transaction_connection) { @@ -249,7 +247,6 @@ pub fn strong_subscription_failure_rolls_back_everything_test() -> Nil { ) let fire = factos.new_subscription( - decision_context: factos.AllEvents, consistency: factos.FireAndForget, handle: fn(_connection, _recorded) { process.send(fire_deliveries, Nil) @@ -266,60 +263,18 @@ pub fn strong_subscription_failure_rolls_back_everything_test() -> Nil { ) let assert Error(error) = result assert error == factos.SubscriptionError(error: "expected strong failure") - assert factos.error_to_string( - error, - fn(_) { "domain" }, - fn(subscription_error) { subscription_error }, - fn(_) { "sqlight" }, - ) - == "subscription error: expected strong failure" assert count_events(connection) == 0 assert projection_rows(connection) == [] let assert Error(Nil) = process.receive(fire_deliveries, within: 100) Nil } -pub fn subscriptions_filter_nonmatching_dispatch_events_test() -> Nil { - use connection <- sqlight.with_connection(":memory:") - execute_migration_file(connection) - let invocations = process.new_subject() - let subscriptions = [ - factos.new_subscription( - decision_context: username_context("maria"), - consistency: factos.StrongConsistency, - handle: fn(_connection, _recorded) { - process.send(invocations, "strong") - Ok(Nil) - }, - ), - factos.new_subscription( - decision_context: username_context("maria"), - consistency: factos.FireAndForget, - handle: fn(_connection, _recorded) { - process.send(invocations, "fire") - Ok(Nil) - }, - ), - ] - - let assert Ok(_) = - dispatch_user_with_subscriptions( - connection, - "renata", - event_id: "filtered", - subscriptions:, - ) - let assert Error(Nil) = process.receive(invocations, within: 150) - assert count_events(connection) == 1 -} - pub fn fire_and_forget_runs_after_commit_without_blocking_test() -> Nil { use connection <- sqlight.with_connection(":memory:") execute_migration_file(connection) let deliveries = process.new_subject() let subscription = factos.new_subscription( - decision_context: factos.AllEvents, consistency: factos.FireAndForget, handle: fn(callback_connection, recorded) { let release = process.new_subject() @@ -366,10 +321,9 @@ pub fn fire_and_forget_continues_after_callback_error_in_append_order_test() -> process.send(ids, "pair-second") let subscription = factos.new_subscription( - decision_context: factos.AllEvents, consistency: factos.FireAndForget, handle: fn(_connection, recorded) { - let UserRegistered(username:) = recorded.event + let UserRegistered(username:) = recorded.event.payload process.send(deliveries, username) Error("ignored") }, @@ -380,7 +334,8 @@ pub fn fire_and_forget_continues_after_callback_error_in_append_order_test() -> connection:, decision_context: factos.NoContext, decider: accepting_decider(), - codec: codec(), + encode: encode_event, + decode: decode_event, ) |> factos.with_subscriptions(subscriptions: [subscription]) |> factos_sqlight.dispatch( @@ -447,47 +402,39 @@ fn empty_decider() -> factos.Decider(Command, State, Event, DomainError) { ) } -fn codec() -> factos.EventCodec(Event, BitArray) { - factos.codec(encode: encode_event, decode: decode_event) -} - -fn hidden_codec() -> factos.EventCodec(Event, BitArray) { - factos.codec( - encode: fn(event) { - let UserRegistered(username:) = event - factos.new_event( - type_: factos.event_type("HiddenEvent"), - version: 1, - data: bit_array.from_string(username), - ) - }, - decode: decode_event, +fn encode_event(event: Event) -> factos.Event(json.Json) { + let UserRegistered(username:) = event + factos.new_event( + type_: factos.event_type("UserRegistered"), + version: 1, + data: json.string(username), ) + |> factos.with_tags(tags: [factos.tag("username:" <> username)]) } -fn encode_event(event: Event) -> factos.Event(BitArray) { +fn hidden_encode_event(event: Event) -> factos.Event(json.Json) { let UserRegistered(username:) = event factos.new_event( - type_: factos.event_type("UserRegistered"), + type_: factos.event_type("HiddenEvent"), version: 1, - data: bit_array.from_string(username), + data: json.string(username), ) - |> factos.with_tags(tags: [factos.tag("username:" <> username)]) } fn decode_event( - stored: factos.Recorded(BitArray), -) -> Result(Event, factos.DecodeError) { - let descriptor = stored.descriptor - case factos.event_type_name(descriptor.type_), descriptor.version { - "UserRegistered", 1 -> { - use username <- result.try( - bit_array.to_string(stored.event) - |> result.replace_error(factos.InvalidData), + stored: factos.Recorded(String), +) -> Result(Event, factos.Recorded(String)) { + case + factos.event_type_to_string(stored.event.descriptor.type_), + stored.event.descriptor.version + { + "UserRegistered", 1 -> + json.parse( + stored.event.payload, + using: decode.string |> decode.map(UserRegistered), ) - Ok(UserRegistered(username:)) - } - _, _ -> Error(factos.UnknownEvent) + |> result.replace_error(stored) + _, _ -> Error(stored) } } @@ -503,15 +450,13 @@ fn dispatch_user( connection: sqlight.Connection, username: String, event_id event_id: String, -) -> Result( - factos.Dispatch(Event), - factos.Error(DomainError, Nil, sqlight.Error), -) { +) -> Result(factos.Dispatch(Event), factos_sqlight.Error(DomainError, Nil)) { factos.new_dispatch( connection:, decision_context: username_context(username), decider: decider(), - codec: codec(), + encode: encode_event, + decode: decode_event, ) |> factos_sqlight.dispatch(RegisterUser(username:), event_id: fn() { event_id @@ -525,15 +470,13 @@ fn dispatch_user_with_subscriptions( subscriptions subscriptions: List( factos.Subscription(Event, String, sqlight.Connection), ), -) -> Result( - factos.Dispatch(Event), - factos.Error(DomainError, String, sqlight.Error), -) { +) -> Result(factos.Dispatch(Event), factos_sqlight.Error(DomainError, String)) { factos.new_dispatch( connection:, decision_context: username_context(username), decider: decider(), - codec: codec(), + encode: encode_event, + decode: decode_event, ) |> factos.with_subscriptions(subscriptions:) |> factos_sqlight.dispatch(RegisterUser(username:), event_id: fn() { @@ -545,7 +488,7 @@ fn insert_projection( connection: sqlight.Connection, recorded: factos.Recorded(Event), ) -> Result(Nil, String) { - let UserRegistered(username:) = recorded.event + let UserRegistered(username:) = recorded.event.payload sqlight.query( "insert into test_projection (event_id, username) values (?, ?)