Current section

Files

Jump to
eventsourcing src eventsourcing@memory_store.erl
Raw

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/2]).
-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(LRO, LRP, LRQ, LRR) :: {memory_store,
gleam@erlang@process:subject(event_message(LRQ)),
gleam@erlang@process:subject(snapshot_message(LRO))} |
{gleam_phantom, LRP, LRR}.
-type event_message(LRS) :: {set_events,
binary(),
list(eventsourcing:event_envelop(LRS))} |
{get_events,
binary(),
gleam@erlang@process:subject(gleam@option:option(list(eventsourcing:event_envelop(LRS))))}.
-type snapshot_message(LRT) :: {set_snapshot,
binary(),
eventsourcing:snapshot(LRT)} |
{get_snapshot,
binary(),
gleam@erlang@process:subject(gleam@option:option(eventsourcing:snapshot(LRT)))}.
-file("src/eventsourcing/memory_store.gleam", 157).
-spec handle_events_message(
gleam@dict:dict(binary(), list(eventsourcing:event_envelop(LTC))),
event_message(LTC)
) -> gleam@otp@actor:next(gleam@dict:dict(binary(), list(eventsourcing:event_envelop(LTC))), 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", 83).
-spec supervised_events_actor(
gleam@erlang@process:subject(gleam@otp@actor:started(gleam@erlang@process:subject(event_message(LXJ))))
) -> gleam@otp@supervision:child_specification(gleam@erlang@process:subject(event_message(LXJ))).
supervised_events_actor(Events_actor_receiver) ->
gleam@otp@supervision:worker(
fun() ->
gleam@result:map(
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),
Started
end
)
end
).
-file("src/eventsourcing/memory_store.gleam", 204).
-spec wrap_events(binary(), integer(), list(LUN), list({binary(), binary()})) -> list(eventsourcing:event_envelop(LUN)).
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", 229).
-spec handle_snapshot_message(
gleam@dict:dict(binary(), eventsourcing:snapshot(LUS)),
snapshot_message(LUS)
) -> gleam@otp@actor:next(gleam@dict:dict(binary(), eventsourcing:snapshot(LUS)), 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", 95).
-spec supervised_snapshot_actor(
gleam@erlang@process:subject(gleam@otp@actor:started(gleam@erlang@process:subject(snapshot_message(LYX))))
) -> gleam@otp@supervision:child_specification(gleam@erlang@process:subject(snapshot_message(LYX))).
supervised_snapshot_actor(Snapshot_actor_receiver) ->
gleam@otp@supervision:worker(
fun() ->
gleam@result:map(
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),
Started
end
)
end
).
-file("src/eventsourcing/memory_store.gleam", 245).
-spec save_snapshot(
memory_store(LUW, any(), any(), LUZ),
eventsourcing:snapshot(LUW)
) -> {ok, nil} | {error, eventsourcing:event_sourcing_error(LUZ)}.
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", 170).
-spec load_events(
memory_store(any(), any(), LTI, LTJ),
any(),
binary(),
integer()
) -> {ok, list(eventsourcing:event_envelop(LTI))} |
{error, eventsourcing:event_sourcing_error(LTJ)}.
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", 188).
-spec commit_events(
memory_store(LTU, LTV, LTW, LTX),
eventsourcing:aggregate(LTU, LTV, LTW, LTX),
list(LTW),
list({binary(), binary()})
) -> {ok, {list(eventsourcing:event_envelop(LTW)), integer()}} |
{error, eventsourcing:event_sourcing_error(LTX)}.
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 => 199,
value => _assert_fail,
start => 5762,
'end' => 5815,
pattern_start => 5773,
pattern_end => 5787})
end,
{Wrapped_events, erlang:element(3, Last_event@1)}
end};
{error, E} ->
{error, E}
end
end.
-file("src/eventsourcing/memory_store.gleam", 256).
-spec load_snapshot(memory_store(LVI, any(), any(), LVL), binary()) -> {ok,
gleam@option:option(eventsourcing:snapshot(LVI))} |
{error, eventsourcing:event_sourcing_error(LVL)}.
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", 50).
?DOC(" Create a new memory store record.\n").
-spec new() -> {ok,
eventsourcing:event_store(memory_store(MDZ, MEA, MEB, MEC), MDZ, MEA, MEB, MEC, memory_store(MDZ, MEA, MEB, MEC))} |
{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", 107).
-spec supervised(
gleam@erlang@process:subject(eventsourcing:event_store(memory_store(LSI, LSJ, LSK, LSL), LSI, LSJ, LSK, LSL, memory_store(LSI, LSJ, LSK, LSL))),
gleam@otp@static_supervisor:strategy()
) -> gleam@otp@supervision:child_specification(gleam@otp@static_supervisor:supervisor()).
supervised(Event_store_receiver, Strategy) ->
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 => 128,
value => _assert_fail,
start => 3380,
'end' => 3453,
pattern_start => 3391,
pattern_end => 3406})
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 => 129,
value => _assert_fail@1,
start => 3456,
'end' => 3534,
pattern_start => 3467,
pattern_end => 3485})
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}
),
_pipe = gleam@otp@static_supervisor:new(Strategy),
_pipe@1 = gleam@otp@static_supervisor:add(
_pipe,
Supervised_events_actor_spec
),
_pipe@2 = gleam@otp@static_supervisor:add(
_pipe@1,
Supervised_snapshot_actor_spec
),
gleam@otp@static_supervisor:supervised(_pipe@2).