Current section
Files
Jump to
Current section
Files
src/mnesplit.erl
%%%---------------------------------------------------------------------------
%%% Copyright (C) 2021, Skulup All Rights Reserved
%%% Unauthorized copy of this file is through any medium strictly not allowed.
%%%
%%% Authors : Alpha Shaw <shawalpha5@gmail.com>
%%% Created : 21 Apr 2021 by <shawalpha5@gmail.com>
%%% 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.