Current section

Files

Jump to
eventsourcing src eventsourcing.erl
Raw

src/eventsourcing.erl

-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.