Current section

Files

Jump to
eventsourcing src eventsourcing@sqlite_store.erl
Raw

src/eventsourcing@sqlite_store.erl

-module(eventsourcing@sqlite_store).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch]).
-export([load_events/2, new/9, load_aggregate_entity/2, create_event_table/1]).
-export_type([sqlite_store/4]).
-opaque sqlite_store(QFO, QFP, QFQ, QFR) :: {sqlite_store,
sqlight:connection(),
eventsourcing:aggregate(QFO, QFP, QFQ, QFR),
fun((QFQ) -> binary()),
fun((binary()) -> {ok, QFQ} |
{error, list(gleam@dynamic:decode_error())}),
binary(),
binary(),
binary()}.
-spec wrap_events(
sqlite_store(any(), any(), QIQ, any()),
binary(),
list(QIQ),
integer()
) -> list(eventsourcing:event_envelop(QIQ)).
wrap_events(Postgres_store, Aggregate_id, Events, Sequence) ->
_pipe = gleam@list:map_fold(
Events,
Sequence,
fun(Sequence@1, Event) ->
Next_sequence = Sequence@1 + 1,
{Next_sequence,
{serialized_event_envelop,
Aggregate_id,
Sequence@1 + 1,
Event,
erlang:element(6, Postgres_store),
erlang:element(7, Postgres_store),
erlang:element(8, Postgres_store)}}
end
),
gleam@pair:second(_pipe).
-spec persist_events(
sqlite_store(any(), any(), QJB, any()),
list(eventsourcing:event_envelop(QJB))
) -> list({ok, list(gleam@dynamic:dynamic_())} | {error, sqlight:error()}).
persist_events(Sqlite_store, Wrapped_events) ->
_pipe = Wrapped_events,
gleam@list:map(
_pipe,
fun(Event) ->
{serialized_event_envelop,
Aggregate_id,
Sequence,
Payload,
Event_type,
Event_version,
Aggregate_type} = case Event of
{serialized_event_envelop, _, _, _, _, _, _} -> Event;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Assertion pattern match failed"/utf8>>,
value => _assert_fail,
module => <<"eventsourcing/sqlite_store"/utf8>>,
function => <<"persist_events"/utf8>>,
line => 217})
end,
sqlight:'query'(
<<"
INSERT INTO event
(aggregate_type, aggregate_id, sequence, event_type, event_version, payload)
VALUES
($1, $2, $3, $4, $5, $6)
"/utf8>>,
erlang:element(2, Sqlite_store),
[sqlight:text(Aggregate_type),
sqlight:text(Aggregate_id),
sqlight:int(Sequence),
sqlight:text(Event_type),
sqlight:text(Event_version),
sqlight:text(
begin
_pipe@1 = Payload,
(erlang:element(4, Sqlite_store))(_pipe@1)
end
)],
fun gleam@dynamic:dynamic/1
)
end
).
-spec commit(
sqlite_store(QIA, QIB, QIC, QID),
eventsourcing:aggregate_context(QIA, QIB, QIC, QID),
list(QIC)
) -> list(eventsourcing:event_envelop(QIC)).
commit(Sqlite_store, Context, Events) ->
{aggregate_context, Aggregate_id, _, Sequence} = Context,
Wrapped_events = wrap_events(Sqlite_store, Aggregate_id, Events, Sequence),
persist_events(Sqlite_store, 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>>
),
Wrapped_events.
-spec load_events(sqlite_store(any(), any(), QHH, any()), binary()) -> {ok,
list(eventsourcing:event_envelop(QHH))} |
{error, sqlight:error()}.
load_events(Sqlite_store, Aggregate_id) ->
gleam@result:map(
sqlight:'query'(
<<"
SELECT aggregate_type, aggregate_id, sequence, event_type, event_version, payload
FROM event
WHERE aggregate_type = $1 AND aggregate_id = $2
ORDER BY sequence
"/utf8>>,
erlang:element(2, Sqlite_store),
[sqlight:text(erlang:element(8, Sqlite_store)),
sqlight:text(Aggregate_id)],
gleam@dynamic:decode6(
fun(Field@0, Field@1, Field@2, Field@3, Field@4, Field@5) -> {serialized_event_envelop, Field@0, Field@1, Field@2, Field@3, Field@4, Field@5} end,
gleam@dynamic:element(1, fun gleam@dynamic:string/1),
gleam@dynamic:element(2, fun gleam@dynamic:int/1),
gleam@dynamic:element(
5,
fun(Dyn) ->
_assert_subject = begin
_pipe = gleam@dynamic:string(Dyn),
gleam@result:map(
_pipe,
erlang:element(5, Sqlite_store)
)
end,
{ok, Payload} = 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/sqlite_store"/utf8>>,
function => <<"load_events"/utf8>>,
line => 130})
end,
Payload
end
),
gleam@dynamic:element(3, fun gleam@dynamic:string/1),
gleam@dynamic:element(4, fun gleam@dynamic:string/1),
gleam@dynamic:element(0, fun gleam@dynamic:string/1)
)
),
fun(Resulted) -> Resulted end
).
-spec load_aggregate(sqlite_store(QHO, QHP, QHQ, QHR), binary()) -> eventsourcing:aggregate_context(QHO, QHP, QHQ, QHR).
load_aggregate(Sqlite_store, Aggregate_id) ->
_assert_subject = load_events(Sqlite_store, Aggregate_id),
{ok, Commited_events} = 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/sqlite_store"/utf8>>,
function => <<"load_aggregate"/utf8>>,
line => 146})
end,
{Aggregate@1, Sequence} = gleam@list:fold(
Commited_events,
{erlang:element(3, Sqlite_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 new(
sqlight:connection(),
QFS,
fun((QFS, QFT) -> {ok, list(QFU)} | {error, QFV}),
fun((QFS, QFU) -> QFS),
fun((QFU) -> binary()),
fun((binary()) -> {ok, QFU} | {error, list(gleam@dynamic:decode_error())}),
binary(),
binary(),
binary()
) -> eventsourcing:event_store(sqlite_store(QFS, QFT, QFU, QFV), QFS, QFT, QFU, QFV).
new(
Sqlight_connection,
Empty_entity,
Handle,
Apply,
Event_encoder,
Event_decoder,
Event_type,
Event_version,
Aggregate_type
) ->
Eventstore = {sqlite_store,
Sqlight_connection,
{aggregate, Empty_entity, Handle, Apply},
Event_encoder,
Event_decoder,
Event_type,
Event_version,
Aggregate_type},
{event_store, Eventstore, fun load_aggregate/2, fun commit/3}.
-spec load_aggregate_entity(sqlite_store(QGX, any(), any(), any()), binary()) -> QGX.
load_aggregate_entity(Sqlite_store, Aggregate_id) ->
erlang:element(
2,
erlang:element(3, load_aggregate(Sqlite_store, Aggregate_id))
).
-spec create_event_table(sqlite_store(any(), any(), any(), any())) -> {ok,
list(gleam@dynamic:dynamic_())} |
{error, sqlight:error()}.
create_event_table(Sqlite_store) ->
sqlight:'query'(
<<"
CREATE TABLE IF NOT EXISTS event
(
aggregate_type text NOT NULL,
aggregate_id text NOT NULL,
sequence bigint CHECK (sequence >= 0) NOT NULL,
event_type text NOT NULL,
event_version text NOT NULL,
payload text NOT NULL,
PRIMARY KEY (aggregate_type, aggregate_id, sequence)
);
"/utf8>>,
erlang:element(2, Sqlite_store),
[],
fun gleam@dynamic:dynamic/1
).