global.erl
来自「OTP是开放电信平台的简称」· ERL 代码 · 共 1,642 行 · 第 1/5 页
ERL
1,642 行
end.lock_still_set(PidOrNode, ExtraInfo, S) -> case ets:lookup(global_locks, ?GLOBAL_RID) of [{?GLOBAL_RID, _LockReqId, PidRefs}] when is_pid(PidOrNode) -> %% Name registration. lists:keymember(PidOrNode, 1, PidRefs); [{?GLOBAL_RID, LockReqId, PidRefs}] when is_atom(PidOrNode) -> case extra_info(lock, ExtraInfo) of {?GLOBAL_RID, LockId} -> % R11B-4 or later LockReqId =:= LockId; undefined -> lock_still_set_old(PidOrNode, LockReqId, PidRefs) end; [] -> %% If the global lock was not removed by a DOWN message %% then we have a node that do not monitor locking pids %% (pre R11B-3), or an R11B-3 node (which does not ensure %% that {new_nodes, ...} arrives before {del_lock, ...}). not S#state.global_lock_down end.%%% The following is probably overkill. It is possible that this node%%% has been locked again, but it is a rare occasion.lock_still_set_old(_Node, ReqId, _PidRefs) when is_pid(ReqId) -> %% Cannot do better than return true. true;lock_still_set_old(Node, ReqId, PidRefs) when is_list(ReqId) -> %% Connection, version > 4, but before R11B-4. [P || {P, _RPid, _Ref} <- PidRefs, node(P) =:= Node] =/= [].extra_info(Tag, ExtraInfo) -> %% ExtraInfo used to be a list of nodes (vsn 2). case catch lists:keysearch(Tag, 1, ExtraInfo) of {value, {Tag, Info}} -> Info; _ -> undefined end.del_name(Ref, S) -> NameL = [{Name, Pid} || {_, Name} <- ets:lookup(global_pid_names, Ref), {_, Pid, _Method, _RPid, Ref1} <- ets:lookup(global_names, Name), Ref1 =:= Ref], ?trace({async_del_name, self(), NameL, Ref}), case NameL of [{Name, Pid}] -> del_names(Name, Pid, S), delete_global_name2(Name, S); [] -> S end.%% Send {async_del_name, ...} to old nodes (pre R11B-3).del_names(Name, Pid, S) -> Send = case ets:lookup(global_names_ext, Name) of [{Name, Pid, RegNode}] -> RegNode =:= node(); [] -> node(Pid) =:= node() end, if Send -> ?trace({del_names, {pid,Pid}, {name,Name}}), S#state.the_deleter ! {delete_name, self(), Name, Pid}; true -> ok end.%% Keeps the entry in global_names for whereis_name/1.delete_global_name_keep_pid(Name, S) -> case ets:lookup(global_names, Name) of [{Name, Pid, _Method, RPid, Ref}] -> delete_global_name2(Name, Pid, RPid, Ref, S); [] -> S end.delete_global_name2(Name, S) -> case ets:lookup(global_names, Name) of [{Name, Pid, _Method, RPid, Ref}] -> true = ets:delete(global_names, Name), delete_global_name2(Name, Pid, RPid, Ref, S); [] -> S end.delete_global_name2(Name, Pid, RPid, Ref, S) -> true = erlang:demonitor(Ref, [flush]), kill_monitor_proc(RPid, Pid), delete_global_name(Name, Pid), ?trace({delete_global_name,{item,Name},{pid,Pid}}), true = ets:delete_object(global_pid_names, {Pid, Name}), true = ets:delete_object(global_pid_names, {Ref, Name}), case ets:lookup(global_names_ext, Name) of [{Name, Pid, RegNode}] -> true = ets:delete(global_names_ext, Name), ?trace({delete_global_name, {name,Name,{pid,Pid},{RegNode,Pid}}}), dounlink_ext(Pid, RegNode); [] -> ?trace({delete_global_name,{name,Name,{pid,Pid},{node(Pid),Pid}}}), ok end, trace_message(S, {del_name, node(Pid)}, [Name, Pid]).%% delete_global_name/2 is traced by the inviso application. %% Do not change.delete_global_name(_Name, _Pid) -> ok.%%-----------------------------------------------------------------%% The locker is a satellite process to global_name_server. When a%% nodeup is received from a new node the global_name_server sends a%% message to the locker. The locker tries to set a lock in our%% partition, i.e. on all nodes known to us. When the lock is set, it%% tells global_name_server about it, and keeps the lock set.%% global_name_server sends a cancel message to the locker when the%% partitions are connected.%% There are two versions of the protocol between lockers on two nodes:%% Version 1: used by unpatched R7.%% Version 2: the messages exchanged between the lockers include the known %% nodes (see OTP-3576).%%------------------------------------------------------------------define(locker_vsn, 2).-record(multi, {local = [], % Requests from nodes on the local host. remote = [], % Other requests. known = [], % Copy of global_name_server's known nodes. It's % faster to keep a copy of known than asking % for it when needed. the_boss, % max([node() | 'known']) just_synced = false, % true if node() synced just a moment ago %% Statistics: do_trace % bool() }).-record(him, {node, locker, vsn, my_tag}).start_the_locker(DoTrace) -> spawn_link(fun() -> init_the_locker(DoTrace) end).init_the_locker(DoTrace) -> process_flag(trap_exit, true), % needed? S0 = #multi{do_trace = DoTrace}, S1 = update_locker_known({add, get_known()}, S0), loop_the_locker(S1), erlang:error(locker_exited).loop_the_locker(S) -> ?trace({loop_the_locker,S}), receive Message when element(1, Message) =/= nodeup -> the_locker_message(Message, S) after 0 -> Timeout = case {S#multi.local, S#multi.remote} of {[],[]} -> infinity; _ -> %% It is important that the timeout is greater %% than zero, or the chance that some other node %% in the partition sets the lock once this node %% has failed after setting the lock is very slim. if S#multi.just_synced -> 0; % no reason to wait after success S#multi.known =:= [] -> 200; % just to get started true -> lists:min([1000 + 100*length(S#multi.known), 3000]) end end, S1 = S#multi{just_synced = false}, receive Message when element(1, Message) =/= nodeup -> the_locker_message(Message, S1) after Timeout -> case is_global_lock_set() of true -> loop_the_locker(S1); false -> select_node(S1) end end end.the_locker_message({his_the_locker, HisTheLocker, HisKnown0, _MyKnown}, S) -> ?trace({his_the_locker, HisTheLocker, {node,node(HisTheLocker)}}), HisVsn = case HisKnown0 of {Vsn0, _} when Vsn0 > 4 -> Vsn0; _ when is_list(HisKnown0) -> 4 end, receive {nodeup, Node, MyTag} when node(HisTheLocker) =:= Node -> ?trace({the_locker_nodeup, {node,Node},{mytag,MyTag}}), Him = #him{node = node(HisTheLocker), my_tag = MyTag, locker = HisTheLocker, vsn = HisVsn}, loop_the_locker(add_node(Him, S)); {cancel, Node, _Tag, no_fun} when node(HisTheLocker) =:= Node -> loop_the_locker(S) after 60000 -> ?trace({nodeupnevercame, node(HisTheLocker)}), error_logger:error_msg("global: nodeup never came ~w ~w\n", [node(), node(HisTheLocker)]), loop_the_locker(S#multi{just_synced = false}) end;the_locker_message({cancel, _Node, undefined, no_fun}, S) -> ?trace({cancel_the_locker, undefined, {node,_Node}}), %% If we actually cancel something when a cancel message with the %% tag 'undefined' arrives, we may be acting on an old nodedown, %% to cancel a new nodeup, so we can't do that. loop_the_locker(S);the_locker_message({cancel, Node, Tag, no_fun}, S) -> ?trace({the_locker, cancel, {multi,S}, {tag,Tag},{node,Node}}), receive {nodeup, Node, Tag} -> ?trace({cancelnodeup2, {node,Node},{tag,Tag}}), ok after 0 -> ok end, loop_the_locker(remove_node(Node, S));the_locker_message({lock_set, _Pid, false, _}, S) -> ?trace({the_locker, spurious, {node,node(_Pid)}}), loop_the_locker(S);the_locker_message({lock_set, Pid, true, _HisKnown}, S) -> Node = node(Pid), ?trace({the_locker, self(), spontaneous, {node,Node}}), case find_node_tag(Node, S) of {true, MyTag, HisVsn} -> LockId = locker_lock_id(Pid, HisVsn), {IsLockSet, S1} = lock_nodes_safely(LockId, [], S), Pid ! {lock_set, self(), IsLockSet, S1#multi.known}, Known2 = [node() | S1#multi.known], ?trace({the_locker, spontaneous, {known2, Known2}, {node,Node}, {is_lock_set,IsLockSet}}), case IsLockSet of true -> gen_server:cast(global_name_server, {lock_is_set, Node, MyTag, LockId}), ?trace({lock_sync_done, {pid,Pid}, {node,node(Pid)}, {me,self()}}), %% Wait for global to tell us to remove lock. %% Should the other locker's node die, %% global_name_server will receive nodedown, and %% then send {cancel, Node, Tag}. receive {cancel, Node, _Tag, Fun} -> ?trace({cancel_the_lock,{node,Node}}), call_fun(Fun), delete_global_lock(LockId, Known2) end, S2 = S1#multi{just_synced = true}, loop_the_locker(remove_node(Node, S2)); false -> loop_the_locker(S1#multi{just_synced = false}) end; false -> ?trace({the_locker, not_there, {node,Node}}), Pid ! {lock_set, self(), false, S#multi.known}, loop_the_locker(S) end;the_locker_message({add_to_known, Nodes}, S) -> S1 = update_locker_known({add, Nodes}, S), loop_the_locker(S1);the_locker_message({remove_from_known, Node}, S) -> S1 = update_locker_known({remove, Node}, S), loop_the_locker(S1);the_locker_message({do_trace, DoTrace}, S) -> loop_the_locker(S#multi{do_trace = DoTrace});the_locker_message(Other, S) -> unexpected_message(Other, locker), ?trace({the_locker, {other_msg, Other}}), loop_the_locker(S).%% Requests from nodes on the local host are chosen before requests%% from other nodes. This should be a safe optimization.select_node(S) -> UseRemote = S#multi.local =:= [], Others1 = if UseRemote -> S#multi.remote; true -> S#multi.local end, Others2 = exclude_known(Others1, S#multi.known), S1 = if UseRemote -> S#multi{remote = Others2}; true -> S#multi{local = Others2} end, if Others2 =:= [] -> loop_the_locker(S1); true -> Him = random_element(Others2), #him{locker = HisTheLocker, vsn = HisVsn, node = Node, my_tag = MyTag} = Him, HisNode = if HisVsn < 5 -> []; true -> [Node] % prevents deadlock; optimization end, Us = [node() | HisNode], LockId = locker_lock_id(HisTheLocker, HisVsn), ?trace({select_node, self(), {us, Us}}), {IsLockSet, S2} = lock_nodes_safely(LockId, HisNode, S1), case IsLockSet of true -> Known1 = Us ++ S2#multi.known, ?trace({sending_lock_set, self(), {his,HisTheLocker}}), HisTheLocker ! {lock_set, self(), true, S2#multi.known}, %% OTP-4902 S3 = lock_set_loop(S2, Him, MyTag, Known1, LockId), loop_the_locker(S3); false -> loop_the_locker(S2) end end.%% Version 5: Both sides use the same requester id. Thereby the nodes%% common to both sides are locked by both locker processes. This%% means that the lock is still there when the 'new_nodes' message is%% received even if
⌨️ 快捷键说明
复制代码Ctrl + C
搜索代码Ctrl + F
全屏模式F11
增大字号Ctrl + =
减小字号Ctrl + -
显示快捷键?