global.erl
来自「OTP是开放电信平台的简称」· ERL 代码 · 共 1,642 行 · 第 1/5 页
ERL
1,642 行
?trace({init_connect,{pre_connect,Node},{histag,HisTag}}), put({pre_connect, Node}, {Vsn, InitMsg, HisTag}) end.%%========================================================================%% In the simple case, we'll get lock_is_set before we get exchange,%% but we may get exchange before we get lock_is_set from our locker.%% If that's the case, we'll have to remember the exchange info, and%% handle it when we get the lock_is_set. We do this by using the%% process dictionary - when the lock_is_set msg is received, we store%% this info. When exchange is received, we can check the dictionary%% if the lock_is_set has been received. If not, we store info about%% the exchange instead. In the lock_is_set we must first check if%% exchange info is stored, in that case we take care of it.%%========================================================================lock_is_set(Node, Resolvers, LockId) -> gen_server:cast({global_name_server, Node}, {exchange, node(), get_names(), _ExtNames = [], get({sync_tag_his, Node})}), put({lock_id, Node}, LockId), %% If both have the lock, continue with exchange. case get({wait_lock, Node}) of {exchange, NameList} -> put({wait_lock, Node}, lock_is_set), exchange(Node, NameList, Resolvers); undefined -> put({wait_lock, Node}, lock_is_set) end.%%========================================================================%% exchange%%========================================================================exchange(Node, NameList, Resolvers) -> ?trace({'####', exchange, {node,Node}, {namelist,NameList}, {resolvers, Resolvers}}), case erase({wait_lock, Node}) of lock_is_set -> {value, {Node, _Tag, Resolver}} = lists:keysearch(Node, 1, Resolvers), Resolver ! {resolve, NameList, Node}; undefined -> put({wait_lock, Node}, {exchange, NameList}) end.resolved(Node, HisResolved, HisKnown, Names_ext, S0) -> Ops = erase({save_ops, Node}) ++ HisResolved, %% Known may have shrunk since the lock was taken (due to nodedowns). Known = S0#state.known, Synced = S0#state.synced, NewNodes = [Node | HisKnown], sync_others(HisKnown), ExtraInfo = [{vsn,get({prot_vsn, Node})}, {lock, get({lock_id, Node})}], S = do_ops(Ops, node(), Names_ext, ExtraInfo, S0), %% I am synced with Node, but not with HisKnown yet lists:foreach(fun(Pid) -> Pid ! {synced, [Node]} end, S#state.syncers), S3 = lists:foldl(fun(Node1, S1) -> F = fun(Tag) -> cancel_locker(Node1,S1,Tag) end, cancel_resolved_locker(Node1, F) end, S, HisKnown), %% The locker that took the lock is asked to send %% the {new_nodes, ...} message. This ensures that %% {del_lock, ...} is received after {new_nodes, ...} %% (except when abcast spawns process(es)...). NewNodesF = fun() -> gen_server:abcast(Known, global_name_server, {new_nodes, node(), Ops, Names_ext, NewNodes, ExtraInfo}) end, F = fun(Tag) -> cancel_locker(Node, S3, Tag, NewNodesF) end, S4 = cancel_resolved_locker(Node, F), %% See (*) below... we're node b in that description AddedNodes = (NewNodes -- Known), NewKnown = Known ++ AddedNodes, S4#state.the_locker ! {add_to_known, AddedNodes}, NewS = trace_message(S4, {added, AddedNodes}, [{new_nodes, NewNodes}, {abcast, Known}, {ops,Ops}]), NewS#state{known = NewKnown, synced = [Node | Synced]}.cancel_resolved_locker(Node, CancelFun) -> Tag = get({sync_tag_my, Node}), ?trace({calling_cancel_locker,Tag,get()}), S = CancelFun(Tag), reset_node_state(Node), S.new_nodes(Ops, ConnNode, Names_ext, Nodes, ExtraInfo, S0) -> Known = S0#state.known, %% (*) This one requires some thought... %% We're node a, other nodes b and c: %% The problem is that {in_sync, a} may arrive before {resolved, [a]} to %% b from c, leading to b sending {new_nodes, [a]} to us (node a). %% Therefore, we make sure we never get duplicates in Known. AddedNodes = lists:delete(node(), Nodes -- Known), sync_others(AddedNodes), S = do_ops(Ops, ConnNode, Names_ext, ExtraInfo, S0), ?trace({added_nodes_in_sync,{added_nodes,AddedNodes}}), S#state.the_locker ! {add_to_known, AddedNodes}, S1 = trace_message(S, {added, AddedNodes}, [{ops,Ops}]), S1#state{known = Known ++ AddedNodes}.do_whereis(Name, From) -> case is_global_lock_set() of false -> gen_server:reply(From, where(Name)); true -> send_again({whereis, Name, From}) end.terminate(_Reason, _S) -> true = ets:delete(global_names), true = ets:delete(global_names_ext), true = ets:delete(global_locks), true = ets:delete(global_pid_names), true = ets:delete(global_pid_ids).code_change(_OldVsn, S, _Extra) -> {ok, S}.%% The resolver runs exchange_names in a separate process. The effect%% is that locks can be used at the same time as name resolution takes%% place.start_resolver(Node, MyTag) -> spawn(fun() -> resolver(Node, MyTag) end).resolver(Node, Tag) -> receive {resolve, NameList, Node} -> ?trace({resolver, {me,self()}, {node,Node}, {namelist,NameList}}), {Ops, Resolved} = exchange_names(NameList, Node, [], []), Exchange = {exchange_ops, Node, Tag, Ops, Resolved}, gen_server:cast(global_name_server, Exchange), exit(normal); _ -> % Ignore garbage. resolver(Node, Tag) end.resend_pre_connect(Node) -> case erase({pre_connect, Node}) of {Vsn, InitMsg, HisTag} -> gen_server:cast(self(), {init_connect, {Vsn, HisTag}, Node, InitMsg}); _ -> ok end.ins_name(Name, Pid, Method, FromPidOrNode, ExtraInfo, S0) -> ?trace({ins_name,insert,{name,Name},{pid,Pid}}), S1 = delete_global_name_keep_pid(Name, S0), S = trace_message(S1, {ins_name, node(Pid)}, [Name, Pid]), insert_global_name(Name, Pid, Method, FromPidOrNode, ExtraInfo, S).ins_name_ext(Name, Pid, Method, RegNode, FromPidOrNode, ExtraInfo, S0) -> ?trace({ins_name_ext, {name,Name}, {pid,Pid}}), S1 = delete_global_name_keep_pid(Name, S0), dolink_ext(Pid, RegNode), S = trace_message(S1, {ins_name_ext, node(Pid)}, [Name, Pid]), true = ets:insert(global_names_ext, {Name, Pid, RegNode}), insert_global_name(Name, Pid, Method, FromPidOrNode, ExtraInfo, S).where(Name) -> case ets:lookup(global_names, Name) of [{_Name, Pid, _Method, _RPid, _Ref}] -> Pid; [] -> undefined end.handle_set_lock(Id, Pid, S) -> ?trace({handle_set_lock, Id, Pid}), case can_set_lock(Id) of {true, PidRefs} -> case pid_is_locking(Pid, PidRefs) of true -> {true, S}; false -> {true, insert_lock(Id, Pid, PidRefs, S)} end; false=Reply -> {Reply, S} end.can_set_lock({ResourceId, LockRequesterId}) -> case ets:lookup(global_locks, ResourceId) of [{ResourceId, LockRequesterId, PidRefs}] -> {true, PidRefs}; [{ResourceId, _LockRequesterId2, _PidRefs}] -> false; [] -> {true, []} end.insert_lock({ResourceId, LockRequesterId}=Id, Pid, PidRefs, S) -> {RPid, Ref} = do_monitor(Pid), true = ets:insert(global_pid_ids, {Pid, ResourceId}), true = ets:insert(global_pid_ids, {Ref, ResourceId}), Lock = {ResourceId, LockRequesterId, [{Pid,RPid,Ref} | PidRefs]}, true = ets:insert(global_locks, Lock), trace_message(S, {ins_lock, node(Pid)}, [Id, Pid]).is_global_lock_set() -> is_lock_set(?GLOBAL_RID).is_lock_set(ResourceId) -> ets:member(global_locks, ResourceId).handle_del_lock({ResourceId, LockReqId}, Pid, S0) -> ?trace({handle_del_lock, {pid,Pid},{id,{ResourceId, LockReqId}}}), case ets:lookup(global_locks, ResourceId) of [{ResourceId, LockReqId, PidRefs}]-> remove_lock(ResourceId, LockReqId, Pid, PidRefs, false, S0); _ -> S0 end.remove_lock(ResourceId, LockRequesterId, Pid, [{Pid,RPid,Ref}], Down, S0) -> ?trace({remove_lock_1, {id,ResourceId},{pid,Pid}}), true = erlang:demonitor(Ref, [flush]), kill_monitor_proc(RPid, Pid), true = ets:delete(global_locks, ResourceId), true = ets:delete_object(global_pid_ids, {Pid, ResourceId}), true = ets:delete_object(global_pid_ids, {Ref, ResourceId}), S = case ResourceId of ?GLOBAL_RID -> S0#state{global_lock_down = Down}; _ -> S0 end, trace_message(S, {rem_lock, node(Pid)}, [{ResourceId, LockRequesterId}, Pid]);remove_lock(ResourceId, LockRequesterId, Pid, PidRefs0, _Down, S) -> ?trace({remove_lock_2, {id,ResourceId},{pid,Pid}}), PidRefs = case lists:keysearch(Pid, 1, PidRefs0) of {value, {Pid, RPid, Ref}} -> true = erlang:demonitor(Ref, [flush]), kill_monitor_proc(RPid, Pid), true = ets:delete_object(global_pid_ids, {Ref, ResourceId}), lists:keydelete(Pid, 1, PidRefs0); false -> PidRefs0 end, Lock = {ResourceId, LockRequesterId, PidRefs}, true = ets:insert(global_locks, Lock), true = ets:delete_object(global_pid_ids, {Pid, ResourceId}), trace_message(S, {rem_lock, node(Pid)}, [{ResourceId, LockRequesterId}, Pid]).kill_monitor_proc(Pid, Pid) -> ok;kill_monitor_proc(RPid, _Pid) -> exit(RPid, kill).do_ops(Ops, ConnNode, Names_ext, ExtraInfo, S0) -> ?trace({do_ops, {ops,Ops}}), XInserts = [{Name, Pid, RegNode, Method} || {Name2, Pid2, RegNode} <- Names_ext, {insert, {Name, Pid, Method}} <- Ops, Name =:= Name2, Pid =:= Pid2], S1 = lists:foldl(fun({Name, Pid, RegNode, Method}, S1) -> ins_name_ext(Name, Pid, Method, RegNode, ConnNode, ExtraInfo, S1) end, S0, XInserts), XNames = [Name || {Name, _Pid, _RegNode, _Method} <- XInserts], Inserts = [{Name, Pid, node(Pid), Method} || {insert, {Name, Pid, Method}} <- Ops, not lists:member(Name, XNames)], S2 = lists:foldl(fun({Name, Pid, _RegNode, Method}, S2) -> ins_name(Name, Pid, Method, ConnNode, ExtraInfo, S2) end, S1, Inserts), DelNames = [Name || {delete, Name} <- Ops], lists:foldl(fun(Name, S) -> delete_global_name2(Name, S) end, S2, DelNames).%% It is possible that a node that was up and running when the%% operations were assembled has since died. The final {in_sync,...}%% messages do not generate nodedown messages for such nodes. To%% compensate "artificial" nodedown messages are created. Since%% monitor_node may take some time processes are spawned to avoid%% locking up the global_name_server. Should somehow double nodedown%% messages occur (one of them artificial), nothing bad can happen%% (the second nodedown is a no-op). It is assumed that there cannot%% be a nodeup before the artificial nodedown.%%%% The extra nodedown messages generated here also take care of the%% case that a nodedown message is received _before_ the operations%% are run.sync_others(Nodes) -> N = case application:get_env(kernel, ?N_CONNECT_RETRIES) of {ok, NRetries} when is_integer(NRetries), NRetries >= 0 -> NRetries; _ -> ?DEFAULT_N_CONNECT_RETRIES end, lists:foreach(fun(Node) -> spawn(fun() -> sync_other(Node, N) end) end, Nodes).sync_other(Node, N) -> erlang:monitor_node(Node, true, [allow_passive_connect]), receive {nodedown, Node} when N > 0 -> sync_other(Node, N - 1); {nodedown, Node} -> ?trace({missing_nodedown, {node, Node}}), error_logger:warning_msg("global: ~w failed to connect to ~w\n", [node(), Node]), global_name_server ! {extra_nodedown, Node} after 0 -> gen_server:cast({global_name_server,Node}, {in_sync,node(),true}) end. % monitor_node(Node, false), % exit(normal).insert_global_name(Name, Pid, Method, FromPidOrNode, ExtraInfo, S) -> {RPid, Ref} = do_monitor(Pid), true = ets:insert(global_names, {Name, Pid, Method, RPid, Ref}), true = ets:insert(global_pid_names, {Pid, Name}), true = ets:insert(global_pid_names, {Ref, Name}), case lock_still_set(FromPidOrNode, ExtraInfo, S) of true -> S; false -> %% The node that took the lock has gone down and then up %% again. The {register, ...} or {new_nodes, ...} message %% was delayed and arrived after nodeup (maybe it caused %% the nodeup). The DOWN signal from the monitor of the %% lock has removed the lock. %% Note: it is assumed here that the DOWN signal arrives %% _before_ nodeup and any message that caused nodeup. %% This is true of Erlang/OTP. delete_global_name2(Name, S)
⌨️ 快捷键说明
复制代码Ctrl + C
搜索代码Ctrl + F
全屏模式F11
增大字号Ctrl + =
减小字号Ctrl + -
显示快捷键?