global.erl
来自「OTP是开放电信平台的简称」· ERL 代码 · 共 1,642 行 · 第 1/5 页
ERL
1,642 行
case get({sync_tag_my, Node}) of MyTag -> Known = S#state.known, gen_server:cast({global_name_server, Node}, {resolved, node(), Resolved, Known, Known,get_names_ext(),get({sync_tag_his,Node})}), case get({save_ops, Node}) of {resolved, HisKnown, Names_ext, HisResolved} -> put({save_ops, Node}, Ops), NewS = resolved(Node, HisResolved, HisKnown, Names_ext,S), {noreply, NewS}; undefined -> put({save_ops, Node}, Ops), {noreply, S} end; _ -> %% Illegal tag, delete the old sync session. NewS = cancel_locker(Node, S, MyTag), {noreply, NewS} end;%%========================================================================%% resolved%%%% Here the name clashes are resolved.%%========================================================================handle_cast({resolved, Node, HisResolved, HisKnown, _HisKnown_v2, Names_ext, MyTag}, S) -> %% Sent from global_name_server at Node. ?trace({'####', resolved, {his_resolved,HisResolved}, {node,Node}}), case get({sync_tag_my, Node}) of MyTag -> %% See the comment at handle_case({exchange_ops, ...}). case get({save_ops, Node}) of Ops when is_list(Ops) -> NewS = resolved(Node, HisResolved, HisKnown, Names_ext, S), {noreply, NewS}; undefined -> Resolved = {resolved, HisKnown, Names_ext, HisResolved}, put({save_ops, Node}, Resolved), {noreply, S} end; _ -> %% Illegal tag, delete the old sync session. NewS = cancel_locker(Node, S, MyTag), {noreply, NewS} end;%%========================================================================%% new_nodes%%%% We get to know the other node's known nodes.%%========================================================================handle_cast({new_nodes, Node, Ops, Names_ext, Nodes, ExtraInfo}, S) -> %% Sent from global_name_server at Node. ?trace({new_nodes, {node,Node},{ops,Ops},{nodes,Nodes},{x,ExtraInfo}}), NewS = new_nodes(Ops, Node, Names_ext, Nodes, ExtraInfo, S), {noreply, NewS};%%========================================================================%% in_sync%%%% We are in sync with this node (from the other node's known world).%%========================================================================handle_cast({in_sync, Node, _IsKnown}, S) -> %% Sent from global_name_server at Node (in the other partition). ?trace({'####', in_sync, {Node, _IsKnown}}), lists:foreach(fun(Pid) -> Pid ! {synced, [Node]} end, S#state.syncers), NewS = cancel_locker(Node, S, get({sync_tag_my, Node})), reset_node_state(Node), NSynced = case lists:member(Node, Synced = NewS#state.synced) of true -> Synced; false -> [Node | Synced] end, {noreply, NewS#state{synced = NSynced}};%% Called when Pid on other node crashedhandle_cast({async_del_name, _Name, _Pid}, S) -> %% Sent from the_deleter at some node in the partition but node(). %% The DOWN message deletes the name. {noreply, S};handle_cast({async_del_lock, _ResourceId, _Pid}, S) -> %% Sent from global_name_server at some node in the partition but node(). %% The DOWN message deletes the lock. {noreply, S};handle_cast(Request, S) -> error_logger:warning_msg("The global_name_server " "received an unexpected message:\n" "handle_cast(~p, _)\n", [Request]), {noreply, S}.handle_info({'EXIT', Deleter, _Reason}=Exit, #state{the_deleter=Deleter}=S) -> {stop, {deleter_died,Exit}, S#state{the_deleter=undefined}};handle_info({'EXIT', Locker, _Reason}=Exit, #state{the_locker=Locker}=S) -> {stop, {locker_died,Exit}, S#state{the_locker=undefined}};handle_info({'EXIT', Registrar, _}=Exit, #state{the_registrar=Registrar}=S) -> {stop, {registrar_died,Exit}, S#state{the_registrar=undefined}};handle_info({'EXIT', Pid, _Reason}, S) when is_pid(Pid) -> ?trace({global_EXIT,_Reason,Pid}), %% The process that died was a synch process started by start_sync %% or a registered process running on an external node (C-node). %% Links to external names are ignored here (there are also DOWN %% signals). Syncers = lists:delete(Pid, S#state.syncers), {noreply, S#state{syncers = Syncers}};handle_info({nodedown, Node}, S) when Node =:= S#state.node_name -> %% Somebody stopped the distribution dynamically - change %% references to old node name (Node) to new node name ('nonode@nohost') {noreply, change_our_node_name(node(), S)};handle_info({nodedown, Node}, S0) -> ?trace({'####', nodedown, {node,Node}}), S1 = trace_message(S0, {nodedown, Node}, []), S = handle_nodedown(Node, S1), {noreply, S};handle_info({extra_nodedown, Node}, S0) -> ?trace({'####', extra_nodedown, {node,Node}}), S1 = trace_message(S0, {extra_nodedown, Node}, []), S = handle_nodedown(Node, S1), {noreply, S};handle_info({nodeup, Node}, S) when Node =:= node() -> ?trace({'####', local_nodeup, {node, Node}}), %% Somebody started the distribution dynamically - change %% references to old node name ('nonode@nohost') to Node. {noreply, change_our_node_name(Node, S)};handle_info({nodeup, Node}, S0) when S0#state.connect_all -> IsKnown = lists:member(Node, S0#state.known) or %% This one is only for double nodeups (shouldn't occur!) lists:keymember(Node, 1, S0#state.resolvers), ?trace({'####', nodeup, {node,Node}, {isknown,IsKnown}}), S1 = trace_message(S0, {nodeup, Node}, []), case IsKnown of true -> {noreply, S1}; false -> resend_pre_connect(Node), %% now() is used as a tag to separate different sycnh sessions %% from each others. Global could be confused at bursty nodeups %% because it couldn't separate the messages between the different %% synch sessions started by a nodeup. MyTag = now(), put({sync_tag_my, Node}, MyTag), ?trace({sending_nodeup_to_locker, {node,Node},{mytag,MyTag}}), S1#state.the_locker ! {nodeup, Node, MyTag}, %% In order to be compatible with unpatched R7 a locker %% process was spawned. Vsn 5 is no longer comptabible with %% vsn 3 nodes, so the locker process is no longer needed. %% The permanent locker takes its place. NotAPid = no_longer_a_pid, Locker = {locker, NotAPid, S1#state.known, S1#state.the_locker}, InitC = {init_connect, {?vsn, MyTag}, node(), Locker}, Rs = S1#state.resolvers, ?trace({casting_init_connect, {node,Node},{initmessage,InitC}, {resolvers,Rs}}), gen_server:cast({global_name_server, Node}, InitC), Resolver = start_resolver(Node, MyTag), S = trace_message(S1, {new_resolver, Node}, [MyTag, Resolver]), {noreply, S#state{resolvers = [{Node, MyTag, Resolver} | Rs]}} end;handle_info({whereis, Name, From}, S) -> do_whereis(Name, From), {noreply, S};handle_info(known, S) -> io:format(">>>> ~p\n",[S#state.known]), {noreply, S};%% "High level trace". For troubleshooting only.handle_info(high_level_trace, S) -> case S of #state{trace = [{Node, _Time, _M, Nodes, _X} | _]} -> send_high_level_trace(), CNode = node(), CNodes = nodes(), case {CNode, CNodes} of {Node, Nodes} -> {noreply, S}; _ -> {New, _, Old} = sofs:symmetric_partition(sofs:set([CNode|CNodes]), sofs:set([Node|Nodes])), M = {nodes_changed, {sofs:to_external(New), sofs:to_external(Old)}}, {noreply, trace_message(S, M, [])} end; _ -> {noreply, S} end;handle_info({trace_message, M}, S) -> {noreply, trace_message(S, M, [])};handle_info({trace_message, M, X}, S) -> {noreply, trace_message(S, M, X)};handle_info({'DOWN', MonitorRef, process, _Pid, _Info}, S0) -> S1 = delete_lock(MonitorRef, S0), S = del_name(MonitorRef, S1), {noreply, S};handle_info(Message, S) -> error_logger:warning_msg("The global_name_server " "received an unexpected message:\n" "handle_info(~p, _)\n", [Message]), {noreply, S}.%%========================================================================%%========================================================================%%=============================== Internal Functions =====================%%========================================================================%%========================================================================-define(HIGH_LEVEL_TRACE_INTERVAL, 500). % mswait_high_level_trace() -> receive high_level_trace -> ok after ?HIGH_LEVEL_TRACE_INTERVAL+1 -> ok end.send_high_level_trace() -> erlang:send_after(?HIGH_LEVEL_TRACE_INTERVAL, self(), high_level_trace).-define(GLOBAL_RID, global).%% Similar to trans(Id, Fun), but always uses global's own lock%% on all nodes known to global, making sure that no new nodes have%% become known while we got the list of known nodes.trans_all_known(Fun) -> Id = {?GLOBAL_RID, self()}, Nodes = set_lock_known(Id, 0), try Fun(Nodes) after delete_global_lock(Id, Nodes) end.set_lock_known(Id, Times) -> Known = get_known(), Nodes = [node() | Known], Boss = the_boss(Nodes), %% Use the same convention (a boss) as lock_nodes_safely. Optimization. case set_lock_on_nodes(Id, [Boss]) of true -> case lock_on_known_nodes(Id, Known, Nodes) of true -> Nodes; false -> del_lock(Id, [Boss]), random_sleep(Times), set_lock_known(Id, Times+1) end; false -> random_sleep(Times), set_lock_known(Id, Times+1) end.lock_on_known_nodes(Id, Known, Nodes) -> case set_lock_on_nodes(Id, Nodes) of true -> (get_known() -- Known) =:= []; false -> false end.set_lock_on_nodes(_Id, []) -> true;set_lock_on_nodes(Id, Nodes) -> case local_lock_check(Id, Nodes) of true -> Msg = {set_lock, Id}, {Replies, _} = gen_server:multi_call(Nodes, global_name_server, Msg), ?trace({set_lock,{me,self()},Id,{nodes,Nodes},{replies,Replies}}), check_replies(Replies, Id, Replies); false=Reply -> Reply end.%% Probe lock on local node to see if one should go on trying other nodes.local_lock_check(_Id, [_] = _Nodes) -> true;local_lock_check(Id, Nodes) -> not lists:member(node(), Nodes) orelse (can_set_lock(Id) =/= false).check_replies([{_Node, true} | T], Id, Replies) -> check_replies(T, Id, Replies);check_replies([{_Node, false=Reply} | _T], _Id, [_]) -> Reply;check_replies([{_Node, false=Reply} | _T], Id, Replies) -> TrueReplyNodes = [N || {N, true} <- Replies], ?trace({check_replies, {true_reply_nodes, TrueReplyNodes}}), gen_server:multi_call(TrueReplyNodes, global_name_server, {del_lock, Id}), Reply;check_replies([], _Id, _Replies) -> true.%%========================================================================%% Another node wants to synchronize its registered names with us.%% Both nodes must have a lock before they are allowed to continue.%%========================================================================init_connect(Vsn, Node, InitMsg, HisTag, Resolvers, S) -> %% It is always the responsibility of newer versions to understand %% older versions of the protocol. put({prot_vsn, Node}, Vsn), put({sync_tag_his, Node}, HisTag), case lists:keysearch(Node, 1, Resolvers) of {value, {Node, MyTag, _Resolver}} -> MyTag = get({sync_tag_my, Node}), % assertion case InitMsg of {locker, HisLocker, HisKnown} -> %% before vsn 5 ?trace({old_init_connect,{histhelocker,HisLocker}}), HisLocker ! {his_locker_new, S#state.the_locker, {HisKnown, S#state.known}}; {locker, _NoLongerAPid, HisKnown, HisTheLocker} -> %% vsn 5 ?trace({init_connect,{histhelocker,HisTheLocker}}), S#state.the_locker ! {his_the_locker, HisTheLocker, {Vsn,HisKnown}, S#state.known} end; false ->
⌨️ 快捷键说明
复制代码Ctrl + C
搜索代码Ctrl + F
全屏模式F11
增大字号Ctrl + =
减小字号Ctrl + -
显示快捷键?