Current section
Files
Jump to
Current section
Files
src/eventsourcing@memory_store.erl
-module(eventsourcing@memory_store).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/eventsourcing/memory_store.gleam").
-export([new/0, supervised/1]).
-export_type([memory_store/4, event_message/1, snapshot_message/1]).
-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 memory_store(NNF, NNG, NNH, NNI) :: {memory_store,
gleam@erlang@process:subject(event_message(NNH)),
gleam@erlang@process:subject(snapshot_message(NNF))} |
{gleam_phantom, NNG, NNI}.
-type event_message(NNJ) :: {set_events,
binary(),
list(eventsourcing:event_envelop(NNJ))} |
{get_events,
binary(),
gleam@erlang@process:subject(gleam@option:option(list(eventsourcing:event_envelop(NNJ))))}.
-type snapshot_message(NNK) :: {set_snapshot,
binary(),
eventsourcing:snapshot(NNK)} |
{get_snapshot,
binary(),
gleam@erlang@process:subject(gleam@option:option(eventsourcing:snapshot(NNK)))}.
-file("src/eventsourcing/memory_store.gleam", 154).
-spec handle_events_message(
gleam@dict:dict(binary(), list(eventsourcing:event_envelop(NOY))),
event_message(NOY)
) -> gleam@otp@actor:next(gleam@dict:dict(binary(), list(eventsourcing:event_envelop(NOY))), any()).
handle_events_message(State, Message) ->
case Message of
{set_events, Key, Value} ->
_pipe = State,
_pipe@1 = gleam@dict:insert(_pipe, Key, Value),
gleam@otp@actor:continue(_pipe@1);
{get_events, Key@1, Response} ->
Value@1 = begin
_pipe@2 = State,
_pipe@3 = gleam_stdlib:map_get(_pipe@2, Key@1),
gleam@option:from_result(_pipe@3)
end,
gleam@otp@actor:send(Response, Value@1),
gleam@otp@actor:continue(State)
end.
-file("src/eventsourcing/memory_store.gleam", 82).
-spec supervised_events_actor(
gleam@erlang@process:subject(gleam@otp@actor:started(gleam@erlang@process:subject(event_message(NTF))))
) -> gleam@otp@supervision:child_specification(gleam@erlang@process:subject(event_message(NTF))).
supervised_events_actor(Events_actor_receiver) ->
gleam@otp@supervision:worker(
fun() ->
gleam@result:'try'(
begin
_pipe = gleam@otp@actor:new(maps:new()),
_pipe@1 = gleam@otp@actor:on_message(
_pipe,
fun handle_events_message/2
),
gleam@otp@actor:start(_pipe@1)
end,
fun(Started) ->
gleam@erlang@process:send(Events_actor_receiver, Started),
{ok, Started}
end
)
end
).
-file("src/eventsourcing/memory_store.gleam", 201).
-spec wrap_events(binary(), integer(), list(NQJ), list({binary(), binary()})) -> list(eventsourcing:event_envelop(NQJ)).
wrap_events(Aggregate_id, Current_sequence, Events, Metadata) ->
_pipe = gleam@list:map_fold(
Events,
Current_sequence,
fun(Sequence, Event) ->
Next_sequence = Sequence + 1,
{Next_sequence,
{memory_store_event_envelop,
Aggregate_id,
Sequence + 1,
Event,
Metadata}}
end
),
gleam@pair:second(_pipe).
-file("src/eventsourcing/memory_store.gleam", 226).
-spec handle_snapshot_message(
gleam@dict:dict(binary(), eventsourcing:snapshot(NQO)),
snapshot_message(NQO)
) -> gleam@otp@actor:next(gleam@dict:dict(binary(), eventsourcing:snapshot(NQO)), any()).
handle_snapshot_message(State, Message) ->
case Message of
{set_snapshot, Key, Value} ->
_pipe = State,
_pipe@1 = gleam@dict:insert(_pipe, Key, Value),
gleam@otp@actor:continue(_pipe@1);
{get_snapshot, Key@1, Response} ->
Value@1 = begin
_pipe@2 = State,
_pipe@3 = gleam_stdlib:map_get(_pipe@2, Key@1),
gleam@option:from_result(_pipe@3)
end,
gleam@otp@actor:send(Response, Value@1),
gleam@otp@actor:continue(State)
end.
-file("src/eventsourcing/memory_store.gleam", 94).
-spec supervised_snapshot_actor(
gleam@erlang@process:subject(gleam@otp@actor:started(gleam@erlang@process:subject(snapshot_message(NUV))))
) -> gleam@otp@supervision:child_specification(gleam@erlang@process:subject(snapshot_message(NUV))).
supervised_snapshot_actor(Snapshot_actor_receiver) ->
gleam@otp@supervision:worker(
fun() ->
gleam@result:'try'(
begin
_pipe = gleam@otp@actor:new(maps:new()),
_pipe@1 = gleam@otp@actor:on_message(
_pipe,
fun handle_snapshot_message/2
),
gleam@otp@actor:start(_pipe@1)
end,
fun(Started) ->
gleam@erlang@process:send(Snapshot_actor_receiver, Started),
{ok, Started}
end
)
end
).
-file("src/eventsourcing/memory_store.gleam", 242).
-spec save_snapshot(
memory_store(NQS, any(), any(), NQV),
eventsourcing:snapshot(NQS)
) -> {ok, nil} | {error, eventsourcing:event_sourcing_error(NQV)}.
save_snapshot(Memory_store, Snapshot) ->
_pipe = gleam@otp@actor:send(
erlang:element(3, Memory_store),
{set_snapshot, erlang:element(2, Snapshot), Snapshot}
),
{ok, _pipe}.
-file("src/eventsourcing/memory_store.gleam", 167).
-spec load_events(
memory_store(any(), any(), NPE, NPF),
any(),
binary(),
integer()
) -> {ok, list(eventsourcing:event_envelop(NPE))} |
{error, eventsourcing:event_sourcing_error(NPF)}.
load_events(Memory_store, _, Aggregate_id, Start_from) ->
Events = gleam@erlang@process:call(
erlang:element(2, Memory_store),
1000,
fun(_capture) -> {get_events, Aggregate_id, _capture} end
),
case Events of
{some, Events@1} ->
{ok,
begin
_pipe = Events@1,
gleam@list:drop(_pipe, Start_from)
end};
none ->
{ok, []}
end.
-file("src/eventsourcing/memory_store.gleam", 185).
-spec commit_events(
memory_store(NPQ, NPR, NPS, NPT),
eventsourcing:aggregate(NPQ, NPR, NPS, NPT),
list(NPS),
list({binary(), binary()})
) -> {ok, {list(eventsourcing:event_envelop(NPS)), integer()}} |
{error, eventsourcing:event_sourcing_error(NPT)}.
commit_events(Memory_store, Aggregate, Events, Metadata) ->
{aggregate, Aggregate_id, _, Sequence} = Aggregate,
Wrapped_events = wrap_events(Aggregate_id, Sequence, Events, Metadata),
begin
Result = load_events(Memory_store, nil, Aggregate_id, 0),
case Result of
{ok, X} ->
{ok,
begin
All_events = lists:append(X, Wrapped_events),
gleam@otp@actor:send(
erlang:element(2, Memory_store),
{set_events, Aggregate_id, All_events}
),
Last_event@1 = case gleam@list:last(Wrapped_events) of
{ok, Last_event} -> Last_event;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"eventsourcing/memory_store"/utf8>>,
function => <<"commit_events"/utf8>>,
line => 196,
value => _assert_fail,
start => 5666,
'end' => 5719,
pattern_start => 5677,
pattern_end => 5691})
end,
{Wrapped_events, erlang:element(3, Last_event@1)}
end};
{error, E} ->
{error, E}
end
end.
-file("src/eventsourcing/memory_store.gleam", 253).
-spec load_snapshot(memory_store(NRE, any(), any(), NRH), binary()) -> {ok,
gleam@option:option(eventsourcing:snapshot(NRE))} |
{error, eventsourcing:event_sourcing_error(NRH)}.
load_snapshot(Memory_store, Aggregate_id) ->
{ok,
gleam@erlang@process:call(
erlang:element(3, Memory_store),
1000,
fun(_capture) -> {get_snapshot, Aggregate_id, _capture} end
)}.
-file("src/eventsourcing/memory_store.gleam", 49).
?DOC(" Create a new memory store record.\n").
-spec new() -> {ok,
eventsourcing:event_store(memory_store(NZZ, OAA, OAB, OAC), NZZ, OAA, OAB, OAC, memory_store(NZZ, OAA, OAB, OAC))} |
{error, gleam@otp@actor:start_error()}.
new() ->
begin
Result = begin
_pipe = gleam@otp@actor:new(maps:new()),
_pipe@1 = gleam@otp@actor:on_message(
_pipe,
fun handle_events_message/2
),
gleam@otp@actor:start(_pipe@1)
end,
case Result of
{ok, X} ->
Event_actor = X,
begin
Result@1 = begin
_pipe@2 = gleam@otp@actor:new(maps:new()),
_pipe@3 = gleam@otp@actor:on_message(
_pipe@2,
fun handle_snapshot_message/2
),
gleam@otp@actor:start(_pipe@3)
end,
case Result@1 of
{ok, X@1} ->
Memory_store = {memory_store,
erlang:element(3, Event_actor),
erlang:element(3, X@1)},
{ok,
{event_store,
fun(F) -> F(Memory_store) end,
fun(F@1) -> F@1(Memory_store) end,
fun(F@2) -> F@2(Memory_store) end,
fun(F@3) -> F@3(Memory_store) end,
fun commit_events/4,
fun load_events/4,
fun load_snapshot/2,
fun save_snapshot/2,
Memory_store}};
{error, E} ->
{error, E}
end
end;
{error, E@1} ->
{error, E@1}
end
end.
-file("src/eventsourcing/memory_store.gleam", 106).
-spec supervised(
gleam@erlang@process:subject(eventsourcing:event_store(memory_store(NNZ, NOA, NOB, NOC), NNZ, NOA, NOB, NOC, memory_store(NNZ, NOA, NOB, NOC)))
) -> {gleam@otp@supervision:child_specification(gleam@erlang@process:subject(event_message(NOB))),
gleam@otp@supervision:child_specification(gleam@erlang@process:subject(snapshot_message(NNZ)))}.
supervised(Event_store_receiver) ->
Events_actor_receiver = gleam@erlang@process:new_subject(),
Supervised_events_actor_spec = supervised_events_actor(
Events_actor_receiver
),
Snapshot_actor_receiver = gleam@erlang@process:new_subject(),
Supervised_snapshot_actor_spec = supervised_snapshot_actor(
Snapshot_actor_receiver
),
Event_actor@1 = case gleam@erlang@process:'receive'(
Events_actor_receiver,
1000
) of
{ok, Event_actor} -> Event_actor;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"eventsourcing/memory_store"/utf8>>,
function => <<"supervised"/utf8>>,
line => 129,
value => _assert_fail,
start => 3405,
'end' => 3478,
pattern_start => 3416,
pattern_end => 3431})
end,
Snapshot_actor@1 = case gleam@erlang@process:'receive'(
Snapshot_actor_receiver,
1000
) of
{ok, Snapshot_actor} -> Snapshot_actor;
_assert_fail@1 ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"eventsourcing/memory_store"/utf8>>,
function => <<"supervised"/utf8>>,
line => 130,
value => _assert_fail@1,
start => 3481,
'end' => 3559,
pattern_start => 3492,
pattern_end => 3510})
end,
Memory_store = {memory_store,
erlang:element(3, Event_actor@1),
erlang:element(3, Snapshot_actor@1)},
gleam@erlang@process:send(
Event_store_receiver,
{event_store,
fun(F) -> F(Memory_store) end,
fun(F@1) -> F@1(Memory_store) end,
fun(F@2) -> F@2(Memory_store) end,
fun(F@3) -> F@3(Memory_store) end,
fun commit_events/4,
fun load_events/4,
fun load_snapshot/2,
fun save_snapshot/2,
Memory_store}
),
{Supervised_events_actor_spec, Supervised_snapshot_actor_spec}.