mnesia_tm.erl

来自「OTP是开放电信平台的简称」· ERL 代码 · 共 1,898 行 · 第 1/5 页

ERL
1,898
字号
			       [Pid | GoodPids], [Pid | SchemaAckPids]);	{?MODULE, _, {acc_pre_commit, Tid, Pid, false}} ->	    rec_acc_pre_commit(Tail, Tid, Store, Commit, Res, DumperMode,			       [Pid | GoodPids], SchemaAckPids);	{?MODULE, _, {acc_pre_commit, Tid, Pid}} ->	    %% Kept for backwards compatibility. Remove after Mnesia 4.x	    rec_acc_pre_commit(Tail, Tid, Store, Commit, Res, DumperMode,			       [Pid | GoodPids], [Pid | SchemaAckPids]);	{?MODULE, _, {do_abort, Tid, Pid, _Reason}} ->	    AbortRes = {do_abort, {bad_commit, node(Pid)}},	    rec_acc_pre_commit(Tail, Tid, Store, Commit, AbortRes, DumperMode,			       GoodPids, SchemaAckPids);	{mnesia_down, Node} when Node == node(Pid) ->	    AbortRes = {do_abort, {bad_commit, Node}},	    catch Pid ! {Tid, AbortRes},  %% Tell him that he has died	    rec_acc_pre_commit(Tail, Tid, Store, Commit, AbortRes, DumperMode,			       GoodPids, SchemaAckPids)    end;rec_acc_pre_commit([], Tid, Store, {Commit,OrigC}, Res, DumperMode, GoodPids, SchemaAckPids) ->    D = Commit#commit.decision,    case Res of	do_commit ->	    %% Now everybody knows that the others	    %% has voted yes. We also know that	    %% everybody are uncertain.	    prepare_sync_schema_commit(Store, SchemaAckPids),	    tell_participants(GoodPids, {Tid, committed}),	    D2 = D#decision{outcome = committed},			    mnesia_recover:log_decision(D2),            ?eval_debug_fun({?MODULE, rec_acc_pre_commit_log_commit},			    [{tid, Tid}]),	    %% Now we have safely logged committed	    %% and we can recover without asking others	    do_commit(Tid, Commit, DumperMode),            ?eval_debug_fun({?MODULE, rec_acc_pre_commit_done_commit},			    [{tid, Tid}]),	    sync_schema_commit(Tid, Store, SchemaAckPids),	    mnesia_locker:release_tid(Tid),	    ?MODULE ! {delete_transaction, Tid};		{do_abort, Reason} ->	    tell_participants(GoodPids, {Tid, {do_abort, Reason}}),	    D2 = D#decision{outcome = aborted},			    mnesia_recover:log_decision(D2),            ?eval_debug_fun({?MODULE, rec_acc_pre_commit_log_abort},			    [{tid, Tid}]),	    do_abort(Tid, OrigC),	    ?eval_debug_fun({?MODULE, rec_acc_pre_commit_done_abort},			    [{tid, Tid}])    end,    Res.%% Note all nodes in case of mnesia_down mgtprepare_sync_schema_commit(_Store, []) ->    ok;prepare_sync_schema_commit(Store, [Pid | Pids]) ->    ?ets_insert(Store, {waiting_for_commit_ack, node(Pid)}),    prepare_sync_schema_commit(Store, Pids).sync_schema_commit(_Tid, _Store, []) ->    ok;sync_schema_commit(Tid, Store, [Pid | Tail]) ->    receive	{?MODULE, _, {schema_commit, Tid, Pid}} ->	    ?ets_match_delete(Store, {waiting_for_commit_ack, node(Pid)}),	    sync_schema_commit(Tid, Store, Tail);	{mnesia_down, Node} when Node == node(Pid) ->	    ?ets_match_delete(Store, {waiting_for_commit_ack, Node}),	    sync_schema_commit(Tid, Store, Tail)    end.tell_participants([Pid | Pids], Msg) ->    Pid ! Msg,    tell_participants(Pids, Msg);tell_participants([], _Msg) ->    ok.%% No need for trapping exits. We are only linked%% to mnesia_tm and if it dies we should also die.%% The same goes for disk_log and dets.commit_participant(Coord, Tid, Bin, DiscNs, RamNs) when binary(Bin) ->    Commit = binary_to_term(Bin),    commit_participant(Coord, Tid, Bin, Commit, DiscNs, RamNs);commit_participant(Coord, Tid, C, DiscNs, RamNs) when record(C, commit) ->    commit_participant(Coord, Tid, C, C, DiscNs, RamNs).commit_participant(Coord, Tid, Bin, C0, DiscNs, _RamNs) ->    ?eval_debug_fun({?MODULE, commit_participant, pre}, [{tid, Tid}]),    case catch mnesia_schema:prepare_commit(Tid, C0, {part, Coord}) of	{Modified, C, DumperMode} when record(C, commit) ->	    %% If we can not find any local unclear decision	    %% we should presume abort at startup recovery	    case lists:member(node(), DiscNs) of		false ->		    ignore;		true ->		    case Modified of			false -> mnesia_log:log(Bin);			true  -> mnesia_log:log(C)		    end	    end,	    ?eval_debug_fun({?MODULE, commit_participant, vote_yes},			    [{tid, Tid}]),	    reply(Coord, {vote_yes, Tid, self()}),	    receive		{Tid, pre_commit} ->		    D = C#commit.decision,		    mnesia_recover:log_decision(D#decision{outcome = unclear}),		    ?eval_debug_fun({?MODULE, commit_participant, pre_commit},				    [{tid, Tid}]),		    Expect_schema_ack = C#commit.schema_ops /= [],		    reply(Coord, {acc_pre_commit, Tid, self(), Expect_schema_ack}),		    %% Now we are vulnerable for failures, since		    %% we cannot decide without asking others		    receive			{Tid, committed} ->			    mnesia_recover:log_decision(D#decision{outcome = committed}),			    ?eval_debug_fun({?MODULE, commit_participant, log_commit},					    [{tid, Tid}]),			    do_commit(Tid, C, DumperMode),			    case Expect_schema_ack of				false -> ignore;				true -> reply(Coord, {schema_commit, Tid, self()})			    end,			    ?eval_debug_fun({?MODULE, commit_participant, do_commit},					    [{tid, Tid}]);						{Tid, {do_abort, _Reason}} ->			    mnesia_recover:log_decision(D#decision{outcome = aborted}),			    ?eval_debug_fun({?MODULE, commit_participant, log_abort},					    [{tid, Tid}]),			    mnesia_schema:undo_prepare_commit(Tid, C0),			    ?eval_debug_fun({?MODULE, commit_participant, undo_prepare},					    [{tid, Tid}]);						{'EXIT', _, _} ->			    mnesia_recover:log_decision(D#decision{outcome = aborted}),			    ?eval_debug_fun({?MODULE, commit_participant, exit_log_abort},					    [{tid, Tid}]),			    mnesia_schema:undo_prepare_commit(Tid, C0),			    ?eval_debug_fun({?MODULE, commit_participant, exit_undo_prepare},					    [{tid, Tid}]);						Msg ->			    verbose("** ERROR ** commit_participant ~p, got unexpected msg: ~p~n",				    [Tid, Msg])		    end;		{Tid, {do_abort, Reason}} ->		    reply(Coord, {do_abort, Tid, self(), Reason}),		    mnesia_schema:undo_prepare_commit(Tid, C0),		    ?eval_debug_fun({?MODULE, commit_participant, pre_commit_undo_prepare},				    [{tid, Tid}]);		{'EXIT', _, Reason} ->		    reply(Coord, {do_abort, Tid, self(), {bad_commit,Reason}}),		    mnesia_schema:undo_prepare_commit(Tid, C0),		    ?eval_debug_fun({?MODULE, commit_participant, pre_commit_undo_prepare}, [{tid, Tid}]);		Msg ->		    reply(Coord, {do_abort, Tid, self(), {bad_commit,internal}}),		    verbose("** ERROR ** commit_participant ~p, got unexpected msg: ~p~n",			    [Tid, Msg])	    end;	{'EXIT', Reason} ->	    ?eval_debug_fun({?MODULE, commit_participant, vote_no},			    [{tid, Tid}]),	    reply(Coord, {vote_no, Tid, Reason}),	    mnesia_schema:undo_prepare_commit(Tid, C0)    end,    mnesia_locker:release_tid(Tid),    ?MODULE ! {delete_transaction, Tid},    unlink(whereis(?MODULE)),    exit(normal).    do_abort(Tid, Bin) when binary(Bin) ->    %% Possible optimization:    %% If we want we could pass arround a flag    %% that tells us whether the binary contains    %% schema ops or not. Only if the binary    %% contains schema ops there are meningful    %% unpack the binary and perform    %% mnesia_schema:undo_prepare_commit/1.    do_abort(Tid, binary_to_term(Bin));do_abort(Tid, Commit) ->    mnesia_schema:undo_prepare_commit(Tid, Commit),     Commit.do_dirty(Tid, Commit) when Commit#commit.schema_ops == [] ->    mnesia_log:log(Commit),    do_commit(Tid, Commit).%% do_commit(Tid, CommitRecord)do_commit(Tid, Bin) when binary(Bin) ->    do_commit(Tid, binary_to_term(Bin));do_commit(Tid, C) ->    do_commit(Tid, C, optional).do_commit(Tid, Bin, DumperMode) when binary(Bin) ->    do_commit(Tid, binary_to_term(Bin), DumperMode);do_commit(Tid, C, DumperMode) ->    mnesia_dumper:update(Tid, C#commit.schema_ops, DumperMode),    R  = do_snmp(Tid, C#commit.snmp),    R2 = do_update(Tid, ram_copies, C#commit.ram_copies, R),    R3 = do_update(Tid, disc_copies, C#commit.disc_copies, R2),    do_update(Tid, disc_only_copies, C#commit.disc_only_copies, R3).%% Update the itemsdo_update(Tid, Storage, [Op | Ops], OldRes) ->    case catch do_update_op(Tid, Storage, Op) of	ok ->	    do_update(Tid, Storage, Ops, OldRes);	{'EXIT', Reason} ->	    %% This may only happen when we recently have	    %% deleted our local replica, changed storage_type	    %% or transformed table	    %% BUGBUG: Updates may be lost if storage_type is changed.	    %%         Determine actual storage type and try again.	    %% BUGBUG: Updates may be lost if table is transformed.	    verbose("do_update in ~w failed: ~p -> {'EXIT', ~p}~n",		    [Tid, Op, Reason]),	    do_update(Tid, Storage, Ops, OldRes); 	NewRes ->	    do_update(Tid, Storage, Ops, NewRes)    end;do_update(_Tid, _Storage, [], Res) ->    Res.do_update_op(Tid, Storage, {{Tab, K}, Obj, write}) ->    commit_write(?catch_val({Tab, commit_work}), Tid, 		 Tab, K, Obj, undefined),    mnesia_lib:db_put(Storage, Tab, Obj);do_update_op(Tid, Storage, {{Tab, K}, Val, delete}) ->    commit_delete(?catch_val({Tab, commit_work}), Tid, Tab, K, Val, undefined),    mnesia_lib:db_erase(Storage, Tab, K);do_update_op(Tid, Storage, {{Tab, K}, {RecName, Incr}, update_counter}) ->    {NewObj, OldObjs} =         case catch mnesia_lib:db_update_counter(Storage, Tab, K, Incr) of            NewVal when integer(NewVal), NewVal >= 0 ->                {{RecName, K, NewVal}, [{RecName, K, NewVal - Incr}]};            _ when Incr > 0 ->                New = {RecName, K, Incr},                mnesia_lib:db_put(Storage, Tab, New),                {New, []};	    _ -> 		Zero = {RecName, K, 0},		mnesia_lib:db_put(Storage, Tab, Zero),		{Zero, []}        end,    commit_update(?catch_val({Tab, commit_work}), Tid, Tab, 		  K, NewObj, OldObjs),    element(3, NewObj);do_update_op(Tid, Storage, {{Tab, Key}, Obj, delete_object}) ->    commit_del_object(?catch_val({Tab, commit_work}), 		      Tid, Tab, Key, Obj, undefined),    mnesia_lib:db_match_erase(Storage, Tab, Obj);do_update_op(Tid, Storage, {{Tab, Key}, Obj, clear_table}) ->    commit_clear(?catch_val({Tab, commit_work}), Tid, Tab, Key, Obj),    mnesia_lib:db_match_erase(Storage, Tab, Obj).commit_write([], _, _, _, _, _) -> ok;commit_write([{checkpoints, CpList}|R], Tid, Tab, K, Obj, Old) ->    mnesia_checkpoint:tm_retain(Tid, Tab, K, write, CpList),    commit_write(R, Tid, Tab, K, Obj, Old);commit_write([H|R], Tid, Tab, K, Obj, Old)   when element(1, H) == subscribers ->    mnesia_subscr:report_table_event(H, Tab, Tid, Obj, write, Old),    commit_write(R, Tid, Tab, K, Obj, Old);commit_write([H|R], Tid, Tab, K, Obj, Old)   when element(1, H) == index ->    mnesia_index:add_index(H, Tab, K, Obj, Old),    commit_write(R, Tid, Tab, K, Obj, Old).commit_update([], _, _, _, _, _) -> ok;commit_update([{checkpoints, CpList}|R], Tid, Tab, K, Obj, _) ->    Old = mnesia_checkpoint:tm_retain(Tid, Tab, K, write, CpList),    commit_update(R, Tid, Tab, K, Obj, Old);commit_update([H|R], Tid, Tab, K, Obj, Old)   when element(1, H) == subscribers ->    mnesia_subscr:report_table_event(H, Tab, Tid, Obj, write, Old),    commit_update(R, Tid, Tab, K, Obj, Old);commit_update([H|R], Tid, Tab, K, Obj, Old)   when element(1, H) == index ->    mnesia_index:add_index(H, Tab, K, Obj, Old),    commit_update(R, Tid, Tab, K, Obj, Old).commit_delete([], _, _, _, _, _) ->  ok;commit_delete([{checkpoints, CpList}|R], Tid, Tab, K, Obj, _) ->    Old = mnesia_checkpoint:tm_retain(Tid, Tab, K, delete, CpList),    commit_delete(R, Tid, Tab, K, Obj, Old);commit_delete([H|R], Tid, Tab, K, Obj, Old)   when element(1, H) == subscribers ->    mnesia_subscr:report_table_event(H, Tab, Tid, Obj, delete, Old),    commit_delete(R, Tid, Tab, K, Obj, Old);commit_delete([H|R], Tid, Tab, K, Obj, Old)   when element(1, H) == index ->    mnesia_index:delete_index(H, Tab, K),    commit_delete(R, Tid, Tab, K, Obj, Old).commit_del_object([], _, _, _, _, _) -> ok;commit_del_object([{checkpoints, CpList}|R], Tid, Tab, K, Obj, _) ->    Old = mnesia_checkpoint:tm_retain(Tid, Tab, K, delete_object, CpList),    commit_del_object(R, Tid, Tab, K, Obj, Old);commit_del_object([H|R], Tid, Tab, K, Obj, Old)   when element(1, H) == subscribers ->     mnesia_subscr:report_table_event(H, Tab, Tid, Obj, delete_object, Old),    commit_del_object(R, Tid, Tab, K, Obj, Old);commit_del_object([H|R], Tid, Tab, K, Obj, Old)   when element(1, H) == index ->     mnesia_index:del_object_index(H, Tab, K, Obj, Old),    commit_del_object(R, Tid, Tab, K, Obj, Old).commit_clear([], _, _, _, _) ->  ok;commit_clear([{checkpoints, CpList}|R], Tid, Tab, K, Obj) ->    mnesia_checkpoint:tm_retain(Tid, Tab, K, clear_table, CpList),    commit_clear(R, Tid, Tab, K, Obj);commit_clear([H|R], Tid, Tab, K, Obj)   when element(1, H) == subscribers ->    mnesia_subscr:report_table_event(H, Tab, Tid, Obj, clear_table, undefined),    commit_clear(R, Tid, Tab, K, Obj);commit_clear([H|R], Tid, Tab, K, Obj)   when element(1, H) == index ->    mnesia_index:clear_index(H, Tab, K, Obj),    commit_clear(R, Tid, Tab, K, Obj).do_snmp(_, []) ->   ok;do_snmp(Tid, [Head | Tail]) ->    case catch mnesia_snmp_hook:update(Head) of	{'EXIT', Reason} ->	    %% This should only happen when we recently have	    %% deleted our local replica or recently deattached	    %% the snmp table 	    verbose("do_snmp in ~w failed: ~p -> {'EXIT', ~p}~n",		    [Tid, Head, Reason]);	ok ->	    ignore    end,    do_snmp(Tid, Tail).commit_nodes([C | Tail], AccD, AccR)         when C#commit.disc_copies == [],             C#commit.disc_only_copies  == [],             C#commit.schema_ops == [] ->    commit_nodes(Tail, AccD, [C#commit.node | AccR]);commit_nodes([C | Tail], AccD, AccR) ->    commit_nodes(Tail, [C#commit.node | AccD], AccR);commit_nodes([], AccD, AccR) ->    {AccD, AccR}.commit_decision(D, [C | Tail], AccD, AccR) ->    N = C#commit.node,    {D2, Tail2} = 	case C#commit.schema_ops of	    [] when C#commit.disc_copies == [],		    C#commit.disc_only_copies  == [] ->		commit_decision(D, Tail, AccD, [N | AccR]);	    [] ->		commit_decision(D, Tail, [N | AccD], AccR);	    Ops ->		case ram_only_ops(N, Ops) of		    true ->			commit_decision(D, Tail, AccD, [N | AccR]);		    false ->			commit_decision(D, Tail, [N | AccD], AccR)		end	end, 

⌨️ 快捷键说明

复制代码Ctrl + C
搜索代码Ctrl + F
全屏模式F11
增大字号Ctrl + =
减小字号Ctrl + -
显示快捷键?