%%%--------------------------------------------------------------------------- %%% Copyright (C) 2021, Skulup All Rights Reserved %%% Unauthorized copy of this file is through any medium strictly not allowed. %%% %%% Authors : Alpha Shaw %%% Created : 21 Apr 2021 by %%% Purpose : %%% %%% %%% Licensed under the Apache License, Version 2.0 (the "License"); %%% you may not use this file except in compliance with the License. %%% You may obtain a copy of the License at %%% %%% http://www.apache.org/licenses/LICENSE-2.0 %%% %%% Unless required by applicable law or agreed to in writing, software %%% distributed under the License is distributed on an "AS IS" BASIS, %%% WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. %%% See the License for the specific language governing permissions and %%% limitations under the License. %%%--------------------------------------------------------------------------- -module(mnesplit). -behaviour(gen_server). -include_lib("stdlib/include/ms_transform.hrl"). -include_lib("erlwater/include/logger.hrl"). -include("mnesplit.hrl"). -export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2, code_change/3]). -export([start_link/0, tracking_tables/0, silent_action_on/2]). -export([is_locally_inserted/2, is_locally_removed/2]). -export([check_inconsistencies/0, report_inconsistency/4]). -export([feed/2]). -type action() :: {'write_local', any()} | {'write_remote', any()} | {'delete_local', any()} | {'delete_remote', any()}. -export_type([action/0]). -record(state, {db, tables = sets:new(), exptimer, stitch_age}). -record(s0, {table, type, attributes, module, function, xargs, remote, modstate}). tracking_tables() -> gen_server:call(?MODULE, tracking_tables). silent_action_on(Action, Table) when is_function(Action); is_atom(Table) -> try gen_server:call(?MODULE, {untrack_table, Table}), Action() catch _:_ -> ok after gen_server:call(?MODULE, {track_table, Table}) end. start_link() -> case wait_mnesia(10) of ok -> gen_server:start_link({local, ?MODULE}, ?MODULE, [[]], []); _ -> {error, {mnesia, not_running}} end. init(_Args) -> {ok, _} = mnesia:subscribe(system), {ok, _} = mnesia:subscribe({table, schema, detailed}), Db = ets:new(?TABLE, [set, named_table]), State = lists:foldl(fun (schema, Acc) -> Acc; (T, Acc) -> Attrs = mnesia:table_info(T, all), case {should_track(T, Attrs), sets:is_element(T, Acc#state.tables)} of {true, true} -> % this table is already tracked, most likely fragment Acc; {true, _} -> track_table(T, Acc); {false, false} -> Acc end end, #state{tables = sets:new()}, mnesia:system_info(tables)), StitchAge = get_env(stitch_age, ?ETS_STITCH_MIN_AGE_MS) * 1000000, logger:info("~p (init): starting; stitch age: ~p ms, tracking the tables: ~p~n", [?MODULE, (StitchAge div 1000000), sets:to_list(State#state.tables)]), Expire = erlang:start_timer(?ETS_ITEM_PURGE_TIMEOUT, ?MODULE, expire), {ok, State#state{db = Db, exptimer = Expire, stitch_age = StitchAge}}. handle_call(tracking_tables, _From, #state{tables = Tables} = State) -> {reply, sets:to_list(Tables), State}; handle_call({track_table, Table}, _From, State) -> {reply, ok, track_table(Table, State)}; handle_call({untrack_table, Table}, _From, State) -> {reply, ok, untrack_table(Table, State)}; handle_call(_Any, _From, State) -> {reply, {error, badcall}, State}. handle_cast(_Any, State) -> {noreply, State}. handle_info({mnesia_table_event, {write, schema, {schema, schema, _Attrs}, _, _ActId}}, State) -> {noreply, State}; handle_info({mnesia_table_event, {write, schema, {schema, Table, Attrs}, _, _ActId}}, State) -> case {should_track(Table, Attrs), sets:is_element(Table, State#state.tables)} of {true, true} -> ?LOG_DEBUG("mnesplit(write, schema): ~p is already tracked", [Table]), {noreply, State}; {true, false} -> ?LOG_DEBUG("mnesplit(write, schema): calling track_table(~p)", [Table]), {noreply, track_table(Table, State)}; {false, true} -> ?LOG_DEBUG("mnesplit(write, schema): calling untrack_table(~p)", [Table]), {noreply, untrack_table(Table, State)}; {false, false} -> ?LOG_DEBUG("mnesplit(write, schema): ~p is not tracked", [Table]), {noreply, State} end; handle_info({mnesia_table_event, {delete, schema, {schema, Table, _Attrs}, _, _ActId}}, State) -> case sets:is_element(Table, State#state.tables) of true -> {noreply, untrack_table(Table, State)}; false -> {noreply, State} end; handle_info({mnesia_table_event, {write, Table, Record, [], _ActId}}, State) -> case sets:is_element(Table, State#state.tables) of false -> ?LOG_DEBUG("mnesplit(write new): table ~p is not tracked~n", [Table]); true -> ?LOG_DEBUG("mnesplit(store): storing {~p, ~p}~n", [Table, element(2, Record)]), ets:insert(?TABLE, {{Table, element(2, Record)}, 'insert', erlang:system_time()}) end, {noreply, State}; handle_info({mnesia_table_event, {write, _Table, _Record, _NonEmptyList, _Act}}, State) -> % this is update of already existing key, may ignore ?LOG_DEBUG("mnesplit(update, ignored): table ~p, ~p~n", [_Table, _Record]), {noreply, State}; handle_info({mnesia_table_event, {delete, Table, {Table, Key}, _Value, _ActId}}, State) -> case sets:is_element(Table, State#state.tables) of false -> ?LOG_DEBUG("mnesplit(delete): table ~p is not tracked~n", [Table]); true -> ?LOG_DEBUG("mnesplit(delete): table: ~p, key: ~p~n", [Table, Key]), ets:insert(?TABLE, {{Table, Key}, 'delete', erlang:system_time()}) end, {noreply, State}; handle_info({mnesia_table_event, {delete, Table, Record, _Old, ActId}}, State) -> Key = element(2, Record), handle_info({mnesia_table_event, {delete, Table, {Table, Key}, Record, ActId}}, State); handle_info({mnesia_system_event, {mnesia_up, Node}}, State) -> logger:info("~p: got mnesia_up at ~p", [?MODULE, Node]), {noreply, State}; handle_info({mnesia_system_event, {mnesia_down, Node}}, State) when node() == Node -> logger:info("~p: got mnesia_down for local node, stop", [?MODULE]), {stop, normal, State}; handle_info({mnesia_system_event, {mnesia_down, Node}}, State) -> logger:info("~p: got mnesia_down at ~p", [?MODULE, Node]), {noreply, State}; handle_info({mnesia_system_event, {inconsistent_database, running_partitioned_network, Node}}, State) -> logger:info("~p: Inconsistency (running_partitioned_network) " "with ~p~n", [?MODULE, Node]), case application:get_env(?MODULE, delay, 0) of 0 -> ok; Value -> ?LOG_DEBUG("mnesplit: sleeping ~p before acquiring lock", [Value]), timer:sleep(Value) end, global:trans({?LOCK, self()}, fun() -> ?LOG_DEBUG("~p: have global lock. mnesia locks: ~p", [?MODULE, mnesia_locker:get_held_locks()]), ?LOG_DEBUG("~p: nodes: ~p,~n running: ~p,~n ~p messages: ~p~n", [?MODULE, mnesia:system_info(db_nodes), mnesia:system_info(running_db_nodes), process_info(self(), message_queue_len), process_info(self(), messages)]), stitch_together(Node) end), {noreply, State}; handle_info({mnesia_system_event, {inconsistent_database, starting_partitioned_network, Node}}, State) -> % this is recovery message sent after merge. logger:info("~p: starting_partitioned_network with ~p", [?MODULE, Node]), {noreply, State}; handle_info({mnesia_system_event, {inconsistent_database, Context, Node}}, State) -> logger:info("~p: mnesia inconsistent_database in ~p with ~p", [?MODULE, Context, Node]), {noreply, State}; handle_info({timeout, Ref, expire}, #state{exptimer = Ref, stitch_age = StitchAge} = State) -> Now = erlang:system_time(), ?LOG_DEBUG("Will exipre all items inserted before: ~p ns", [Now - StitchAge]), Items = ets:select(?TABLE, ets:fun2ms( fun({_, _, Ts} = I) when (Ts + StitchAge) < Now -> I end )), lists:foreach( fun(I) -> ets:delete(?TABLE, element(1, I)) end, Items), case Items of [] -> ok; _ -> error_logger:warning_msg("~p: deletes ~p cached table stiches.", [?MODULE, length(Items)]) end, ETimer = erlang:start_timer(?ETS_ITEM_PURGE_TIMEOUT, self(), expire), {noreply, State#state{exptimer = ETimer}}; handle_info({timeout, _, {subscribe, T}}, #state{} = State) -> case {should_track(T), sets:is_element(T, State#state.tables)} of {false, false} -> {noreply, State}; {true, false} -> {noreply, track_table(T, State)}; {true, true} -> {noreply, State} end; handle_info({mnesia_system_event, {mnesia_info, _, _}} = _E, State) -> logger:info("unhandled mnesia event ~p", [_E]), {noreply, State}; handle_info(Any, State) -> logger:info("~p: unhandled info ~p~n", [?MODULE, Any]), {noreply, State}. terminate(Reason, _State) -> logger:info("~p: terminating (~p)", [?MODULE, Reason]), ok. code_change(_Old, State, _Extra) -> {ok, State}. track_table(Table, State) -> Pre = erlang:system_time(microsecond), case mnesia:wait_for_tables([Table], ?WAIT_TIMEOUT) of ok -> Post = erlang:system_time(microsecond), case mnesia:subscribe({table, Table, detailed}) of {ok, _} -> logger:info("~p: started tracking ~p, " "elapsed: ~p microsecs", [?MODULE, Table, (Post - Pre)]), Ns = sets:add_element(Table, State#state.tables), track_fragments(Table, State#state{tables = Ns}); {error, {not_active_local, Table}} -> erlang:start_timer(?RESUBSCRIBE_TIMEOUT, self(), {subscribe, Table}), State; {error, {no_exists, Table}} -> logger:info("~p: table ~p disappeared while " "subscribe", [?MODULE, Table]), State; {error, Other} -> logger:info("~p: error subscribing ~p: ~p", [?MODULE, Table, Other]), State end; {timeout, [Table]} -> erlang:start_timer(?RESUBSCRIBE_TIMEOUT, ?MODULE, {subscribe, Table}), State end. track_fragments(Table, State) when is_atom(Table) -> case mnesia:activity(async_dirty, fun() -> mnesia:table_info(Table, frag_names) end, [], mnesia_frag) of [Table] -> State; [Table | Frags] -> track_fragments(Frags, State) end; track_fragments([], State) -> State; track_fragments([Fragment | Next], State) -> case {should_track(Fragment), sets:is_element(Fragment, State#state.tables)} of {true, true} -> track_fragments(Next, State); {true, false} -> track_fragments(Next, track_table(Fragment, State)); {false, false} -> track_fragments(Next, State) end. untrack_fragments(Table, State) when is_atom(Table) -> case mnesia:activity(async_dirty, fun() -> mnesia:table_info(Table, frag_names) end, [], mnesia_frag) of [Table] -> State; [Table | Frags] -> untrack_fragments(Frags, State) end; untrack_fragments([], State) -> State; untrack_fragments([Fragment | Next], State) -> case {should_track(Fragment), sets:is_element(Fragment, State#state.tables)} of {false, false} -> untrack_fragments(Next, State); {false, true} -> untrack_fragments(Next, untrack_table(Fragment, State)); {true, true} -> untrack_fragments(Next, State) end. untrack_table(Table, State) -> logger:info("~p: stop tracking ~p", [?MODULE, Table]), mnesia:unsubscribe({table, Table, detailed}), Ns = sets:del_element(Table, State#state.tables), untrack_fragments(Table, State#state{tables = Ns}). should_track(T) -> try mnesia:table_info(T, all) of Attrs -> should_track(T, Attrs) catch _:_ -> false end. should_track(T, Attr) -> LocalContent = proplists:get_value(local_content, Attr), Type = proplists:get_value(type, Attr), AllNodes = proplists:get_value(disc_copies, Attr, []) ++ proplists:get_value(ram_copies, Attr, []) ++ proplists:get_value(disc_only_copies, Attr, []), Member = lists:member(node(), AllNodes), Compare = case proplists:get_value(frag_properties, Attr, []) of [] -> proplists:get_value(mnesplit_compare, proplists:get_value(user_properties, Attr, [])); Frag -> case proplists:get_value(base_table, Frag) of T -> proplists:get_value(mnesplit_compare, proplists:get_value(user_properties, Attr, [])); DiffT -> get_method(DiffT, default_method()) end end, if LocalContent == true -> ?LOG_DEBUG("mnesplit: should not track ~p: local_content only~n", [T]), false; Member == false -> ?LOG_DEBUG("mnesplit: should not track ~p: this node does not have " "a copy", [T]), false; Type == bag -> % sets and ordered_sets are ok ?LOG_DEBUG("mnesplit: should not track ~p: bag~n", [T]), false; length(AllNodes) == 1 -> ?LOG_DEBUG("mnesplit: should not track ~p: all_nodes ~p~n", [T, AllNodes]), false; Compare == ignore -> ?LOG_DEBUG("mnesplit: should not track ~p: mnesplit_compare is ignore~n", [T]), false; true -> ?LOG_DEBUG("mnesplit: should track ~p~n", [T]), true end. stitch_together(Node) -> case rpc:call(Node, mnesia, system_info, [is_running]) of yes -> do_stitch_together(Node); Other -> logger:info("~p: node ~p: mnesia not running (~p), not " "stitching~n", [?MODULE, Node, Other]), ok end. do_stitch_together(Node) -> IslandB = case rpc:call(Node, mnesia, system_info, [running_db_nodes]) of {badrpc, Reason} -> logger:info("~p: unable to ask mnesia:system_info(" "running_db_nodes) on ~p: ~p", [?MODULE, Node, Reason]), []; Answer -> Answer end, TabsAndNodes = affected_tables(IslandB), Tables = [T || {T, _} <- TabsAndNodes], TabMethods = [{T, Ns, get_method(T, default_method())} || {T, Ns} <- TabsAndNodes], logger:info("~p: will attempts to stitch the tables: ~p at :~p", [?MODULE, Tables, Node]), stitch_tabs(TabMethods, Node). stitch_tabs(TabMethods, Node) -> [do_stitch(TM, Node) || TM <- TabMethods]. do_stitch({Tab, _Nodes, {M, F, Xargs}}, Node) -> Type = case mnesia:table_info(Tab, type) of ordered_set -> set; S -> S end, Attrs = mnesia:table_info(Tab, attributes), try M:F(init, {Tab, Type, Attrs, Xargs}, Node) of {ok, Ms} -> S0 = #s0{module = M, function = F, xargs = Xargs, table = Tab, type = Type, attributes = Attrs, remote = Node, modstate = Ms}, logger:info("~p: starting table ~p with ~p", [?MODULE, Tab, Node]), case rpc:call(Node, ?MODULE, feed, [Tab, self()]) of {ok, Ref} -> try run_feed(S0, Ref) of ok -> logger:info("~p: finished table ~p " "(feed mode)", [?MODULE, Tab]), ok catch throw:?DONE -> ok; Error:Code -> logger:info("~p: exception ~p:~p " "merging ~p with ~p (feed)", [?MODULE, Error, Code, Tab, Node]), ok end; {badrpc, _} -> try run_stitch(S0) of ok -> logger:info("~p: finished table ~p " "(key-by-key mode)", [?MODULE, Tab]), ok catch throw:?DONE -> ok; Error:Code -> logger:info("~p: exception ~p:~p " "merging ~p with ~p (key-by-key)", [?MODULE, Error, Code, Tab, Node]), ok end end; Other -> logger:error("~p: unexpected answer ~p on init state " "~p:~p(init, {~p, ~p, ~p, ~p}, ~p)", [?MODULE, Other, M, F, Tab, Type, Attrs, Xargs, Node]), ok catch Error:Code -> logger:info("~p: exception ~p:~p on init state " "~p:~p(init, {~p, ~p, ~p, ~p}, ~p)", [?MODULE, Error, Code, M, F, Tab, Type, Attrs, Xargs, Node]), ok end; do_stitch({Tab, _Nodes, ignore}, _Node) -> logger:info("~p: ignoring table ~p (configuration)", [?MODULE, Tab]), ok. run_feed(#s0{table = Tab} = S0, Ref) -> LocalKeys = sets:from_list(mnesia:dirty_all_keys(Tab)), run_feed(S0, Ref, LocalKeys). run_feed(#s0{module = M, function = F, table = Tab, remote = Remote, type = Type, modstate = MSt} = S0, Ref, LocalKeys) -> receive {mnesplit_feed, Tab, Key, B, Inserted} -> A = mnesia:dirty_read({Tab, Key}), case {A, B} of {[], []} -> run_feed(S0, Ref, sets:del_element(Key, LocalKeys)); {A, A} -> run_feed(S0, Ref, sets:del_element(Key, LocalKeys)); {[], [Bb]} when Type == set, Inserted == false -> case is_locally_removed(Tab, Key) of {true, _} -> ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p ~p): delete_remote " "(locally removed)", [Tab, Key, A, B]), delete(Remote, Bb), run_feed(S0, Ref, sets:del_element(Key, LocalKeys)); false -> ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p ~p): write_local " "(not locally removed)", [Tab, Key, A, B]), write(Bb), run_feed(S0, Ref, sets:del_element(Key, LocalKeys)) end; {[], [Bb]} when Type == set -> {true, RemoteTime} = Inserted, case is_locally_removed(Tab, Key) of false -> ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p ~p): write_local", [Tab, Key, A, B]), write(Bb); {true, LocalTime} when LocalTime < RemoteTime -> ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p ~p): write_local " "(remote timer won)", [Tab, Key, A, B]), write(Bb); {true, _} -> ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p ~p): delete_remote " "(local timer won)", [Tab, Key, A, B]), delete(Remote, Bb) end, run_feed(S0, Ref, sets:del_element(Key, LocalKeys)); {A, B} -> try M:F(A, B, MSt) of {ok, Actions, Snext} -> ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p ~p): actions ~p", [Tab, Key, A, B, Actions]), do_actions(Actions, Remote), run_feed(S0#s0{modstate = Snext}, Ref, sets:del_element(Key, LocalKeys)); {inconsistency, Error, Snext} -> ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p ~p): inconsistency" " ~p", [Tab, Key, A, B, Error]), report_inconsistency(Remote, Tab, Key, Error), run_feed(S0#s0{modstate = Snext}, Ref, sets:del_element(Key, LocalKeys)); Other -> logger:error("~p(stitch ~p ~p ~p ~p): " "unexpected result ~p, ignoring", [?MODULE, Tab, Key, A, B, Other]), run_feed(S0, Ref, sets:del_element(Key, LocalKeys)) catch Error:Code:Stacktrace -> logger:error("~p(stitch ~p ~p ~p ~p): " "caught ~p:~p in ~p", [?MODULE, Tab, Key, A, B, Error, Code, Stacktrace]), run_feed(S0, Ref, sets:del_element(Key, LocalKeys)) end end; {mnesplit_feed, Ref, eof} -> end_feed(S0, LocalKeys) after 30000 -> logger:error("~p: timeout waiting for key or eof from ~p", [?MODULE, Remote]), ok end. end_feed(#s0{module = M, function = F, table = Tab, remote = Remote, modstate = MSt, type = Type}, LocalKeys) -> lists:foldl(fun(K, Sx) -> A = mnesia:dirty_read({Tab, K}), case A of [] -> Sx; % key was deleted locally too [Aa] when Type == set -> case is_locally_inserted(Tab, K) of {true, LocalTime} -> case rpc:call(Remote, ?MODULE, is_locally_removed, [Tab, K]) of false -> ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p []): " "write_remote", [Tab, K, A]), write(Remote, Aa); {true, RemoteTime} when LocalTime < RemoteTime -> ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p []): " "delete_local (timer won)", [Tab, K, A]), delete(Aa); _ -> ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p []): " "write_remote", [Tab, K, A]), write(Remote, Aa) end; false -> ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p, []): delete_local " "(not inserted locally)", [Tab, K, A]), delete(Aa) end, Sx; A -> try M:F(A, [], Sx) of {ok, Actions, Snext} -> ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p, []): actions ~p", [Tab, K, A, Actions]), do_actions(Actions, Remote), Snext; {inconsistency, Error, Snext} -> ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p): inconsistency" " ~p", [Tab, K, A, Error]), report_inconsistency(Remote, Tab, K, Error), Snext; Other -> logger:error("~p(stitch ~p ~p ~p []): " "unexpected result ~p, ignoring", [?MODULE, Tab, K, A, Other]), Sx catch Error:Code:Stacktrace -> logger:error("~p(stitch ~p ~p ~p []): " "caught ~p:~p in ~p", [?MODULE, Tab, K, A, Error, Code, Stacktrace]), Sx end end end, MSt, sets:to_list(LocalKeys)), try M:F(done, MSt, Remote) of _ -> ok catch Error:Code -> logger:info("~p: caught ~p:~p finalizing " "~p:~p(done, ~p, ~p), ignored", [?MODULE, Error, Code, M, F, MSt, Remote]), ok end, ok. run_stitch(#s0{module = M, function = F, table = Tab, remote = Remote, type = Type, modstate = MSt}) -> LocalKeys = mnesia:dirty_all_keys(Tab), % usort used to remove duplicates Keys = lists:usort(lists:concat([LocalKeys, remote_keys(Remote, Tab)])), lists:foldl(fun(K, Sx) -> A = mnesia:dirty_read({Tab, K}), B = remote_object(Remote, Tab, K), case {A, B} of {[], []} -> % element is not present anymore Sx; {Aa, Aa} -> % elements are the same ?LOG_DEBUG("mnesplit(~p, ~p): same element ~p", [Tab, K, Aa]), Sx; {[Aa], []} when Type == set -> % remote element is not present case is_locally_inserted(Tab, K) of {true, LocalTime} -> case rpc:call(Remote, mnesplit, is_locally_removed, [Tab, K]) of false -> ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p ~p): write " "remote (not removed)~n", [Tab, Type, A, B]), write(Remote, Aa); {true, RemoteTime} when LocalTime < RemoteTime -> ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p ~p): delete " "local (remote timer won)~n", [Tab, Type, A, B]), delete(Aa); _ -> % can be an error too ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p ~p): write " "remote (error or timer)~n", [Tab, Type, A, B]), write(Remote, Aa) end; false -> ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p ~p): delete local~n", [Tab, Type, A, B]), delete(Aa) end, Sx; {[], [Bb]} when Type == set -> % local element not present case is_locally_removed(Tab, K) of {true, LocalTime} -> case rpc:call(Remote, ?MODULE, is_locally_inserted, [Tab, K]) of false -> ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p ~p): delete " "remote (not remotely inserted)~n", [Tab, Type, A, B]), delete(Remote, Bb); {true, RemoteTime} when LocalTime < RemoteTime -> ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p ~p): write " "local (remote timer won)~n", [Tab, Type, A, B]), write(Bb); _ -> ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p ~p): delete " "remote (error or timer)~n", [Tab, Type, A, B]), delete(Remote, Bb) end; false -> ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p ~p): write local~n", [Tab, Type, A, B]), write(Bb) end, Sx; {A, B} -> Sn = try M:F(A, B, Sx) of {ok, Actions, Sr} -> ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p ~p): ~p" "both~n", [Tab, Type, A, B]), do_actions(Actions, Remote), Sr; {inconsistency, Error, Sr} -> ?LOG_DEBUG("mnesplit(stitch ~p ~p ~p ~p): inconsistency " "~p~n", [Tab, Type, A, B, Error]), report_inconsistency(Remote, Tab, K, Error), Sr; Other -> logger:info("~p: ~p:~p(~p, ~p, ~p): bad " "return value ~p, ignoring", [?MODULE, M, F, A, B, Sx, Other]), Sx catch Error:Code -> logger:info("~p: ~p:~p(~p, ~p, ~p): caught " "~p:~p, ignoring", [?MODULE, M, F, A, B, Sx, Error, Code]), Sx end, Sn end end, MSt, Keys), try M:F(done, MSt, Remote) of _ -> ok catch Error:Code -> logger:info("~p: caught ~p:~p finalizing " "~p:~p(done, ~p, ~p), ignored", [?MODULE, Error, Code, M, F, MSt, Remote]), ok end, ok. do_actions([], _) -> ok; do_actions([{write_local, Ae} | Next], Remote) -> write(Ae), do_actions(Next, Remote); do_actions({write_local, Ae}, _Remote) -> write(Ae); do_actions([{delete_local, Ae} | Next], Remote) -> delete(Ae), do_actions(Next, Remote); do_actions({delete_local, Ae}, _Remote) -> delete(Ae); do_actions([{write_remote, Ae} | Next], Remote) -> write(Remote, Ae), do_actions(Next, Remote); do_actions({write_remote, Ae}, Remote) -> write(Remote, Ae); do_actions([{delete_remote, Ae} | Next], Remote) -> delete(Remote, Ae), do_actions(Next, Remote); do_actions({delete_remote, Ae}, Remote) -> delete(Remote, Ae); do_actions([A | Next], Remote) -> logger:error("~p: invalid action ~p merging with ~p", [?MODULE, A, Remote]), do_actions(Next, Remote); do_actions(A, Remote) -> logger:error("~p: invalid action ~p merging with ~p", [?MODULE, A, Remote]). affected_tables(IslandB) -> IslandA = mnesia:system_info(running_db_nodes), Tables = mnesia:system_info(tables) -- [schema], lists:foldl(fun(T, Acc) -> Nodes = mnesia:table_info(T, all_nodes), Attrs = mnesia:table_info(T, all), Should = should_track(T, Attrs), case {intersection(IslandA, Nodes), intersection(IslandB, Nodes)} of {[_ | _], [_ | _]} when Should == true -> [{T, Nodes} | Acc]; _ -> Acc end end, [], Tables). write(Remote, A) -> rpc:call(Remote, mnesia, dirty_write, [A]). write(A) when is_list(A) -> lists:foreach(fun(E) -> mnesia:dirty_write(E) end, A); write(A) -> mnesia:dirty_write(A). delete(Remote, A) -> rpc:call(Remote, mnesia, dirty_delete_object, [A]). delete(A) when is_list(A) -> lists:foreach(fun(E) -> mnesia:dirty_delete_object(E) end, A); delete(A) -> mnesia:dirty_delete_object(A). remote_keys(Remote, Tab) -> case rpc:call(Remote, mnesia, dirty_all_keys, [Tab]) of {badrpc, {'EXIT', {aborted, {no_exists, [Tab | _]}}}} -> logger:error("~p: tab ~p does not exists on ~p", [?MODULE, Tab, Remote]), throw(?DONE); {badrpc, Reason} -> logger:error("~p: error querying dirty_all_keys(~p) " "on ~p:~n ~p", [?MODULE, Tab, Remote, Reason]), mnesia:abort({badrpc, Remote, Reason}); Keys -> Keys end. remote_object(Remote, Tab, Key) -> case rpc:call(Remote, mnesia, dirty_read, [Tab, Key]) of {badrpc, Reason} -> logger:error("?p: error fetching {~p, ~p} on ~p:~n ~p", [?MODULE, Tab, Key, Remote, Reason]), mnesia:abort({badrpc, Remote, Reason}); Object -> Object end. is_locally_inserted(Tab, Key) -> case ets:lookup(?TABLE, {Tab, Key}) of [] -> ?LOG_DEBUG("mnesplit(is_locally_inserted): {~p, ~p} not in ets~n", [Tab, Key]), false; [{{Tab, Key}, 'insert', T}] -> ?LOG_DEBUG("mnesplit(is_locally_inserted): {~p, ~p} was inserted~n", [Tab, Key]), {true, T}; [{{Tab, Key}, 'delete', _}] -> ?LOG_DEBUG("mnesplit(is_locally_inserted): {~p, ~p} was deleted~n", [Tab, Key]), false end. is_locally_removed(Tab, Key) -> case ets:lookup(?TABLE, {Tab, Key}) of [] -> ?LOG_DEBUG("mnesplit(is_locally_removed): {~p, ~p} not in ets~n", [Tab, Key]), false; [{{Tab, Key}, 'insert', _}] -> ?LOG_DEBUG("mnesplit(is_locally_removed): {~p, ~p} was inserted~n", [Tab, Key]), false; [{{Tab, Key}, 'delete', T}] -> ?LOG_DEBUG("mnesplit(is_locally_removed): {~p, ~p} was deleted~n", [Tab, Key]), {true, T} end. feed(Tab, Pid) -> R = erlang:make_ref(), spawn(fun() -> lists:foreach(fun(K) -> Pid ! {mnesplit_feed, Tab, K, mnesia:dirty_read(Tab, K), is_locally_inserted(Tab, K)} end, mnesia:dirty_all_keys(Tab)), Pid ! {mnesplit_feed, R, eof} end), {ok, R}. check_inconsistencies() -> lists:foreach(fun({{?MODULE, inconsistency, Remote, Tab, Key}, _}) -> case lists:member(Remote, nodes()) of false -> ok; true -> case {mnesia:dirty_read(Tab, Key), rpc:call(Remote, mnesia, dirty_read, [Tab, Key])} of {A, A} -> alarm_handler:clear_alarm({?MODULE, inconsistency, Remote, Tab, Key}); _ -> ok end end; (_) -> ok end, alarm_handler:get_alarms()). report_inconsistency(Remote, Tab, Key, Error) -> alarm_handler:set_alarm({{?MODULE, inconsistency, Remote, Tab, Key}, Error}). default_method() -> get_env(default_method, ?DEFAULT_METHOD). get_method(Table, Default) -> try mnesia:read_table_property(Table, mnesplit_compare) of {mnesplit_compare, Method} -> Method catch exit:_ -> Default end. get_env(Env, Default) -> case application:get_env(Env) of undefined -> Default; {ok, undefined} -> Default; {ok, Value} -> Value end. intersection(A, B) -> A -- (A -- B). wait_mnesia(N) -> case mnesia:system_info(is_running) of yes -> ok; _ -> case N > 0 of true -> timer:sleep(100), wait_mnesia(N - 1); false -> {error, not_running} end end.