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 + -
显示快捷键?