-module(eventsourcing). -compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]). -define(FILEPATH, "src/eventsourcing.gleam"). -export([execute/3, execute_with_metadata/4, load_aggregate/2, load_events/2, load_events_from/3, system_stats/1, aggregate_stats/2, latest_snapshot/2, supervised/7, timeout/1, frequency/1]). -export_type([timeout_/0, frequency/0, aggregate/4, snapshot/1, snapshot_config/0, event_envelop/1, event_sourcing_error/1, system_stats/0, aggregate_stats/0, query_actor/1, query_message/1, aggregate_message/4, manager_message/4, manager_state/6, event_sourcing/6, event_store/6, aggregate_actor_state/6]). -if(?OTP_RELEASE >= 27). -define(MODULEDOC(Str), -moduledoc(Str)). -define(DOC(Str), -doc(Str)). -else. -define(MODULEDOC(Str), -compile([])). -define(DOC(Str), -compile([])). -endif. -opaque timeout_() :: {timeout, integer()}. -opaque frequency() :: {frequency, integer()}. -type aggregate(JRT, JRU, JRV, JRW) :: {aggregate, binary(), JRT, integer()} | {gleam_phantom, JRU, JRV, JRW}. -type snapshot(JRX) :: {snapshot, binary(), JRX, integer(), gleam@time@timestamp:timestamp()}. -type snapshot_config() :: {snapshot_config, frequency()}. -type event_envelop(JRY) :: {memory_store_event_envelop, binary(), integer(), JRY, list({binary(), binary()})} | {serialized_event_envelop, binary(), integer(), JRY, list({binary(), binary()}), binary(), binary(), binary()}. -type event_sourcing_error(JRZ) :: {domain_error, JRZ} | {event_store_error, binary()} | non_positive_argument | entity_not_found | transaction_failed | transaction_rolled_back | {actor_timeout, binary(), integer()}. -type system_stats() :: {system_stats, integer(), integer()}. -type aggregate_stats() :: {aggregate_stats, binary(), integer(), integer(), boolean()}. -type query_actor(JSA) :: {query_actor, gleam@erlang@process:subject(query_message(JSA))}. -type query_message(JSB) :: {process_events, binary(), list(event_envelop(JSB))}. -type aggregate_message(JSC, JSD, JSE, JSF) :: {execute_command, binary(), JSD, list({binary(), binary()})} | {load_aggregate, binary(), gleam@erlang@process:subject({ok, aggregate(JSC, JSD, JSE, JSF)} | {error, event_sourcing_error(JSF)})} | {load_all_events, binary(), gleam@erlang@process:subject({ok, list(event_envelop(JSE))} | {error, event_sourcing_error(JSF)})} | {load_events, binary(), integer(), gleam@erlang@process:subject({ok, list(event_envelop(JSE))} | {error, event_sourcing_error(JSF)})} | {load_latest_snapshot, binary(), gleam@erlang@process:subject({ok, gleam@option:option(snapshot(JSC))} | {error, event_sourcing_error(JSF)})} | {get_system_stats, gleam@erlang@process:subject(system_stats())} | {get_aggregate_stats, binary(), gleam@erlang@process:subject({ok, aggregate_stats()} | {error, event_sourcing_error(JSF)})}. -type manager_message(JSG, JSH, JSI, JSJ) :: {query_actor_started, query_actor(JSI)} | {get_event_sourcing_actor, gleam@erlang@process:subject(gleam@erlang@process:subject(aggregate_message(JSG, JSH, JSI, JSJ)))}. -type manager_state(JSK, JSL, JSM, JSN, JSO, JSP) :: {manager_state, gleam@option:option(gleam@erlang@process:subject(aggregate_message(JSL, JSM, JSN, JSO))), integer(), integer(), event_store(JSK, JSL, JSM, JSN, JSO, JSP), fun((JSL, JSM) -> {ok, list(JSN)} | {error, JSO}), fun((JSL, JSN) -> JSL), JSL}. -opaque event_sourcing(JSQ, JSR, JSS, JST, JSU, JSV) :: {event_sourcing, event_store(JSQ, JSR, JSS, JST, JSU, JSV), list(gleam@erlang@process:name(query_message(JST))), fun((JSR, JSS) -> {ok, list(JST)} | {error, JSU}), fun((JSR, JST) -> JSR), JSR, gleam@option:option(snapshot_config()), gleam@time@timestamp:timestamp(), integer()}. -type event_store(JSW, JSX, JSY, JSZ, JTA, JTB) :: {event_store, fun((fun((JTB) -> {ok, nil} | {error, event_sourcing_error(JTA)})) -> {ok, nil} | {error, event_sourcing_error(JTA)}), fun((fun((JTB) -> {ok, aggregate(JSX, JSY, JSZ, JTA)} | {error, event_sourcing_error(JTA)})) -> {ok, aggregate(JSX, JSY, JSZ, JTA)} | {error, event_sourcing_error(JTA)}), fun((fun((JTB) -> {ok, list(event_envelop(JSZ))} | {error, event_sourcing_error(JTA)})) -> {ok, list(event_envelop(JSZ))} | {error, event_sourcing_error(JTA)}), fun((fun((JTB) -> {ok, gleam@option:option(snapshot(JSX))} | {error, event_sourcing_error(JTA)})) -> {ok, gleam@option:option(snapshot(JSX))} | {error, event_sourcing_error(JTA)}), fun((JTB, aggregate(JSX, JSY, JSZ, JTA), list(JSZ), list({binary(), binary()})) -> {ok, {list(event_envelop(JSZ)), integer()}} | {error, event_sourcing_error(JTA)}), fun((JSW, JTB, binary(), integer()) -> {ok, list(event_envelop(JSZ))} | {error, event_sourcing_error(JTA)}), fun((JTB, binary()) -> {ok, gleam@option:option(snapshot(JSX))} | {error, event_sourcing_error(JTA)}), fun((JTB, snapshot(JSX)) -> {ok, nil} | {error, event_sourcing_error(JTA)}), JSW}. -type aggregate_actor_state(JTC, JTD, JTE, JTF, JTG, JTH) :: {aggregate_actor_state, aggregate(JTD, JTE, JTF, JTG), event_sourcing(JTC, JTD, JTE, JTF, JTG, JTH)}. -file("src/eventsourcing.gleam", 364). ?DOC( " Executes a command against an aggregate in the event sourcing system.\n" " The command will be validated, events generated if successful, and the events\n" " will be persisted and sent to all registered query actors. Commands that violate\n" " business rules will be rejected without affecting system stability.\n" "\n" " ## Example\n" " ```gleam\n" " eventsourcing.execute(actor, \"bank-account-123\", OpenAccount(\"123\"))\n" " // Command is processed asynchronously via message passing\n" " ```\n" ). -spec execute( gleam@erlang@process:subject(aggregate_message(any(), JVA, any(), any())), binary(), JVA ) -> nil. execute(Eventsourcing_actor, Aggregate_id, Command) -> gleam@erlang@process:send( Eventsourcing_actor, {execute_command, Aggregate_id, Command, []} ). -file("src/eventsourcing.gleam", 383). ?DOC( " Executes a command against an aggregate with additional metadata.\n" " The metadata will be stored with the generated events and can be used for\n" " tracking, auditing, or enriching events with contextual information.\n" "\n" " ## Example\n" " ```gleam\n" " let metadata = [(\"user_id\", \"alice\"), (\"session_id\", \"abc123\")]\n" " eventsourcing.execute_with_metadata(actor, \"bank-123\", DepositMoney(100.0), metadata)\n" " ```\n" ). -spec execute_with_metadata( gleam@erlang@process:subject(aggregate_message(any(), JVJ, any(), any())), binary(), JVJ, list({binary(), binary()}) ) -> nil. execute_with_metadata(Eventsourcing_actor, Aggregate_id, Command, Metadata) -> gleam@erlang@process:send( Eventsourcing_actor, {execute_command, Aggregate_id, Command, Metadata} ). -file("src/eventsourcing.gleam", 637). -spec describe_error(event_sourcing_error(any())) -> binary(). describe_error(Error) -> case Error of {domain_error, Domainerror} -> <<"Domain error: "/utf8, (gleam@string:inspect(Domainerror))/binary>>; {event_store_error, Msg} -> <<"Event store error: "/utf8, Msg/binary>>; non_positive_argument -> <<"Non-positive argument"/utf8>>; entity_not_found -> <<"Entity not found"/utf8>>; transaction_failed -> <<"Transaction failed"/utf8>>; transaction_rolled_back -> <<"Transaction rolled back"/utf8>>; {actor_timeout, Operation, Timeout_ms} -> <<<<<<<<"Actor timeout: "/utf8, Operation/binary>>/binary, " failed after "/utf8>>/binary, (erlang:integer_to_binary(Timeout_ms))/binary>>/binary, "ms"/utf8>> end. -file("src/eventsourcing.gleam", 654). -spec load_aggregate_or_create_new( event_sourcing(any(), JYC, JYD, JYE, JYF, JYG), JYG, binary() ) -> {ok, aggregate(JYC, JYD, JYE, JYF)} | {error, event_sourcing_error(JYF)}. load_aggregate_or_create_new(Eventsourcing, Tx, Aggregate_id) -> begin Result = case erlang:element(7, Eventsourcing) of none -> {ok, none}; {some, _} -> (erlang:element(8, erlang:element(2, Eventsourcing)))( Tx, Aggregate_id ) end, case Result of {ok, X} -> {Starting_state, Starting_sequence} = case X of none -> {erlang:element(6, Eventsourcing), 0}; {some, Snapshot} -> {erlang:element(3, Snapshot), erlang:element(4, Snapshot)} end, begin Result@1 = (erlang:element( 7, erlang:element(2, Eventsourcing) ))( erlang:element(10, erlang:element(2, Eventsourcing)), Tx, Aggregate_id, Starting_sequence ), case Result@1 of {ok, X@1} -> {ok, begin {Instance, Sequence@1} = begin _pipe = X@1, gleam@list:fold( _pipe, {Starting_state, Starting_sequence}, fun( Aggregate_and_sequence, Event_envelop ) -> {Aggregate, Sequence} = Aggregate_and_sequence, {(erlang:element( 5, Eventsourcing ))( Aggregate, erlang:element( 4, Event_envelop ) ), Sequence + 1} end ) end, {aggregate, Aggregate_id, Instance, Sequence@1} end}; {error, E} -> {error, E} end end; {error, E@1} -> {error, E@1} end end. -file("src/eventsourcing.gleam", 430). -spec on_message( event_sourcing(JWX, JWY, JWZ, JXA, JXB, JXC), aggregate_message(JWY, JWZ, JXA, JXB) ) -> gleam@otp@actor:next(event_sourcing(JWX, JWY, JWZ, JXA, JXB, JXC), aggregate_message(JWY, JWZ, JXA, JXB)). on_message(State, Message) -> case Message of {load_latest_snapshot, Aggregate_id, Reply_to} -> Result = begin (erlang:element(5, erlang:element(2, State)))( fun(Tx) -> case erlang:element(7, State) of none -> {ok, none}; {some, _} -> (erlang:element(8, erlang:element(2, State)))( Tx, Aggregate_id ) end end ) end, gleam@erlang@process:send(Reply_to, Result), gleam@otp@actor:continue(State); {load_events, Aggregate_id@1, Start_from, Reply_to@1} -> Result@1 = begin (erlang:element(4, erlang:element(2, State)))( fun(Tx@1) -> (erlang:element(7, erlang:element(2, State)))( erlang:element(10, erlang:element(2, State)), Tx@1, Aggregate_id@1, Start_from ) end ) end, gleam@erlang@process:send(Reply_to@1, Result@1), gleam@otp@actor:continue(State); {load_all_events, Aggregate_id@2, Reply_to@2} -> Result@2 = begin (erlang:element(4, erlang:element(2, State)))( fun(Tx@2) -> (erlang:element(7, erlang:element(2, State)))( erlang:element(10, erlang:element(2, State)), Tx@2, Aggregate_id@2, 0 ) end ) end, gleam@erlang@process:send(Reply_to@2, Result@2), gleam@otp@actor:continue(State); {load_aggregate, Aggregate_id@3, Reply_to@3} -> Result@3 = begin (erlang:element(3, erlang:element(2, State)))( fun(Tx@3) -> _pipe = load_aggregate_or_create_new( State, Tx@3, Aggregate_id@3 ), gleam@result:'try'( _pipe, fun(Aggregate) -> case erlang:element(3, Aggregate) =:= erlang:element( 6, State ) of true -> {error, entity_not_found}; false -> {ok, Aggregate} end end ) end ) end, gleam@erlang@process:send(Reply_to@3, Result@3), gleam@otp@actor:continue(State); {execute_command, Aggregate_id@4, Command, Metadata} -> Result@4 = begin (erlang:element(2, erlang:element(2, State)))( fun(Tx@4) -> gleam@result:'try'( load_aggregate_or_create_new( State, Tx@4, Aggregate_id@4 ), fun(Aggregate@1) -> gleam@result:'try'( begin _pipe@1 = (erlang:element(4, State))( erlang:element(3, Aggregate@1), Command ), gleam@result:map_error( _pipe@1, fun(Error) -> {domain_error, Error} end ) end, fun(Events) -> Aggregate@2 = {aggregate, erlang:element(2, Aggregate@1), begin _pipe@2 = Events, gleam@list:fold( _pipe@2, erlang:element( 3, Aggregate@1 ), fun(Entity, Event) -> (erlang:element( 5, State ))(Entity, Event) end ) end, erlang:element(4, Aggregate@1)}, gleam@result:'try'( (erlang:element( 6, erlang:element(2, State) ))( Tx@4, Aggregate@2, Events, Metadata ), fun(_use0) -> {Commited_events, Sequence} = _use0, gleam@result:'try'( case erlang:element( 7, State ) of {some, Config} when (Sequence rem erlang:element( 2, erlang:element( 2, Config ) )) =:= 0 -> Snapshot = {snapshot, erlang:element( 2, Aggregate@2 ), erlang:element( 3, Aggregate@2 ), Sequence, gleam@time@timestamp:system_time( )}, (erlang:element( 9, erlang:element( 2, State ) ))(Tx@4, Snapshot); _ -> {ok, nil} end, fun(_) -> _pipe@3 = erlang:element( 3, State ), gleam@list:each( _pipe@3, fun(Query) -> gleam@erlang@process:send( gleam@erlang@process:named_subject( Query ), {process_events, Aggregate_id@4, Commited_events} ) end ), {ok, nil} end ) end ) end ) end ) end ) end, case Result@4 of {ok, _} -> Updated_state = {event_sourcing, erlang:element(2, State), erlang:element(3, State), erlang:element(4, State), erlang:element(5, State), erlang:element(6, State), erlang:element(7, State), erlang:element(8, State), erlang:element(9, State) + 1}, gleam@otp@actor:continue(Updated_state); {error, {domain_error, _}} -> Updated_state@1 = {event_sourcing, erlang:element(2, State), erlang:element(3, State), erlang:element(4, State), erlang:element(5, State), erlang:element(6, State), erlang:element(7, State), erlang:element(8, State), erlang:element(9, State) + 1}, gleam@otp@actor:continue(Updated_state@1); {error, Error@1} -> gleam@otp@actor:stop_abnormal(describe_error(Error@1)) end; {get_system_stats, Reply_to@4} -> Stats = {system_stats, erlang:length(erlang:element(3, State)), erlang:element(9, State)}, gleam@erlang@process:send(Reply_to@4, Stats), gleam@otp@actor:continue(State); {get_aggregate_stats, Aggregate_id@5, Reply_to@5} -> Events_result = begin (erlang:element(4, erlang:element(2, State)))( fun(Tx@5) -> (erlang:element(7, erlang:element(2, State)))( erlang:element(10, erlang:element(2, State)), Tx@5, Aggregate_id@5, 0 ) end ) end, Result@5 = case Events_result of {ok, Events@1} -> Event_count = erlang:length(Events@1), Current_sequence = case Events@1 of [] -> 0; _ -> _pipe@4 = Events@1, _pipe@5 = gleam@list:last(_pipe@4), _pipe@6 = case _pipe@5 of {ok, X} -> {ok, erlang:element(3, X)}; {error, E} -> {error, E} end, gleam@result:unwrap(_pipe@6, 0) end, Has_snapshot = case erlang:element(7, State) of {some, _} -> case begin (erlang:element(5, erlang:element(2, State)))( fun(Tx_snap) -> (erlang:element( 8, erlang:element(2, State) ))(Tx_snap, Aggregate_id@5) end ) end of {ok, _} -> true; {error, _} -> false end; none -> false end, Stats@1 = {aggregate_stats, Aggregate_id@5, Event_count, Current_sequence, Has_snapshot}, {ok, Stats@1}; {error, Error@2} -> {error, Error@2} end, gleam@erlang@process:send(Reply_to@5, Result@5), gleam@otp@actor:continue(State) end. -file("src/eventsourcing.gleam", 397). -spec start( gleam@erlang@process:name(aggregate_message(JVS, JVT, JVU, JVV)), event_store(any(), JVS, JVT, JVU, JVV, any()), fun((JVS, JVT) -> {ok, list(JVU)} | {error, JVV}), list(gleam@erlang@process:name(query_message(JVU))), fun((JVS, JVU) -> JVS), JVS, gleam@option:option(snapshot_config()) ) -> {ok, gleam@otp@actor:started(gleam@erlang@process:subject(aggregate_message(JVS, JVT, JVU, JVV)))} | {error, gleam@otp@actor:start_error()}. start( Name, Eventstore, Handle, Query_actors, Apply, Empty_state, Snapshot_config ) -> _pipe = gleam@otp@actor:new( {event_sourcing, Eventstore, Query_actors, Handle, Apply, Empty_state, Snapshot_config, gleam@time@timestamp:system_time(), 0} ), _pipe@1 = gleam@otp@actor:on_message(_pipe, fun on_message/2), _pipe@2 = gleam@otp@actor:named(_pipe@1, Name), gleam@otp@actor:start(_pipe@2). -file("src/eventsourcing.gleam", 705). ?DOC( " Loads the current state of an aggregate asynchronously by replaying all its events.\n" " Returns a subject that will receive the aggregate with its current entity state and sequence number.\n" " Use process.receive() to get the result. Returns EntityNotFound error if aggregate doesn't exist.\n" "\n" " ## Example\n" " ```gleam\n" " let result = eventsourcing.load_aggregate(actor, \"bank-account-123\")\n" " let assert Ok(aggregate) = process.receive(result, 1000)\n" " // aggregate.entity contains current state, aggregate.sequence shows current version\n" " ```\n" ). -spec load_aggregate( gleam@erlang@process:subject(aggregate_message(JYU, JYV, JYW, JYX)), binary() ) -> gleam@erlang@process:subject({ok, aggregate(JYU, JYV, JYW, JYX)} | {error, event_sourcing_error(JYX)}). load_aggregate(Eventsourcing, Aggregate_id) -> Receiver = gleam@erlang@process:new_subject(), gleam@erlang@process:send( Eventsourcing, {load_aggregate, Aggregate_id, Receiver} ), Receiver. -file("src/eventsourcing.gleam", 728). ?DOC( " Loads all events for a specific aggregate asynchronously from the beginning.\n" " Returns a subject that will receive a chronologically ordered list of all events\n" " that have occurred for the aggregate. Use process.receive() to get the result.\n" "\n" " ## Example\n" " ```gleam\n" " let result = eventsourcing.load_events(eventsourcing, \"bank-account-123\")\n" " let assert Ok(events) = process.receive(result, 1000)\n" " // events contains all EventEnvelop items for this aggregate\n" " ```\n" ). -spec load_events( gleam@erlang@process:subject(aggregate_message(any(), any(), JZN, JZO)), binary() ) -> gleam@erlang@process:subject({ok, list(event_envelop(JZN))} | {error, event_sourcing_error(JZO)}). load_events(Eventsourcing, Aggregate_id) -> Receiver = gleam@erlang@process:new_subject(), gleam@erlang@process:send( Eventsourcing, {load_all_events, Aggregate_id, Receiver} ), Receiver. -file("src/eventsourcing.gleam", 751). ?DOC( " Loads events for an aggregate asynchronously starting from a specific sequence number.\n" " Returns a subject that will receive the events list. Useful for pagination or continuing\n" " event processing from a known point. Use process.receive() to get the result.\n" "\n" " ## Example\n" " ```gleam\n" " let result = eventsourcing.load_events_from(eventsourcing, \"bank-account-123\", start_from: 10)\n" " let assert Ok(events) = process.receive(result, 1000)\n" " // events contains EventEnvelop items starting from sequence 10\n" " ```\n" ). -spec load_events_from( gleam@erlang@process:subject(aggregate_message(any(), any(), KAC, KAD)), binary(), integer() ) -> gleam@erlang@process:subject({ok, list(event_envelop(KAC))} | {error, event_sourcing_error(KAD)}). load_events_from(Eventsourcing, Aggregate_id, Start_from) -> Receiver = gleam@erlang@process:new_subject(), gleam@erlang@process:send( Eventsourcing, {load_events, Aggregate_id, Start_from, Receiver} ), Receiver. -file("src/eventsourcing.gleam", 776). ?DOC( " Gets system statistics including command count and query actor health.\n" " Returns a subject that will receive the SystemStats. Use process.receive() to get the result.\n" " Useful for monitoring system health and performance in production.\n" "\n" " ## Example\n" " ```gleam\n" " let stats_subject = eventsourcing.system_stats(eventsourcing_actor)\n" " let assert Ok(stats) = process.receive(stats_subject, 1000)\n" " io.println(\"Query actors: \" <> int.to_string(stats.query_actors_count))\n" " io.println(\"Commands processed: \" <> int.to_string(stats.total_commands_processed))\n" " ```\n" ). -spec system_stats( gleam@erlang@process:subject(aggregate_message(any(), any(), any(), any())) ) -> gleam@erlang@process:subject(system_stats()). system_stats(Eventsourcing) -> Receiver = gleam@erlang@process:new_subject(), gleam@erlang@process:send(Eventsourcing, {get_system_stats, Receiver}), Receiver. -file("src/eventsourcing.gleam", 796). ?DOC( " Gets statistics for a specific aggregate including event count and snapshot status.\n" " Returns a subject that will receive the AggregateStats result. Use process.receive() to get the result.\n" " Useful for debugging and monitoring individual aggregate health.\n" "\n" " ## Example\n" " ```gleam\n" " let result = eventsourcing.aggregate_stats(eventsourcing, \"bank-account-123\")\n" " let assert Ok(stats) = process.receive(result, 1000)\n" " io.println(\"Events: \" <> int.to_string(stats.event_count))\n" " ```\n" ). -spec aggregate_stats( gleam@erlang@process:subject(aggregate_message(any(), any(), any(), KBC)), binary() ) -> gleam@erlang@process:subject({ok, aggregate_stats()} | {error, event_sourcing_error(KBC)}). aggregate_stats(Eventsourcing, Aggregate_id) -> Receiver = gleam@erlang@process:new_subject(), gleam@erlang@process:send( Eventsourcing, {get_aggregate_stats, Aggregate_id, Receiver} ), Receiver. -file("src/eventsourcing.gleam", 817). ?DOC( " Retrieves the most recent snapshot for an aggregate asynchronously if snapshots are enabled.\n" " Returns a subject that will receive the snapshot option. Snapshots provide a point-in-time\n" " capture of aggregate state for faster reconstruction. Use process.receive() to get the result.\n" "\n" " ## Example\n" " ```gleam\n" " let result = eventsourcing.latest_snapshot(eventsourcing, \"bank-account-123\")\n" " let assert Ok(Some(snapshot)) = process.receive(result, 1000)\n" " // snapshot.entity contains the saved state, snapshot.sequence shows version\n" " ```\n" ). -spec latest_snapshot( gleam@erlang@process:subject(aggregate_message(KBM, any(), any(), KBP)), binary() ) -> gleam@erlang@process:subject({ok, gleam@option:option(snapshot(KBM))} | {error, event_sourcing_error(KBP)}). latest_snapshot(Eventsourcing, Aggregate_id) -> Receiver = gleam@erlang@process:new_subject(), gleam@erlang@process:send( Eventsourcing, {load_latest_snapshot, Aggregate_id, Receiver} ), Receiver. -file("src/eventsourcing.gleam", 830). -spec start_query( gleam@erlang@process:name(query_message(KCB)), fun((binary(), list(event_envelop(KCB))) -> nil) ) -> {ok, gleam@otp@actor:started(gleam@erlang@process:subject(query_message(KCB)))} | {error, gleam@otp@actor:start_error()}. start_query(Name, Query) -> _pipe = gleam@otp@actor:new(nil), _pipe@1 = gleam@otp@actor:on_message( _pipe, fun(_, Message) -> case Message of {process_events, Aggregate_id, Events} -> Query(Aggregate_id, Events), gleam@otp@actor:continue(nil) end end ), _pipe@2 = gleam@otp@actor:named(_pipe@1, Name), gleam@otp@actor:start(_pipe@2). -file("src/eventsourcing.gleam", 303). ?DOC( " Creates a supervised event sourcing architecture with fault tolerance.\n" " Sets up a supervision tree where the main event sourcing actor and query actors\n" " are managed by a supervisor that can restart them if they fail. This is the \n" " recommended approach for production applications.\n" " For the queries to work you have to register them after the supervisor is started using register_queries().\n" "\n" " ## Example\n" " ```gleam\n" " // First set up memory store\n" " let events_name = process.new_name(\"events_actor\")\n" " let snapshot_name = process.new_name(\"snapshot_actor\")\n" " let #(store, _) = memory_store.supervised(events_name, snapshot_name, static_supervisor.OneForOne)\n" " \n" " // Then create event sourcing system\n" " let balance_query = #(process.new_name(\"balance_query\"), fn(aggregate_id, events) { /* update read model */ })\n" " let assert Ok(spec) = eventsourcing.supervised(\n" " name: process.new_name(\"eventsourcing_actor\"),\n" " eventstore: store,\n" " handle: my_handle,\n" " apply: my_apply,\n" " empty_state: MyState,\n" " queries: [balance_query],\n" " snapshot_config: None\n" " )\n" " ```\n" ). -spec supervised( gleam@erlang@process:name(aggregate_message(JTU, JTV, JTW, JTX)), event_store(any(), JTU, JTV, JTW, JTX, any()), fun((JTU, JTV) -> {ok, list(JTW)} | {error, JTX}), fun((JTU, JTW) -> JTU), JTU, list({gleam@erlang@process:name(query_message(JTW)), fun((binary(), list(event_envelop(JTW))) -> nil)}), gleam@option:option(snapshot_config()) ) -> {ok, gleam@otp@supervision:child_specification(gleam@otp@static_supervisor:supervisor())} | {error, nil}. supervised( Name, Eventstore, Handle, Apply, Empty_state, Queries, Snapshot_config ) -> Queries@1 = gleam@list:map( Queries, fun(Query) -> {Name@1, Query@1} = Query, {Name@1, gleam@otp@supervision:worker( fun() -> start_query(Name@1, Query@1) end )} end ), Names = gleam@list:map( Queries@1, fun(Query@2) -> erlang:element(1, Query@2) end ), Specs = gleam@list:map( Queries@1, fun(Query@3) -> erlang:element(2, Query@3) end ), Eventsourcing_spec = gleam@otp@supervision:worker( fun() -> gleam@result:'try'( start( Name, Eventstore, Handle, Names, Apply, Empty_state, Snapshot_config ), fun(Eventsourcing) -> {ok, Eventsourcing} end ) end ), Supervisor@1 = begin _pipe = gleam@otp@static_supervisor:new(one_for_one), _pipe@1 = gleam@otp@static_supervisor:add(_pipe, Eventsourcing_spec), _pipe@2 = gleam@list:fold( Specs, _pipe@1, fun(Supervisor, Spec) -> gleam@otp@static_supervisor:add(Supervisor, Spec) end ), _pipe@3 = gleam@otp@static_supervisor:supervised(_pipe@2), {ok, _pipe@3} end, Supervisor@1. -file("src/eventsourcing.gleam", 853). ?DOC( " Creates a validated timeout value for use with process operations.\n" " Ensures that timeout values are positive, preventing invalid configurations\n" " that could cause system operations to behave unexpectedly.\n" "\n" " ## Example\n" " ```gleam\n" " let assert Ok(timeout) = eventsourcing.timeout(5000)\n" " let result = process.receive(subject, timeout)\n" " ```\n" ). -spec timeout(integer()) -> {ok, timeout_()} | {error, event_sourcing_error(any())}. timeout(Ms) -> case Ms =< 0 of true -> {error, non_positive_argument}; false -> {ok, {timeout, Ms}} end. -file("src/eventsourcing.gleam", 869). ?DOC( " Creates a validated frequency value for snapshot configuration.\n" " Snapshots will be created every N events when this frequency is used.\n" " Ensures that frequency values are positive to prevent division by zero or infinite loops.\n" "\n" " ## Example\n" " ```gleam\n" " let assert Ok(freq) = eventsourcing.frequency(5)\n" " let config = eventsourcing.SnapshotConfig(freq)\n" " ```\n" ). -spec frequency(integer()) -> {ok, frequency()} | {error, event_sourcing_error(any())}. frequency(N) -> case N =< 0 of true -> {error, non_positive_argument}; false -> {ok, {frequency, N}} end.