mnesia_locker.erl

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

ERL
1,176
字号
%% ``The contents of this file are subject to the Erlang Public License,%% Version 1.1, (the "License"); you may not use this file except in%% compliance with the License. You should have received a copy of the%% Erlang Public License along with this software. If not, it can be%% retrieved via the world wide web at http://www.erlang.org/.%% %% Software distributed under the License is distributed on an "AS IS"%% basis, WITHOUT WARRANTY OF ANY KIND, either express or implied. See%% the License for the specific language governing rights and limitations%% under the License.%% %% The Initial Developer of the Original Code is Ericsson Utvecklings AB.%% Portions created by Ericsson are Copyright 1999, Ericsson Utvecklings%% AB. All Rights Reserved.''%% %%     $Id$%%-module(mnesia_locker).-export([	 get_held_locks/0,	 get_lock_queue/0,	 global_lock/5,	 ixrlock/5,	 init/1,	 mnesia_down/2,	 release_tid/1,	 async_release_tid/2,	 send_release_tid/2,	 receive_release_tid_acc/2,	 rlock/3,	 rlock_table/3,	 rwlock/3,	 sticky_rwlock/3,	 start/0,	 sticky_wlock/3,	 sticky_wlock_table/3,	 wlock/3,	 wlock_no_exist/4,	 wlock_table/3	]).%% sys callback functions-export([system_continue/3,	 system_terminate/4,	 system_code_change/4	]).-include("mnesia.hrl").-import(mnesia_lib, [dbg_out/2, error/2, verbose/2]).-define(dbg(S,V), ok).%-define(dbg(S,V), dbg_out("~p:~p: " ++ S, [?MODULE, ?LINE] ++ V)).-define(ALL, '______WHOLETABLE_____').-define(STICK, '______STICK_____').-define(GLOBAL, '______GLOBAL_____').-record(state, {supervisor}).-record(queue, {oid, tid, op, pid, lucky}).%% mnesia_held_locks: contain       {Oid, Op, Tid} entries  (bag)-define(match_oid_held_locks(Oid),  {Oid, '_', '_'}).%% mnesia_tid_locks: contain        {Tid, Oid, Op} entries  (bag) -define(match_oid_tid_locks(Tid),   {Tid, '_', '_'}).%% mnesia_sticky_locks: contain     {Oid, Node} entries and {Tab, Node} entries (set)-define(match_oid_sticky_locks(Oid),{Oid, '_'}).%% mnesia_lock_queue: contain       {queue, Oid, Tid, Op, ReplyTo, WaitForTid} entries (bag)-define(match_oid_lock_queue(Oid),  #queue{oid=Oid, tid='_', op = '_', pid = '_', lucky = '_'}). %% mnesia_lock_counter:             {{write, Tab}, Number} &&%%                                  {{read, Tab}, Number} entries  (set)start() ->    mnesia_monitor:start_proc(?MODULE, ?MODULE, init, [self()]).init(Parent) ->    register(?MODULE, self()),    process_flag(trap_exit, true),    proc_lib:init_ack(Parent, {ok, self()}),    case ?catch_val(pid_sort_order) of	r9b_plain -> put(pid_sort_order, r9b_plain);	standard ->  put(pid_sort_order, standard);	_ -> ignore    end,    loop(#state{supervisor = Parent}).val(Var) ->    case ?catch_val(Var) of	{'EXIT', _ReASoN_} -> mnesia_lib:other_val(Var, _ReASoN_); 	_VaLuE_ -> _VaLuE_     end.reply(From, R) ->    From ! {?MODULE, node(), R}.l_request(Node, X, Store) ->    {?MODULE, Node} ! {self(), X},    l_req_rec(Node, Store).l_req_rec(Node, Store) ->    ?ets_insert(Store, {nodes, Node}),    receive 	{?MODULE, Node, Reply} -> 	    Reply;	{mnesia_down, Node} -> 	    {not_granted, {node_not_running, Node}}    end.release_tid(Tid) ->    ?MODULE ! {release_tid, Tid}.async_release_tid(Nodes, Tid) ->    rpc:abcast(Nodes, ?MODULE, {release_tid, Tid}).send_release_tid(Nodes, Tid) ->    rpc:abcast(Nodes, ?MODULE, {self(), {sync_release_tid, Tid}}).receive_release_tid_acc([Node | Nodes], Tid) ->    receive 	{?MODULE, Node, {tid_released, Tid}} -> 	    receive_release_tid_acc(Nodes, Tid);	{mnesia_down, Node} -> 	    receive_release_tid_acc(Nodes, Tid)    end;receive_release_tid_acc([], _Tid) ->    ok.loop(State) ->    receive	{From, {write, Tid, Oid}} ->	    try_sticky_lock(Tid, write, From, Oid),	    loop(State);	%% If Key == ?ALL it's a request to lock the entire table	%%	{From, {read, Tid, Oid}} ->	    try_sticky_lock(Tid, read, From, Oid),	    loop(State);	%% Really do a  read, but get hold of a write lock	%% used by mnesia:wread(Oid).		{From, {read_write, Tid, Oid}} ->	    try_sticky_lock(Tid, read_write, From, Oid),	    loop(State);		%% Tid has somehow terminated, clear up everything	%% and pass locks on to queued processes.	%% This is the purpose of the mnesia_tid_locks table		{release_tid, Tid} ->	    do_release_tid(Tid),	    loop(State);		%% stick lock, first tries this to the where_to_read Node	{From, {test_set_sticky, Tid, {Tab, _} = Oid, Lock}} ->	    case ?ets_lookup(mnesia_sticky_locks, Tab) of		[] -> 		    reply(From, not_stuck),		    loop(State);		[{_,Node}] when Node == node() ->		    %% Lock is stuck here, see now if we can just set 		    %% a regular write lock		    try_lock(Tid, Lock, From, Oid),		    loop(State);		[{_,Node}] ->		    reply(From, {stuck_elsewhere, Node}),		    loop(State)	    end;	%% If test_set_sticky fails, we send this to all nodes	%% after aquiring a real write lock on Oid	{stick, {Tab, _}, N} ->	    ?ets_insert(mnesia_sticky_locks, {Tab, N}),	    loop(State);	%% The caller which sends this message, must have first 	%% aquired a write lock on the entire table	{unstick, Tab} ->	    ?ets_delete(mnesia_sticky_locks, Tab),	    loop(State);	{From, {ix_read, Tid, Tab, IxKey, Pos}} ->	    case catch mnesia_index:get_index_table(Tab, Pos) of		{'EXIT', _} ->		    reply(From, {not_granted, {no_exists, Tab, {index, [Pos]}}}),		    loop(State);		Index ->		    Rk = mnesia_lib:elems(2,mnesia_index:db_get(Index, IxKey)),		    %% list of real keys		    case ?ets_lookup(mnesia_sticky_locks, Tab) of			[] ->			    set_read_lock_on_all_keys(Tid, From,Tab,Rk,Rk, 						      []),			    loop(State);			[{_,N}] when N == node() ->			    set_read_lock_on_all_keys(Tid, From,Tab,Rk,Rk, 						      []),			    loop(State);			[{_,N}] ->			    Req = {From, {ix_read, Tid, Tab, IxKey, Pos}},			    From ! {?MODULE, node(), {switch, N, Req}},			    loop(State)		    end	    end;	{From, {sync_release_tid, Tid}} ->	    do_release_tid(Tid),	    reply(From, {tid_released, Tid}),	    loop(State);		{release_remote_non_pending, Node, Pending} ->	    release_remote_non_pending(Node, Pending),	    mnesia_monitor:mnesia_down(?MODULE, Node),	    loop(State);	{'EXIT', Pid, _} when Pid == State#state.supervisor ->	    do_stop();	{system, From, Msg} ->	    verbose("~p got {system, ~p, ~p}~n", [?MODULE, From, Msg]),	    Parent = State#state.supervisor,	    sys:handle_system_msg(Msg, From, Parent, ?MODULE, [], State);		Msg ->	    error("~p got unexpected message: ~p~n", [?MODULE, Msg]),	    loop(State)    end.set_lock(Tid, Oid, Op) ->    ?dbg("Granted ~p ~p ~p~n", [Tid,Oid,Op]),    ?ets_insert(mnesia_held_locks, {Oid, Op, Tid}),    ?ets_insert(mnesia_tid_locks, {Tid, Oid, Op}).%%%%%%%%%%%%%%%%%%%%%%%%%%%%% Acquire lockstry_sticky_lock(Tid, Op, Pid, {Tab, _} = Oid) ->    case ?ets_lookup(mnesia_sticky_locks, Tab) of	[] ->	    try_lock(Tid, Op, Pid, Oid);	[{_,N}] when N == node() ->	    try_lock(Tid, Op, Pid, Oid);	[{_,N}] ->	    Req = {Pid, {Op, Tid, Oid}},	    Pid ! {?MODULE, node(), {switch, N, Req}}    end.try_lock(Tid, read_write, Pid, Oid) ->    try_lock(Tid, read_write, read, write, Pid, Oid);try_lock(Tid, Op, Pid, Oid) ->    try_lock(Tid, Op, Op, Op, Pid, Oid).try_lock(Tid, Op, SimpleOp, Lock, Pid, Oid) ->    case can_lock(Tid, Lock, Oid, {no, bad_luck}) of	yes ->	    Reply = grant_lock(Tid, SimpleOp, Lock, Oid),	    reply(Pid, Reply);	{no, Lucky} ->	    C = #cyclic{op = SimpleOp, lock = Lock, oid = Oid, lucky = Lucky},	    ?dbg("Rejected ~p ~p ~p ~p ~n", [Tid, Oid, Lock, Lucky]),	    reply(Pid, {not_granted, C});	{queue, Lucky} ->	    ?dbg("Queued ~p ~p ~p ~p ~n", [Tid, Oid, Lock, Lucky]),	    %% Append to queue: Nice place for trace output	    ?ets_insert(mnesia_lock_queue, 			#queue{oid = Oid, tid = Tid, op = Op, 			       pid = Pid, lucky = Lucky}),	    ?ets_insert(mnesia_tid_locks, {Tid, Oid, {queued, Op}})    end.grant_lock(Tid, read, Lock, {Tab, Key})  when Key /= ?ALL, Tab /= ?GLOBAL ->    case node(Tid#tid.pid) == node() of	true ->	    set_lock(Tid, {Tab, Key}, Lock),	    {granted, lookup_in_client};	false ->	    case catch mnesia_lib:db_get(Tab, Key) of %% lookup as well		{'EXIT', _Reason} ->		    %% Table has been deleted from this node,		    %% restart the transaction.		    C = #cyclic{op = read, lock = Lock, oid = {Tab, Key},				lucky = nowhere},		    {not_granted, C};		Val -> 		    set_lock(Tid, {Tab, Key}, Lock),		    {granted, Val}	    end    end;grant_lock(Tid, read, Lock, Oid) ->    set_lock(Tid, Oid, Lock),    {granted, ok};grant_lock(Tid, write, Lock, Oid) ->    set_lock(Tid, Oid, Lock),    granted.%% 1) Impose an ordering on all transactions favour old (low tid) transactions%%    newer (higher tid) transactions may never wait on older ones,%% 2) When releasing the tids from the queue always begin with youngest (high tid)%%    because of 1) it will avoid the deadlocks.%% 3) TabLocks is the problem :-) They should not starve and not deadlock %%    handle tablocks in queue as they had locks on unlocked records.can_lock(Tid, read, {Tab, Key}, AlreadyQ) when Key /= ?ALL ->    %% The key is bound, no need for the other BIF    Oid = {Tab, Key},     ObjLocks = ?ets_match_object(mnesia_held_locks, {Oid, write, '_'}),    TabLocks = ?ets_match_object(mnesia_held_locks, {{Tab, ?ALL}, write, '_'}),    check_lock(Tid, Oid, ObjLocks, TabLocks, yes, AlreadyQ, read);can_lock(Tid, read, Oid, AlreadyQ) -> % Whole tab    Tab = element(1, Oid),    ObjLocks = ?ets_match_object(mnesia_held_locks, {{Tab, '_'}, write, '_'}),    check_lock(Tid, Oid, ObjLocks, [], yes, AlreadyQ, read);can_lock(Tid, write, {Tab, Key}, AlreadyQ) when Key /= ?ALL ->     Oid = {Tab, Key},    ObjLocks = ?ets_lookup(mnesia_held_locks, Oid),    TabLocks = ?ets_lookup(mnesia_held_locks, {Tab, ?ALL}),    check_lock(Tid, Oid, ObjLocks, TabLocks, yes, AlreadyQ, write);can_lock(Tid, write, Oid, AlreadyQ) -> % Whole tab    Tab = element(1, Oid),    ObjLocks = ?ets_match_object(mnesia_held_locks, ?match_oid_held_locks({Tab, '_'})),    check_lock(Tid, Oid, ObjLocks, [], yes, AlreadyQ, write).%% Check held locks for conflicting lockscheck_lock(Tid, Oid, [Lock | Locks], TabLocks, X, AlreadyQ, Type) ->    case element(3, Lock) of	Tid ->	    check_lock(Tid, Oid, Locks, TabLocks, X, AlreadyQ, Type);	WaitForTid ->	    Queue = allowed_to_be_queued(WaitForTid,Tid),	    if Queue == true -> 		    check_lock(Tid, Oid, Locks, TabLocks, {queue, WaitForTid}, AlreadyQ, Type);	       Tid#tid.pid == WaitForTid#tid.pid ->		    dbg_out("Spurious lock conflict ~w ~w: ~w -> ~w~n",			    [Oid, Lock, Tid, WaitForTid]),  		    %% Test..		    {Tab, _Key} = Oid,		    HaveQ = (ets:lookup(mnesia_lock_queue, Oid) /= []) 			orelse (ets:lookup(mnesia_lock_queue,{Tab,?ALL}) /= []),		    if 			HaveQ -> 			    {no, WaitForTid};			true -> 			    check_lock(Tid,Oid,Locks,TabLocks,{queue,WaitForTid},AlreadyQ,Type)		    end;		    %%{no, WaitForTid};  Safe solution 	       true ->		    {no, WaitForTid}	    end    end;check_lock(_, _, [], [], X, {queue, bad_luck}, _) ->    X;  %% The queue should be correct already no need to check it againcheck_lock(_, _, [], [], X = {queue, _Tid}, _AlreadyQ, _) ->    X;  check_lock(Tid, Oid, [], [], X, AlreadyQ, Type) ->    {Tab, Key} = Oid,    if	Type == write ->	    check_queue(Tid, Tab, X, AlreadyQ);	Key == ?ALL ->	    %% hmm should be solvable by a clever select expr but not today...	    check_queue(Tid, Tab, X, AlreadyQ);	true ->	    %% If there is a queue on that object, read_lock shouldn't be granted	    ObjLocks = ets:lookup(mnesia_lock_queue, Oid),	    case max(ObjLocks) of		empty -> 		    check_queue(Tid, Tab, X, AlreadyQ);		ObjL ->		    case allowed_to_be_queued(ObjL,Tid) of			false ->			    %% Starvation Preemption (write waits for read)			    {no, ObjL};			true ->			    check_queue(Tid, Tab, {queue, ObjL}, AlreadyQ)		    end	    end    end;check_lock(Tid, Oid, [], TabLocks, X, AlreadyQ, Type) ->    check_lock(Tid, Oid, TabLocks, [], X, AlreadyQ, Type).

⌨️ 快捷键说明

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