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]).
-export([load_events/2, load_aggregate_entity/2, new/3]).
-export_type([memory_store/4, message/1]).
-opaque memory_store(OIU, OIV, OIW, OIX) :: {memory_store,
gleam@erlang@process:subject(message(OIW)),
eventsourcing:aggregate(OIU, OIV, OIW, OIX)}.
-type message(OIY) :: {set, binary(), list(eventsourcing:event_envelop(OIY))} |
{get,
binary(),
gleam@erlang@process:subject({ok,
list(eventsourcing:event_envelop(OIY))} |
{error, nil})}.
-spec handle_message(
message(OKH),
gleam@dict:dict(binary(), list(eventsourcing:event_envelop(OKH)))
) -> gleam@otp@actor:next(any(), gleam@dict:dict(binary(), list(eventsourcing:event_envelop(OKH)))).
handle_message(Message, State) ->
case Message of
{set, Key, Value} ->
_pipe = State,
_pipe@1 = gleam@dict:insert(_pipe, Key, Value),
gleam@otp@actor:continue(_pipe@1);
{get, Key@1, Response} ->
Value@1 = begin
_pipe@2 = State,
gleam@dict:get(_pipe@2, Key@1)
end,
gleam@otp@actor:send(Response, Value@1),
gleam@otp@actor:continue(State)
end.
-spec load_commited_events(memory_store(any(), any(), OKN, any()), binary()) -> list(eventsourcing:event_envelop(OKN)).
load_commited_events(Memory_store, Aggregate_id) ->
_pipe = gleam@otp@actor:call(
erlang:element(2, Memory_store),
fun(_capture) -> {get, Aggregate_id, _capture} end,
10000
),
gleam@result:unwrap(_pipe, []).
-spec load_events(memory_store(any(), any(), OJR, any()), binary()) -> list(eventsourcing:event_envelop(OJR)).
load_events(Memory_store, Aggregate_id) ->
_pipe = load_commited_events(Memory_store, Aggregate_id),
(fun(Events) ->
gleam@io:println(
<<<<<<<<"loading: "/utf8,
(begin
_pipe@1 = Events,
_pipe@2 = erlang:length(_pipe@1),
gleam@int:to_string(_pipe@2)
end)/binary>>/binary,
" events for Aggregate ID '"/utf8>>/binary,
Aggregate_id/binary>>/binary,
"'"/utf8>>
),
Events
end)(_pipe).
-spec load_aggregate(memory_store(OKU, OKV, OKW, OKX), binary()) -> eventsourcing:aggregate_context(OKU, OKV, OKW, OKX).
load_aggregate(Memory_store, Aggregate_id) ->
Commited_events = load_events(Memory_store, Aggregate_id),
{Aggregate@1, Sequence} = gleam@list:fold(
Commited_events,
{erlang:element(3, Memory_store), 0},
fun(Aggregate_and_sequence, Event_envelop) ->
{Aggregate, _} = Aggregate_and_sequence,
{erlang:setelement(
2,
Aggregate,
(erlang:element(4, Aggregate))(
erlang:element(2, Aggregate),
erlang:element(4, Event_envelop)
)
),
erlang:element(3, Event_envelop)}
end
),
{aggregate_context, Aggregate_id, Aggregate@1, Sequence}.
-spec load_aggregate_entity(memory_store(OJZ, any(), any(), any()), binary()) -> OJZ.
load_aggregate_entity(Memory_store, Aggregate_id) ->
erlang:element(
2,
erlang:element(3, load_aggregate(Memory_store, Aggregate_id))
).
-spec wrap_events(binary(), integer(), list(OLV)) -> list(eventsourcing:event_envelop(OLV)).
wrap_events(Aggregate_id, Current_sequence, Events) ->
_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}}
end
),
gleam@pair:second(_pipe).
-spec commit(
memory_store(OLG, OLH, OLI, OLJ),
eventsourcing:aggregate_context(OLG, OLH, OLI, OLJ),
list(OLI)
) -> list(eventsourcing:event_envelop(OLI)).
commit(Memory_store, Context, Events) ->
{aggregate_context, Aggregate_id, _, Sequence} = Context,
Wrapped_events = wrap_events(Aggregate_id, Sequence, Events),
Past_events = load_commited_events(Memory_store, Aggregate_id),
Events@1 = lists:append(Past_events, Wrapped_events),
gleam@io:println(
<<<<<<<<"storing: "/utf8,
(begin
_pipe = Wrapped_events,
_pipe@1 = erlang:length(_pipe),
gleam@int:to_string(_pipe@1)
end)/binary>>/binary,
" events for Aggregate ID '"/utf8>>/binary,
Aggregate_id/binary>>/binary,
"'"/utf8>>
),
gleam@otp@actor:send(
erlang:element(2, Memory_store),
{set, Aggregate_id, Events@1}
),
Wrapped_events.
-spec new(
OJE,
fun((OJE, OJF) -> {ok, list(OJG)} | {error, OJH}),
fun((OJE, OJG) -> OJE)
) -> eventsourcing:event_store(memory_store(OJE, OJF, OJG, OJH), OJE, OJF, OJG, OJH).
new(Empty_entity, Handle, Apply) ->
_assert_subject = begin
_pipe = gleam@otp@actor:start(gleam@dict:new(), fun handle_message/2),
gleam@result:'try'(
_pipe,
fun(Subject) ->
{ok,
{memory_store,
Subject,
{aggregate, Empty_entity, Handle, Apply}}}
end
)
end,
{ok, Actor} = case _assert_subject of
{ok, _} -> _assert_subject;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Assertion pattern match failed"/utf8>>,
value => _assert_fail,
module => <<"eventsourcing/memory_store"/utf8>>,
function => <<"new"/utf8>>,
line => 42})
end,
{event_store, Actor, fun load_aggregate/2, fun commit/3}.