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