Skip to content

Commit 3383976

Browse files
committed
Fix mnesia:force_load_table/1 getting stuck
1 parent fffa95e commit 3383976

3 files changed

Lines changed: 233 additions & 86 deletions

File tree

lib/mnesia/src/mnesia_controller.erl

Lines changed: 62 additions & 58 deletions
Original file line numberDiff line numberDiff line change
@@ -1184,71 +1184,75 @@ handle_info(#dumper_done{worker_pid=Pid, worker_res=Res}, State) ->
11841184
{stop, fatal, State}
11851185
end;
11861186

1187-
handle_info(Done = #loader_done{worker_pid=WPid, table_name=Tab}, State0) ->
1187+
handle_info(Done = #loader_done{worker_pid=WPid, table_name=Tab, reply=Reply}, State0) ->
11881188
LateQueue0 = State0#state.late_loader_queue,
1189-
State1 = State0#state{loader_pid = lists:keydelete(WPid,1,get_loaders(State0))},
1189+
{value, {WPid, Worker}, Loaders} = lists:keytake(WPid, 1, get_loaders(State0)),
1190+
State1 = State0#state{loader_pid = Loaders},
11901191

11911192
State2 =
1192-
case Done#loader_done.is_loaded of
1193-
true ->
1194-
%% Optional table announcement
1195-
if
1196-
Done#loader_done.needs_announce == true,
1197-
Done#loader_done.needs_reply == true ->
1198-
i_have_tab(Tab),
1199-
%% Should be {dumper,{add_table_copy, _}} only
1200-
reply(Done#loader_done.reply_to,
1201-
Done#loader_done.reply);
1202-
Done#loader_done.needs_reply == true ->
1203-
%% Should be {dumper,{add_table_copy,_}} only
1204-
reply(Done#loader_done.reply_to,
1205-
Done#loader_done.reply);
1206-
Done#loader_done.needs_announce == true, Tab == schema ->
1207-
i_have_tab(Tab);
1208-
Done#loader_done.needs_announce == true ->
1209-
i_have_tab(Tab),
1210-
%% Local node needs to perform user_sync_tab/1
1211-
Ns = val({current, db_nodes}),
1212-
abcast(Ns, {i_have_tab, Tab, node()});
1213-
Tab == schema ->
1214-
ignore;
1215-
true ->
1216-
%% Local node needs to perform user_sync_tab/1
1217-
Ns = val({current, db_nodes}),
1218-
AlreadyKnows = val({Tab, active_replicas}),
1219-
abcast(Ns -- AlreadyKnows, {i_have_tab, Tab, node()})
1220-
end,
1221-
%% Optional user sync
1222-
case Done#loader_done.needs_sync of
1223-
true -> user_sync_tab(Tab);
1224-
false -> ignore
1225-
end,
1226-
State1#state{late_loader_queue=gb_trees:delete_any(Tab, LateQueue0)};
1227-
false ->
1228-
%% Either the node went down or table was not
1229-
%% loaded remotly yet
1230-
case Done#loader_done.needs_reply of
1231-
true ->
1232-
reply(Done#loader_done.reply_to,
1233-
Done#loader_done.reply);
1234-
false ->
1235-
ignore
1236-
end,
1237-
1238-
case {?catch_val({Tab, storage_type}), val({Tab, active_replicas})} of
1239-
{unknown, _} -> %% Should not have a local copy anymore
1193+
case Done#loader_done.is_loaded of
1194+
true ->
1195+
%% Optional table announcement
1196+
if
1197+
Done#loader_done.needs_announce == true,
1198+
Done#loader_done.needs_reply == true ->
1199+
i_have_tab(Tab),
1200+
%% Should be {dumper,{add_table_copy, _}} only
1201+
reply(Done#loader_done.reply_to, Reply);
1202+
Done#loader_done.needs_reply == true ->
1203+
%% Should be {dumper,{add_table_copy,_}} only
1204+
reply(Done#loader_done.reply_to, Reply);
1205+
Done#loader_done.needs_announce == true, Tab == schema ->
1206+
i_have_tab(Tab);
1207+
Done#loader_done.needs_announce == true ->
1208+
i_have_tab(Tab),
1209+
%% Local node needs to perform user_sync_tab/1
1210+
Ns = val({current, db_nodes}),
1211+
abcast(Ns, {i_have_tab, Tab, node()});
1212+
Tab == schema ->
1213+
ignore;
1214+
true ->
1215+
%% Local node needs to perform user_sync_tab/1
1216+
Ns = val({current, db_nodes}),
1217+
AlreadyKnows = val({Tab, active_replicas}),
1218+
abcast(Ns -- AlreadyKnows, {i_have_tab, Tab, node()})
1219+
end,
1220+
%% Optional user sync
1221+
case Done#loader_done.needs_sync of
1222+
true -> user_sync_tab(Tab);
1223+
false -> ignore
1224+
end,
1225+
State1#state{late_loader_queue=gb_trees:delete_any(Tab, LateQueue0)};
1226+
false ->
1227+
%% Either the node went down or table was not
1228+
%% loaded remotly yet
1229+
case Done#loader_done.needs_reply of
1230+
true ->
1231+
reply(Done#loader_done.reply_to, Reply);
1232+
false ->
1233+
ignore
1234+
end,
1235+
1236+
case {?catch_val({Tab, storage_type}), val({Tab, active_replicas}), ?catch_val({Tab, load_by_force})} of
1237+
{unknown, _, _} -> %% Should not have a local copy anymore
12401238
State1#state{late_loader_queue=gb_trees:delete_any(Tab, LateQueue0)};
1241-
{_, [_|_]} -> % still available elsewhere
1242-
{value,{_,Worker}} = lists:keysearch(WPid,1,get_loaders(State0)),
1243-
add_loader(Tab,Worker,State1);
1244-
{ram_copies, []} ->
1245-
DelState = State1#state{late_loader_queue=gb_trees:delete_any(Tab, LateQueue0)},
1239+
{_, [_|_], _} -> % still available elsewhere
1240+
add_loader(Tab,Worker,State1);
1241+
{ram_copies, [], _} ->
1242+
DelState = State1#state{late_loader_queue=gb_trees:delete_any(Tab, LateQueue0)},
12461243
cast({disc_load, Tab, ram_only}),
12471244
DelState;
1248-
{_, []} -> %% Table deleted or not loaded anywhere
1245+
{_, [], true} when is_record(Worker, net_load), Reply =:= {not_loaded, none_active} ->
1246+
%% Network load could have been aborted by Sender node going down,
1247+
%% if user has forced the load, we retry loading from disc
1248+
DelState = State1#state{late_loader_queue=gb_trees:delete_any(Tab, LateQueue0)},
1249+
cast({disc_load, Tab, forced_by_user}),
1250+
DelState;
1251+
{_, [], _} ->
1252+
%% Table deleted or not loaded anywhere
12491253
State1#state{late_loader_queue=gb_trees:delete_any(Tab, LateQueue0)}
1250-
end
1251-
end,
1254+
end
1255+
end,
12521256
State3 = opt_start_worker(State2),
12531257
noreply(State3);
12541258

lib/mnesia/src/mnesia_loader.erl

Lines changed: 83 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,7 @@
3333
net_load_table/4,
3434
send_table/4]).
3535

36-
-export([spawned_receiver/8]). %% Spawned lock taking process
36+
-export([spawned_receiver/7]). %% Spawned lock taking process
3737

3838
-import(mnesia_lib, [set/2, fatal/2, verbose/2, dbg_out/2]).
3939

@@ -281,9 +281,8 @@ init_receiver(Node, Tab,Storage,Cs,Reason) ->
281281
true = lists:member(Node, Active),
282282
{SenderPid, TabSize, DetsData} =
283283
start_remote_sender(Node,Tab,Storage,load),
284-
Init = table_init_fun(SenderPid, Storage),
285-
Args = [self(),Tab,Storage,Cs,SenderPid,
286-
TabSize,DetsData,Init],
284+
Args = [self(),Tab,Storage,Cs,SenderPid,
285+
TabSize,DetsData],
287286
Pid = spawn_link(?MODULE, spawned_receiver, Args),
288287
put(mnesia_real_loader, Pid),
289288
wait_on_load_complete(Pid)
@@ -339,7 +338,18 @@ start_remote_sender(Node,Tab,Storage, Why) ->
339338
end.
340339

341340
table_init_fun(SenderPid, Storage) ->
341+
Parent = self(),
342342
fun(read) ->
343+
case Storage of
344+
disc_only_copies ->
345+
%% For disc_only_copies this is running inside dets_server process,
346+
%% Parent is the process which is registered as receiver, and must
347+
%% have our pid in order to forward us any {copier_done, Node} sent to it.
348+
put(mnesia_table_dets_receiver, {Parent, SenderPid}),
349+
Parent ! {dets_receiver_pid, self()};
350+
_ ->
351+
ok
352+
end,
343353
% We want to store subscribed mnesia table events received during
344354
% table copying for later processing to not let receiver message queue
345355
% to grow too much (which in consequence would slow down the whole copying process)
@@ -364,7 +374,8 @@ start_receiver(Tab,Storage,Cs,SenderPid,TabSize,DetsData,{dumper,{add_table_copy
364374
Else
365375
end.
366376

367-
spawned_receiver(ReplyTo,Tab,Storage,Cs, SenderPid,TabSize,DetsData, Init) ->
377+
spawned_receiver(ReplyTo,Tab,Storage,Cs, SenderPid,TabSize,DetsData) ->
378+
Init = table_init_fun(SenderPid, Storage),
368379
process_flag(trap_exit, true),
369380
Done = do_init_table(Tab,Storage,Cs,
370381
SenderPid,TabSize,DetsData,
@@ -546,26 +557,43 @@ init_table(Tab, {ext,Alias,Mod}, Fun, State, Sender) ->
546557
ext_init_table(Alias, Mod, Tab, Fun, State, Sender);
547558
init_table(Tab, disc_only_copies, Fun, DetsInfo,Sender) ->
548559
ErtsVer = erlang:system_info(version),
549-
case DetsInfo of
550-
{ErtsVer, DetsData} ->
551-
try dets:is_compatible_bchunk_format(Tab, DetsData) of
552-
false ->
553-
Sender ! {self(), {old_protocol, Tab}},
554-
dets:init_table(Tab, Fun); %% Old dets version
555-
true ->
556-
dets:init_table(Tab, Fun, [{format, bchunk}])
557-
catch
558-
error:{undef,[{dets,_,_,_}|_]} ->
559-
Sender ! {self(), {old_protocol, Tab}},
560-
dets:init_table(Tab, Fun); %% Old dets version
561-
error:What ->
562-
What
563-
end;
564-
Old when Old /= false ->
565-
Sender ! {self(), {old_protocol, Tab}},
566-
dets:init_table(Tab, Fun); %% Old dets version
567-
_ ->
568-
dets:init_table(Tab, Fun)
560+
Parent = self(),
561+
DetsInit = fun(Opts) ->
562+
fun() ->
563+
put(mnesia_dets_worker, {Tab, node(Sender), Sender}),
564+
Res = dets:init_table(Tab, Fun, Opts),
565+
Parent ! {self(), Res},
566+
unlink(Parent),
567+
exit(normal)
568+
end
569+
end,
570+
Worker =
571+
case DetsInfo of
572+
{ErtsVer, DetsData} ->
573+
try dets:is_compatible_bchunk_format(Tab, DetsData) of
574+
false ->
575+
Sender ! {self(), {old_protocol, Tab}},
576+
spawn_link(DetsInit([])); %% Old dets version
577+
true ->
578+
spawn_link(DetsInit([{format, bchunk}]))
579+
catch
580+
error:{undef,[{dets,_,_,_}|_]} ->
581+
Sender ! {self(), {old_protocol, Tab}},
582+
spawn_link(DetsInit([])); %% Old dets version
583+
error:What ->
584+
What
585+
end;
586+
Old when Old /= false ->
587+
Sender ! {self(), {old_protocol, Tab}},
588+
spawn_link(DetsInit([])); %% Old dets version
589+
_ ->
590+
spawn_link(DetsInit([]))
591+
end,
592+
case is_pid(Worker) of
593+
true ->
594+
init_dets_receiver(Worker);
595+
_ ->
596+
Worker
569597
end;
570598
init_table(Tab, _, Fun, _DetsInfo,_) ->
571599
try
@@ -574,6 +602,36 @@ init_table(Tab, _, Fun, _DetsInfo,_) ->
574602
catch _:Else:Stacktrace -> {Else, Stacktrace}
575603
end.
576604

605+
init_dets_receiver(Worker) ->
606+
receive
607+
{dets_receiver_pid, Receiver} ->
608+
%% We got the pid of the dets receiver, we can start forwarding {copier_done, Node}
609+
handle_dets_receiver(Worker, Receiver);
610+
{Worker, Result} ->
611+
%% Worker has finished before we got receiver pid
612+
Result;
613+
{'EXIT', Pid, Reason} ->
614+
%% If some local pid crashes we also crash
615+
handle_exit(Pid, Reason),
616+
init_dets_receiver(Pid)
617+
end.
618+
619+
handle_dets_receiver(Worker, Receiver) ->
620+
receive
621+
{Worker, Result} ->
622+
%% Dets worker has finished with ok | {error, Reason}
623+
Result;
624+
{copier_done, Node} ->
625+
%% Forward any {copier_done, Node} blindly
626+
%% they will either be ignored, or cause Worker
627+
%% to finish with ok | {error, Reason}
628+
Receiver ! {copier_done, Node},
629+
handle_dets_receiver(Worker, Receiver);
630+
{'EXIT', Pid, Reason} ->
631+
%% If some local pid crashes we also crash
632+
handle_exit(Pid, Reason),
633+
handle_dets_receiver(Worker, Receiver)
634+
end.
577635

578636
finish_copy(Storage,Tab,Cs,SenderPid,DatBin,OrigTabRec) ->
579637
TabRef = {Storage, Tab},

0 commit comments

Comments
 (0)