-module(aarondb@transactor). -compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]). -define(FILEPATH, "src/aarondb/transactor.gleam"). -export([compute_next_state/4, start_with_timeout/2, start/1, start_named/2, start_distributed/2, retract_entity/3, log_query/2, get_state/1, set_schema/3, set_schema_with_timeout/4, register_function/3, register_composite/2, register_predicate/3, store_rule/2, set_config/2, transact/2, transact_with_timeout/3, retract/2, subscribe/1]). -export_type([message/0]). -type message() :: {transact, list({aarondb@fact:eid(), binary(), aarondb@fact:value()}), gleam@option:option(integer()), gleam@erlang@process:subject({ok, aarondb@shared@state:db_state()} | {error, binary()})} | {retract, list({aarondb@fact:eid(), binary(), aarondb@fact:value()}), gleam@option:option(integer()), gleam@erlang@process:subject({ok, aarondb@shared@state:db_state()} | {error, binary()})} | {get_state, gleam@erlang@process:subject(aarondb@shared@state:db_state())} | {set_schema, binary(), aarondb@fact:attribute_config(), gleam@erlang@process:subject({ok, nil} | {error, binary()})} | {register_function, binary(), fun((aarondb@shared@state:db_state(), integer(), integer(), list(aarondb@fact:value())) -> list({aarondb@fact:eid(), binary(), aarondb@fact:value()})), gleam@erlang@process:subject(nil)} | {register_predicate, binary(), fun((aarondb@fact:value()) -> boolean()), gleam@erlang@process:subject(nil)} | {register_composite, list(binary()), gleam@erlang@process:subject({ok, nil} | {error, binary()})} | {store_rule, aarondb@shared@ast:rule(), gleam@erlang@process:subject({ok, nil} | {error, binary()})} | {set_reactive, gleam@erlang@process:subject(aarondb@shared@state:reactive_message())} | {join, gleam@erlang@process:pid_()} | {sync_datoms, list(aarondb@fact:datom())} | {raft_msg, aarondb@raft:raft_message()} | {compact, gleam@erlang@process:subject(nil)} | {set_config, aarondb@shared@state:config(), gleam@erlang@process:subject(nil)} | {sync, gleam@erlang@process:subject(nil)} | {boot, gleam@option:option(binary()), aarondb@storage:storage_adapter(), gleam@erlang@process:subject(nil)} | {register_index_adapter, aarondb@shared@state:index_adapter(), gleam@erlang@process:subject(nil)} | {create_index, binary(), binary(), binary(), gleam@erlang@process:subject({ok, nil} | {error, binary()})} | {create_b_m25_index, binary(), gleam@erlang@process:subject({ok, nil} | {error, binary()})} | {subscribe, gleam@erlang@process:subject(list(aarondb@fact:datom()))} | {prune, integer(), list(binary()), gleam@erlang@process:subject(integer())} | {retract_entity, aarondb@fact:entity_id(), gleam@erlang@process:subject({ok, aarondb@shared@state:db_state()} | {error, binary()})} | tick | {log_query, aarondb@shared@state:query_context(), gleam@erlang@process:subject(nil)}. -file("src/aarondb/transactor.gleam", 163). -spec lifecycle_loop(gleam@erlang@process:subject(message())) -> any(). lifecycle_loop(Parent) -> gleam_erlang_ffi:sleep(5000), gleam@erlang@process:send(Parent, tick), lifecycle_loop(Parent). -file("src/aarondb/transactor.gleam", 302). -spec compute_next_state( aarondb@shared@state:db_state(), list({aarondb@fact:eid(), binary(), aarondb@fact:value()}), gleam@option:option(integer()), aarondb@fact:operation() ) -> {ok, {aarondb@shared@state:db_state(), list(aarondb@fact:datom())}} | {error, binary()}. compute_next_state(State, Facts, Valid_time, Op) -> Tx_id = erlang:element(6, State) + 1, Vt = gleam@option:unwrap(Valid_time, Tx_id), Resolved_facts = aarondb@transactor@apply:resolve_transaction_functions( State, Tx_id, Vt, Facts ), Datoms_res = gleam@list:fold_until( Resolved_facts, {ok, []}, fun(Acc_res, F) -> Acc@1 = case Acc_res of {ok, Acc} -> Acc; _assert_fail -> erlang:error(#{gleam_error => let_assert, message => <<"Pattern match failed, no pattern matched the value."/utf8>>, file => <>, module => <<"aarondb/transactor"/utf8>>, function => <<"compute_next_state"/utf8>>, line => 318, value => _assert_fail, start => 8764, 'end' => 8792, pattern_start => 8775, pattern_end => 8782}) end, Eid_res = case erlang:element(1, F) of {uid, Id} -> {ok, Id}; {lookup, Lu} -> {A, V} = Lu, case A =:= <<"db/fn"/utf8>> of true -> {error, <<"Unresolved transaction function: "/utf8, (gleam@string:inspect(V))/binary>>}; false -> _pipe = aarondb@index:get_entity_by_av( erlang:element(5, State), A, V ), gleam@result:replace_error( _pipe, <<"Lookup failed for "/utf8, A/binary>> ) end end, case Eid_res of {ok, Eid} -> D = {datom, Eid, erlang:element(2, F), erlang:element(3, F), Tx_id, erlang:length(Acc@1), Vt, Op}, {continue, {ok, [D | Acc@1]}}; {error, E} -> {stop, {error, E}} end end ), case Datoms_res of {ok, Datoms} -> Datoms@1 = lists:reverse(Datoms), {Final_state, All_datoms, _} = gleam@list:fold( Datoms@1, {State, [], 0}, fun(Acc@2, D@1) -> {Curr_state, Collected, Next_idx} = Acc@2, {New_state, Side_effects, Updated_idx} = aarondb@transactor@apply:apply_datom( Curr_state, {datom, erlang:element(2, D@1), erlang:element(3, D@1), erlang:element(4, D@1), erlang:element(5, D@1), Next_idx, erlang:element(7, D@1), erlang:element(8, D@1)}, Next_idx ), {New_state, lists:append(Side_effects, Collected), Updated_idx} end ), All_datoms@1 = lists:reverse(All_datoms), Validate_res = gleam@list:fold_until( All_datoms@1, {ok, nil}, fun(_, D@2) -> case aarondb@transactor@validation:validate_datom( State, All_datoms@1, D@2 ) of {ok, _} -> {continue, {ok, nil}}; {error, E@1} -> {stop, {error, E@1}} end end ), case Validate_res of {ok, _} -> {ok, {{db_state, erlang:element(2, Final_state), erlang:element(3, Final_state), erlang:element(4, Final_state), erlang:element(5, Final_state), Tx_id, erlang:element(7, Final_state), erlang:element(8, Final_state), erlang:element(9, Final_state), erlang:element(10, Final_state), erlang:element(11, Final_state), erlang:element(12, Final_state), erlang:element(13, Final_state), erlang:element(14, Final_state), erlang:element(15, Final_state), erlang:element(16, Final_state), erlang:element(17, Final_state), erlang:element(18, Final_state), erlang:element(19, Final_state), erlang:element(20, Final_state), erlang:element(21, Final_state), erlang:element(22, Final_state), erlang:element(23, Final_state), erlang:element(24, Final_state), erlang:element(25, Final_state), erlang:element(26, Final_state)}, All_datoms@1}}; {error, E@2} -> {error, E@2} end; {error, E@3} -> {error, E@3} end. -file("src/aarondb/transactor.gleam", 393). -spec handle_message(aarondb@shared@state:db_state(), message()) -> gleam@otp@actor:next(aarondb@shared@state:db_state(), message()). handle_message(State, Msg) -> case Msg of {log_query, Ctx, Reply} -> aarondb@transactor@messages:log_query(State, Ctx, Reply); tick -> gleam@otp@actor:continue( aarondb@transactor@lifecycle:handle_tick(State) ); {boot, Ets_name, _, Reply@1} -> case Ets_name of {some, Name} -> aarondb@index@ets:init_tables(Name); none -> nil end, _ = aarondb_mnesia_ffi:init(), New_state = aarondb@transactor@runtime:recover_state(State), gleam@erlang@process:send(Reply@1, nil), gleam@otp@actor:continue(New_state); {transact, Facts, Vt, Reply_to} -> aarondb@transactor@runtime:do_handle_transact( State, Facts, Vt, assert, Reply_to, fun compute_next_state/4 ); {retract, Facts@1, Vt@1, Reply_to@1} -> aarondb@transactor@runtime:do_handle_transact( State, Facts@1, Vt@1, retract, Reply_to@1, fun compute_next_state/4 ); {retract_entity, Eid, Reply_to@2} -> Datoms = case erlang:element(14, State) of {some, Name@1} -> aarondb@index@ets:lookup_datoms( <>, Eid ); none -> aarondb@index:filter_by_entity( erlang:element(3, State), Eid ) end, Facts@2 = gleam@list:map( Datoms, fun(D) -> {{uid, erlang:element(2, D)}, erlang:element(3, D), erlang:element(4, D)} end ), aarondb@transactor@runtime:do_handle_transact( State, Facts@2, none, retract, Reply_to@2, fun compute_next_state/4 ); {get_state, Reply_to@3} -> gleam@erlang@process:send(Reply_to@3, State), gleam@otp@actor:continue(State); {set_schema, Attr, Config, Reply_to@4} -> Error = case erlang:element(2, Config) of true -> aarondb@transactor@schema:validate_unique(State, Attr); false -> none end, Error@1 = case Error of none -> case erlang:element(5, Config) =:= one of true -> aarondb@transactor@schema:validate_cardinality_one( State, Attr ); false -> none end; {some, E} -> {some, E} end, aarondb@transactor@messages:set_schema( State, Attr, Config, Error@1, Reply_to@4 ); {register_function, Name@2, Func, Reply_to@5} -> aarondb@transactor@messages:register_function( State, Name@2, Func, Reply_to@5 ); {register_predicate, Name@3, Pred, Reply_to@6} -> aarondb@transactor@messages:register_predicate( State, Name@3, Pred, Reply_to@6 ); {register_composite, Attrs, Reply_to@7} -> aarondb@transactor@messages:register_composite( State, Attrs, aarondb@transactor@schema:validate_composite(State, Attrs), Reply_to@7 ); {store_rule, Rule, Reply_to@8} -> aarondb@transactor@messages:store_rule( State, Rule, Reply_to@8, fun compute_next_state/4 ); {subscribe, Reply_to@9} -> aarondb@transactor@messages:subscribe(State, Reply_to@9); {set_config, Config@1, Reply_to@10} -> aarondb@transactor@messages:set_config(State, Config@1, Reply_to@10); _ -> gleam@otp@actor:continue(State) end. -file("src/aarondb/transactor.gleam", 93). -spec do_start_named( aarondb@storage:storage_adapter(), boolean(), gleam@option:option(binary()) ) -> {ok, gleam@erlang@process:subject(message())} | {error, gleam@otp@actor:start_error()}. do_start_named(Store, Is_distributed, Ets_name) -> Reactive_subject@1 = case aarondb@reactive:start_link() of {ok, Reactive_subject} -> Reactive_subject; _assert_fail -> erlang:error(#{gleam_error => let_assert, message => <<"Pattern match failed, no pattern matched the value."/utf8>>, file => <>, module => <<"aarondb/transactor"/utf8>>, function => <<"do_start_named"/utf8>>, line => 98, value => _assert_fail, start => 2983, 'end' => 3038, pattern_start => 2994, pattern_end => 3014}) end, Base_state = {db_state, Store, aarondb@index:new_index(), aarondb@index:new_aindex(), aarondb@index:new_avindex(), 0, [], maps:new(), maps:new(), [], Reactive_subject@1, [], Is_distributed, Ets_name, aarondb@raft:new([]), aarondb@vec_index:new(), maps:new(), aarondb@index@art:new(), maps:new(), maps:new(), maps:new(), [], maps:new(), maps:new(), {config, 1000, 1000, false, 10000}, []}, Res = begin _pipe = gleam@otp@actor:new(Base_state), _pipe@1 = gleam@otp@actor:on_message(_pipe, fun handle_message/2), gleam@otp@actor:start(_pipe@1) end, case Res of {ok, Started} -> Subj = erlang:element(3, Started), Reply = gleam@erlang@process:new_subject(), gleam@erlang@process:send(Subj, {boot, Ets_name, Store, Reply}), _ = gleam@erlang@process:'receive'(Reply, 600000), Pid = aarondb_process_ffi:subject_to_pid(Subj), _ = case Is_distributed of true -> nil; false -> _ = aarondb_global_ffi:register( <<"aarondb_leader"/utf8>>, Pid ), nil end, _ = proc_lib:spawn_link(fun() -> lifecycle_loop(Subj) end), {ok, Subj}; {error, E} -> {error, E} end. -file("src/aarondb/transactor.gleam", 86). -spec start_with_timeout(aarondb@storage:storage_adapter(), integer()) -> {ok, gleam@erlang@process:subject(message())} | {error, gleam@otp@actor:start_error()}. start_with_timeout(Store, _) -> do_start_named(Store, false, none). -file("src/aarondb/transactor.gleam", 66). -spec start(aarondb@storage:storage_adapter()) -> {ok, gleam@erlang@process:subject(message())} | {error, gleam@otp@actor:start_error()}. start(Store) -> start_with_timeout(Store, 1000). -file("src/aarondb/transactor.gleam", 72). -spec start_named(binary(), aarondb@storage:storage_adapter()) -> {ok, gleam@erlang@process:subject(message())} | {error, gleam@otp@actor:start_error()}. start_named(Name, Store) -> do_start_named(Store, false, {some, Name}). -file("src/aarondb/transactor.gleam", 79). -spec start_distributed(binary(), aarondb@storage:storage_adapter()) -> {ok, gleam@erlang@process:subject(message())} | {error, gleam@otp@actor:start_error()}. start_distributed(Name, Store) -> do_start_named(Store, true, {some, Name}). -file("src/aarondb/transactor.gleam", 169). -spec retract_entity( gleam@erlang@process:subject(message()), aarondb@fact:entity_id(), gleam@erlang@process:subject({ok, aarondb@shared@state:db_state()} | {error, binary()}) ) -> nil. retract_entity(Subj, Eid, Reply) -> gleam@erlang@process:send(Subj, {retract_entity, Eid, Reply}). -file("src/aarondb/transactor.gleam", 177). -spec log_query( gleam@erlang@process:subject(message()), aarondb@shared@state:query_context() ) -> nil. log_query(Subj, Ctx) -> Reply = gleam@erlang@process:new_subject(), gleam@erlang@process:send(Subj, {log_query, Ctx, Reply}), _ = gleam@erlang@process:'receive'(Reply, 100), nil. -file("src/aarondb/transactor.gleam", 188). -spec get_state(gleam@erlang@process:subject(message())) -> aarondb@shared@state:db_state(). get_state(Subj) -> Reply = gleam@erlang@process:new_subject(), gleam@erlang@process:send(Subj, {get_state, Reply}), State@1 = case gleam@erlang@process:'receive'(Reply, 5000) of {ok, State} -> State; _assert_fail -> erlang:error(#{gleam_error => let_assert, message => <<"Pattern match failed, no pattern matched the value."/utf8>>, file => <>, module => <<"aarondb/transactor"/utf8>>, function => <<"get_state"/utf8>>, line => 191, value => _assert_fail, start => 5401, 'end' => 5452, pattern_start => 5412, pattern_end => 5421}) end, State@1. -file("src/aarondb/transactor.gleam", 195). -spec set_schema( gleam@erlang@process:subject(message()), binary(), aarondb@fact:attribute_config() ) -> {ok, nil} | {error, binary()}. set_schema(Subj, Attr, Config) -> Reply = gleam@erlang@process:new_subject(), gleam@erlang@process:send(Subj, {set_schema, Attr, Config, Reply}), Res@1 = case gleam@erlang@process:'receive'(Reply, 5000) of {ok, Res} -> Res; _assert_fail -> erlang:error(#{gleam_error => let_assert, message => <<"Pattern match failed, no pattern matched the value."/utf8>>, file => <>, module => <<"aarondb/transactor"/utf8>>, function => <<"set_schema"/utf8>>, line => 202, value => _assert_fail, start => 5683, 'end' => 5732, pattern_start => 5694, pattern_end => 5701}) end, Res@1. -file("src/aarondb/transactor.gleam", 206). -spec set_schema_with_timeout( gleam@erlang@process:subject(message()), binary(), aarondb@fact:attribute_config(), integer() ) -> {ok, nil} | {error, binary()}. set_schema_with_timeout(Subj, Attr, Config, Timeout_ms) -> Reply = gleam@erlang@process:new_subject(), gleam@erlang@process:send(Subj, {set_schema, Attr, Config, Reply}), case gleam@erlang@process:'receive'(Reply, Timeout_ms) of {ok, Res} -> Res; {error, _} -> {error, <<"Timeout setting schema"/utf8>>} end. -file("src/aarondb/transactor.gleam", 220). -spec register_function( gleam@erlang@process:subject(message()), binary(), fun((aarondb@shared@state:db_state(), integer(), integer(), list(aarondb@fact:value())) -> list({aarondb@fact:eid(), binary(), aarondb@fact:value()})) ) -> nil. register_function(Subj, Name, Func) -> Reply = gleam@erlang@process:new_subject(), gleam@erlang@process:send(Subj, {register_function, Name, Func, Reply}), case gleam@erlang@process:'receive'(Reply, 5000) of {ok, nil} -> nil; _assert_fail -> erlang:error(#{gleam_error => let_assert, message => <<"Pattern match failed, no pattern matched the value."/utf8>>, file => <>, module => <<"aarondb/transactor"/utf8>>, function => <<"register_function"/utf8>>, line => 227, value => _assert_fail, start => 6332, 'end' => 6381, pattern_start => 6343, pattern_end => 6350}) end, nil. -file("src/aarondb/transactor.gleam", 231). -spec register_composite( gleam@erlang@process:subject(message()), list(binary()) ) -> {ok, nil} | {error, binary()}. register_composite(Subj, Attrs) -> Reply = gleam@erlang@process:new_subject(), gleam@erlang@process:send(Subj, {register_composite, Attrs, Reply}), Res@1 = case gleam@erlang@process:'receive'(Reply, 5000) of {ok, Res} -> Res; _assert_fail -> erlang:error(#{gleam_error => let_assert, message => <<"Pattern match failed, no pattern matched the value."/utf8>>, file => <>, module => <<"aarondb/transactor"/utf8>>, function => <<"register_composite"/utf8>>, line => 237, value => _assert_fail, start => 6594, 'end' => 6643, pattern_start => 6605, pattern_end => 6612}) end, Res@1. -file("src/aarondb/transactor.gleam", 241). -spec register_predicate( gleam@erlang@process:subject(message()), binary(), fun((aarondb@fact:value()) -> boolean()) ) -> nil. register_predicate(Subj, Name, Pred) -> Reply = gleam@erlang@process:new_subject(), gleam@erlang@process:send(Subj, {register_predicate, Name, Pred, Reply}), case gleam@erlang@process:'receive'(Reply, 5000) of {ok, nil} -> nil; _assert_fail -> erlang:error(#{gleam_error => let_assert, message => <<"Pattern match failed, no pattern matched the value."/utf8>>, file => <>, module => <<"aarondb/transactor"/utf8>>, function => <<"register_predicate"/utf8>>, line => 248, value => _assert_fail, start => 6870, 'end' => 6919, pattern_start => 6881, pattern_end => 6888}) end, nil. -file("src/aarondb/transactor.gleam", 252). -spec store_rule( gleam@erlang@process:subject(message()), aarondb@shared@ast:rule() ) -> {ok, nil} | {error, binary()}. store_rule(Subj, Rule) -> Reply = gleam@erlang@process:new_subject(), gleam@erlang@process:send(Subj, {store_rule, Rule, Reply}), Res@1 = case gleam@erlang@process:'receive'(Reply, 5000) of {ok, Res} -> Res; _assert_fail -> erlang:error(#{gleam_error => let_assert, message => <<"Pattern match failed, no pattern matched the value."/utf8>>, file => <>, module => <<"aarondb/transactor"/utf8>>, function => <<"store_rule"/utf8>>, line => 258, value => _assert_fail, start => 7110, 'end' => 7159, pattern_start => 7121, pattern_end => 7128}) end, Res@1. -file("src/aarondb/transactor.gleam", 262). -spec set_config( gleam@erlang@process:subject(message()), aarondb@shared@state:config() ) -> nil. set_config(Subj, Config) -> Reply = gleam@erlang@process:new_subject(), gleam@erlang@process:send(Subj, {set_config, Config, Reply}), case gleam@erlang@process:'receive'(Reply, 5000) of {ok, nil} -> nil; _assert_fail -> erlang:error(#{gleam_error => let_assert, message => <<"Pattern match failed, no pattern matched the value."/utf8>>, file => <>, module => <<"aarondb/transactor"/utf8>>, function => <<"set_config"/utf8>>, line => 265, value => _assert_fail, start => 7335, 'end' => 7384, pattern_start => 7346, pattern_end => 7353}) end, nil. -file("src/aarondb/transactor.gleam", 269). -spec transact( gleam@erlang@process:subject(message()), list({aarondb@fact:eid(), binary(), aarondb@fact:value()}) ) -> {ok, aarondb@shared@state:db_state()} | {error, binary()}. transact(Subj, Facts) -> Reply = gleam@erlang@process:new_subject(), gleam@erlang@process:send(Subj, {transact, Facts, none, Reply}), Res@1 = case gleam@erlang@process:'receive'(Reply, 5000) of {ok, Res} -> Res; _assert_fail -> erlang:error(#{gleam_error => let_assert, message => <<"Pattern match failed, no pattern matched the value."/utf8>>, file => <>, module => <<"aarondb/transactor"/utf8>>, function => <<"transact"/utf8>>, line => 275, value => _assert_fail, start => 7597, 'end' => 7646, pattern_start => 7608, pattern_end => 7615}) end, Res@1. -file("src/aarondb/transactor.gleam", 279). -spec transact_with_timeout( gleam@erlang@process:subject(message()), list({aarondb@fact:eid(), binary(), aarondb@fact:value()}), integer() ) -> {ok, aarondb@shared@state:db_state()} | {error, binary()}. transact_with_timeout(Subj, Facts, Timeout_ms) -> Reply = gleam@erlang@process:new_subject(), gleam@erlang@process:send(Subj, {transact, Facts, none, Reply}), case gleam@erlang@process:'receive'(Reply, Timeout_ms) of {ok, Res} -> Res; {error, _} -> {error, <<"Transaction timeout"/utf8>>} end. -file("src/aarondb/transactor.gleam", 292). -spec retract( gleam@erlang@process:subject(message()), list({aarondb@fact:eid(), binary(), aarondb@fact:value()}) ) -> {ok, aarondb@shared@state:db_state()} | {error, binary()}. retract(Subj, Facts) -> Reply = gleam@erlang@process:new_subject(), gleam@erlang@process:send(Subj, {retract, Facts, none, Reply}), Res@1 = case gleam@erlang@process:'receive'(Reply, 5000) of {ok, Res} -> Res; _assert_fail -> erlang:error(#{gleam_error => let_assert, message => <<"Pattern match failed, no pattern matched the value."/utf8>>, file => <>, module => <<"aarondb/transactor"/utf8>>, function => <<"retract"/utf8>>, line => 298, value => _assert_fail, start => 8205, 'end' => 8254, pattern_start => 8216, pattern_end => 8223}) end, Res@1. -file("src/aarondb/transactor.gleam", 498). -spec subscribe(gleam@erlang@process:subject(message())) -> gleam@erlang@process:subject(list(aarondb@fact:datom())). subscribe(Subj) -> Reply = gleam@erlang@process:new_subject(), gleam@erlang@process:send(Subj, {subscribe, Reply}), Reply.