Current section
Files
Jump to
Current section
Files
src/eventsourcing.erl
-module(eventsourcing).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/eventsourcing.gleam").
-export([register_queries/2, 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/8, 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@otp@actor:started(gleam@erlang@process:subject(query_message(JSA))),
fun((binary(), list(event_envelop(JSA))) -> nil)}.
-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()})} |
{register_query_actor, query_actor(JSE)} |
{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(query_actor(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", 372).
?DOC(
" Registers query actors with the event sourcing system after supervisor startup.\n"
" Queries must be added beforehand to the supervised() function to ensure their actors are started and supervised.\n"
"\n"
" ## Example\n"
" ```gleam\n"
" let queries = [projection_query, analytics_query]\n"
" let assert Ok(query_actors) = list.try_map(queries, fn(_) { \n"
" process.receive(query_receiver, 1000) \n"
" })\n"
" eventsourcing.register_queries(eventsourcing_actor, query_actors)\n"
" ```\n"
).
-spec register_queries(
gleam@otp@actor:started(gleam@erlang@process:subject(aggregate_message(any(), any(), JVD, any()))),
list(query_actor(JVD))
) -> nil.
register_queries(Eventsourcing_actor, Query_actors) ->
_pipe = Query_actors,
gleam@list:each(
_pipe,
fun(Query_actor) ->
gleam@erlang@process:send(
erlang:element(3, Eventsourcing_actor),
{register_query_actor, Query_actor}
)
end
).
-file("src/eventsourcing.gleam", 394).
?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@otp@actor:started(gleam@erlang@process:subject(aggregate_message(any(), JVO, any(), any()))),
binary(),
JVO
) -> nil.
execute(Eventsourcing_actor, Aggregate_id, Command) ->
gleam@erlang@process:send(
erlang:element(3, Eventsourcing_actor),
{execute_command, Aggregate_id, Command, []}
).
-file("src/eventsourcing.gleam", 416).
?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@otp@actor:started(gleam@erlang@process:subject(aggregate_message(any(), JVY, any(), any()))),
binary(),
JVY,
list({binary(), binary()})
) -> nil.
execute_with_metadata(Eventsourcing_actor, Aggregate_id, Command, Metadata) ->
gleam@erlang@process:send(
erlang:element(3, Eventsourcing_actor),
{execute_command, Aggregate_id, Command, Metadata}
).
-file("src/eventsourcing.gleam", 675).
-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", 692).
-spec load_aggregate_or_create_new(
event_sourcing(any(), JYO, JYP, JYQ, JYR, JYS),
JYS,
binary()
) -> {ok, aggregate(JYO, JYP, JYQ, JYR)} | {error, event_sourcing_error(JYR)}.
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", 461).
-spec on_message(
event_sourcing(JXJ, JXK, JXL, JXM, JXN, JXO),
aggregate_message(JXK, JXL, JXM, JXN)
) -> gleam@otp@actor:next(event_sourcing(JXJ, JXK, JXL, JXM, JXN, JXO), aggregate_message(JXK, JXL, JXM, JXN)).
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);
{register_query_actor, Query_actor} ->
New_state = {event_sourcing,
erlang:element(2, State),
[Query_actor | 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)},
gleam@otp@actor:continue(New_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(
erlang:element(
3,
erlang:element(
2,
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", 430).
-spec start(
event_store(any(), JWJ, JWK, JWL, JWM, any()),
fun((JWJ, JWK) -> {ok, list(JWL)} | {error, JWM}),
list(query_actor(JWL)),
fun((JWJ, JWL) -> JWJ),
JWJ,
gleam@option:option(snapshot_config())
) -> {ok,
gleam@otp@actor:started(gleam@erlang@process:subject(aggregate_message(JWJ, JWK, JWL, JWM)))} |
{error, gleam@otp@actor:start_error()}.
start(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),
gleam@otp@actor:start(_pipe@1).
-file("src/eventsourcing.gleam", 743).
?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@otp@actor:started(gleam@erlang@process:subject(aggregate_message(JZG, JZH, JZI, JZJ))),
binary()
) -> gleam@erlang@process:subject({ok, aggregate(JZG, JZH, JZI, JZJ)} |
{error, event_sourcing_error(JZJ)}).
load_aggregate(Eventsourcing_actor, Aggregate_id) ->
Receiver = gleam@erlang@process:new_subject(),
gleam@erlang@process:send(
erlang:element(3, Eventsourcing_actor),
{load_aggregate, Aggregate_id, Receiver}
),
Receiver.
-file("src/eventsourcing.gleam", 766).
?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(actor, \"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@otp@actor:started(gleam@erlang@process:subject(aggregate_message(any(), any(), KAA, KAB))),
binary()
) -> gleam@erlang@process:subject({ok, list(event_envelop(KAA))} |
{error, event_sourcing_error(KAB)}).
load_events(Eventsourcing_actor, Aggregate_id) ->
Receiver = gleam@erlang@process:new_subject(),
gleam@erlang@process:send(
erlang:element(3, Eventsourcing_actor),
{load_all_events, Aggregate_id, Receiver}
),
Receiver.
-file("src/eventsourcing.gleam", 789).
?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(actor, \"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@otp@actor:started(gleam@erlang@process:subject(aggregate_message(any(), any(), KAQ, KAR))),
binary(),
integer()
) -> gleam@erlang@process:subject({ok, list(event_envelop(KAQ))} |
{error, event_sourcing_error(KAR)}).
load_events_from(Eventsourcing_actor, Aggregate_id, Start_from) ->
Receiver = gleam@erlang@process:new_subject(),
gleam@erlang@process:send(
erlang:element(3, Eventsourcing_actor),
{load_events, Aggregate_id, Start_from, Receiver}
),
Receiver.
-file("src/eventsourcing.gleam", 816).
?DOC(
" Gets system statistics including uptime, 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 result = eventsourcing.system_stats(actor)\n"
" let assert Ok(stats) = process.receive(result, 1000)\n"
" io.println(\"Uptime: \" <> int.to_string(stats.uptime_seconds) <> \" seconds\")\n"
" ```\n"
).
-spec system_stats(
gleam@otp@actor:started(gleam@erlang@process:subject(aggregate_message(any(), any(), any(), any())))
) -> gleam@erlang@process:subject(system_stats()).
system_stats(Eventsourcing_actor) ->
Receiver = gleam@erlang@process:new_subject(),
gleam@erlang@process:send(
erlang:element(3, Eventsourcing_actor),
{get_system_stats, Receiver}
),
Receiver.
-file("src/eventsourcing.gleam", 836).
?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(actor, \"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@otp@actor:started(gleam@erlang@process:subject(aggregate_message(any(), any(), any(), KBS))),
binary()
) -> gleam@erlang@process:subject({ok, aggregate_stats()} |
{error, event_sourcing_error(KBS)}).
aggregate_stats(Eventsourcing_actor, Aggregate_id) ->
Receiver = gleam@erlang@process:new_subject(),
gleam@erlang@process:send(
erlang:element(3, Eventsourcing_actor),
{get_aggregate_stats, Aggregate_id, Receiver}
),
Receiver.
-file("src/eventsourcing.gleam", 860).
?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(actor, \"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@otp@actor:started(gleam@erlang@process:subject(aggregate_message(KCD, any(), any(), KCG))),
binary()
) -> gleam@erlang@process:subject({ok, gleam@option:option(snapshot(KCD))} |
{error, event_sourcing_error(KCG)}).
latest_snapshot(Eventsourcing_actor, Aggregate_id) ->
Receiver = gleam@erlang@process:new_subject(),
gleam@erlang@process:send(
erlang:element(3, Eventsourcing_actor),
{load_latest_snapshot, Aggregate_id, Receiver}
),
Receiver.
-file("src/eventsourcing.gleam", 876).
-spec start_query(fun((binary(), list(event_envelop(KCT))) -> nil)) -> {ok,
query_actor(KCT)} |
{error, gleam@otp@actor:start_error()}.
start_query(Query) ->
begin
Result = begin
_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
),
gleam@otp@actor:start(_pipe@1)
end,
case Result of
{ok, X} ->
{ok, {query_actor, X, Query}};
{error, E} ->
{error, E}
end
end.
-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"
"\n"
" ## Example\n"
" ```gleam\n"
" let assert Ok(eventstore) = memory_store.new()\n"
" let eventsourcing_actor_receiver = process.new_subject()\n"
" let query_actor_receiver = process.new_subject()\n"
" let assert Ok(spec) = eventsourcing.supervised(\n"
" eventstore:,\n"
" handle: my_handle,\n"
" apply: my_apply,\n"
" empty_state: MyState,\n"
" queries: [],\n"
" eventsourcing_actor_receiver:,\n"
" query_actors_receiver:,\n"
" snapshot_config: None\n"
" )\n"
" ```\n"
).
-spec supervised(
event_store(any(), JTV, JTW, JTX, JTY, any()),
fun((JTV, JTW) -> {ok, list(JTX)} | {error, JTY}),
fun((JTV, JTX) -> JTV),
JTV,
list(fun((binary(), list(event_envelop(JTX))) -> nil)),
gleam@erlang@process:subject(gleam@otp@actor:started(gleam@erlang@process:subject(aggregate_message(JTV, JTW, JTX, JTY)))),
gleam@erlang@process:subject(query_actor(JTX)),
gleam@option:option(snapshot_config())
) -> {ok,
gleam@otp@supervision:child_specification(gleam@otp@static_supervisor:supervisor())} |
{error, nil}.
supervised(
Eventstore,
Handle,
Apply,
Empty_state,
Queries,
Eventsourcing_actor_receiver,
Query_actors_receiver,
Snapshot_config
) ->
Queries_child_specs = gleam@list:map(
Queries,
fun(Query) ->
gleam@otp@supervision:worker(
fun() ->
gleam@result:'try'(
start_query(Query),
fun(Query_actor) ->
gleam@erlang@process:send(
Query_actors_receiver,
Query_actor
),
{ok, erlang:element(2, Query_actor)}
end
)
end
)
end
),
Eventsourcing_spec = gleam@otp@supervision:worker(
fun() ->
gleam@result:'try'(
start(
Eventstore,
Handle,
[],
Apply,
Empty_state,
Snapshot_config
),
fun(Eventsourcing_actor) ->
gleam@erlang@process:send(
Eventsourcing_actor_receiver,
Eventsourcing_actor
),
{ok, Eventsourcing_actor}
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(
Queries_child_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", 903).
?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", 919).
?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.