-module(eventsourcing@postgres_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([postgres_store/4]). -opaque postgres_store(PRE, PRF, PRG, PRH) :: {postgres_store, gleam@pgo:connection(), eventsourcing:aggregate(PRE, PRF, PRG, PRH), fun((PRG) -> binary()), fun((binary()) -> {ok, PRG} | {error, list(gleam@dynamic:decode_error())}), binary(), binary(), binary()}. -spec wrap_events( postgres_store(any(), any(), PUL, any()), binary(), list(PUL), integer() ) -> list(eventsourcing:event_envelop(PUL)). 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( postgres_store(any(), any(), PUW, any()), list(eventsourcing:event_envelop(PUW)) ) -> list({ok, gleam@pgo:returned(gleam@dynamic:dynamic_())} | {error, gleam@pgo:query_error()}). persist_events(Postgres_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/postgres_store"/utf8>>, function => <<"persist_events"/utf8>>, line => 217}) end, gleam@pgo:execute( <<" INSERT INTO event (aggregate_type, aggregate_id, sequence, event_type, event_version, payload) VALUES ($1, $2, $3, $4, $5, $6) "/utf8>>, erlang:element(2, Postgres_store), [gleam_pgo_ffi:coerce(Aggregate_type), gleam_pgo_ffi:coerce(Aggregate_id), gleam_pgo_ffi:coerce(Sequence), gleam_pgo_ffi:coerce(Event_type), gleam_pgo_ffi:coerce(Event_version), gleam_pgo_ffi:coerce( begin _pipe@1 = Payload, (erlang:element(4, Postgres_store))(_pipe@1) end )], fun gleam@dynamic:dynamic/1 ) end ). -spec commit( postgres_store(PTV, PTW, PTX, PTY), eventsourcing:aggregate_context(PTV, PTW, PTX, PTY), list(PTX) ) -> list(eventsourcing:event_envelop(PTX)). commit(Postgres_store, Context, Events) -> {aggregate_context, Aggregate_id, _, Sequence} = Context, Wrapped_events = wrap_events(Postgres_store, Aggregate_id, Events, Sequence), persist_events(Postgres_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(postgres_store(any(), any(), PSZ, any()), binary()) -> {ok, list(eventsourcing:event_envelop(PSZ))} | {error, gleam@pgo:query_error()}. load_events(Postgres_store, Aggregate_id) -> gleam@result:map( gleam@pgo:execute( <<" 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, Postgres_store), [gleam_pgo_ffi:coerce(erlang:element(8, Postgres_store)), gleam_pgo_ffi:coerce(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, Postgres_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/postgres_store"/utf8>>, function => <<"load_events"/utf8>>, line => 129}) 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) -> erlang:element(3, Resulted) end ). -spec load_aggregate(postgres_store(PTJ, PTK, PTL, PTM), binary()) -> eventsourcing:aggregate_context(PTJ, PTK, PTL, PTM). load_aggregate(Postgres_store, Aggregate_id) -> _assert_subject = load_events(Postgres_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/postgres_store"/utf8>>, function => <<"load_aggregate"/utf8>>, line => 145}) end, {Aggregate@1, Sequence} = gleam@list:fold( Commited_events, {erlang:element(3, Postgres_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( gleam@pgo:config(), PRI, fun((PRI, PRJ) -> {ok, list(PRK)} | {error, PRL}), fun((PRI, PRK) -> PRI), fun((PRK) -> binary()), fun((binary()) -> {ok, PRK} | {error, list(gleam@dynamic:decode_error())}), binary(), binary(), binary() ) -> eventsourcing:event_store(postgres_store(PRI, PRJ, PRK, PRL), PRI, PRJ, PRK, PRL). new( Pgo_config, Empty_entity, Handle, Apply, Event_encoder, Event_decoder, Event_type, Event_version, Aggregate_type ) -> Db = gleam_pgo_ffi:connect(Pgo_config), Eventstore = {postgres_store, Db, {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(postgres_store(PSP, any(), any(), any()), binary()) -> PSP. load_aggregate_entity(Postgres_store, Aggregate_id) -> erlang:element( 2, erlang:element(3, load_aggregate(Postgres_store, Aggregate_id)) ). -spec create_event_table(postgres_store(any(), any(), any(), any())) -> {ok, gleam@pgo:returned(gleam@dynamic:dynamic_())} | {error, gleam@pgo:query_error()}. create_event_table(Postgres_store) -> gleam@pgo:execute( <<" 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, Postgres_store), [], fun gleam@dynamic:dynamic/1 ).