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
341340table_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 );
547558init_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 ;
570598init_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+ % % Dets worker has finished with ok | {error, Reason} 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 (Worker )
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
578636finish_copy (Storage ,Tab ,Cs ,SenderPid ,DatBin ,OrigTabRec ) ->
579637 TabRef = {Storage , Tab },
0 commit comments