fix(archive): async peer start в archive_manager, без краша gen_server. Refs EventHub/EventHubBack#33 [skip ci]
This commit is contained in:
+120
-35
@@ -1,10 +1,17 @@
|
|||||||
-module(archive_manager).
|
-module(archive_manager).
|
||||||
-behaviour(gen_server).
|
-behaviour(gen_server).
|
||||||
|
-compile([{nowarn_deprecated_function, [{slave, start, 3}]}]).
|
||||||
|
|
||||||
|
%% Peer start must not run inside handle_call: peer:start_it timeout
|
||||||
|
%% exits the gen_server and cascades via infra_sup (seen under IFT load).
|
||||||
|
%% Starts are async single-flight; callers wait or get {error, _}.
|
||||||
|
|
||||||
-export([start_link/0, get_archive_node/1]).
|
-export([start_link/0, get_archive_node/1]).
|
||||||
-export([init/1, handle_call/3, handle_cast/2, handle_info/2]).
|
-export([init/1, handle_call/3, handle_cast/2, handle_info/2,
|
||||||
|
terminate/2, code_change/3]).
|
||||||
|
|
||||||
-define(TIMEOUT, 30000). % 30 секунд неактивности
|
-define(IDLE_MS, 30000).
|
||||||
|
-define(PEER_CONN_MS, 15000).
|
||||||
|
|
||||||
start_link() ->
|
start_link() ->
|
||||||
gen_server:start_link({local, ?MODULE}, ?MODULE, [], []).
|
gen_server:start_link({local, ?MODULE}, ?MODULE, [], []).
|
||||||
@@ -13,46 +20,119 @@ get_archive_node(Day) ->
|
|||||||
gen_server:call(?MODULE, {get_node, Day}).
|
gen_server:call(?MODULE, {get_node, Day}).
|
||||||
|
|
||||||
init([]) ->
|
init([]) ->
|
||||||
{ok, #{}}.
|
process_flag(trap_exit, true),
|
||||||
|
{ok, #{nodes => #{}, starting => #{}}}.
|
||||||
|
|
||||||
handle_call({get_node, Day}, _From, State) ->
|
handle_call({get_node, Day}, From, State) ->
|
||||||
Node = archive_node_name(Day),
|
Node = archive_node_name(Day),
|
||||||
case maps:find(Node, State) of
|
Nodes = maps:get(nodes, State),
|
||||||
{ok, _Info} -> % переменная не используется, ок
|
Starting = maps:get(starting, State),
|
||||||
{reply, {ok, Node}, update_access(State, Node)};
|
case maps:find(Node, Nodes) of
|
||||||
|
{ok, _} ->
|
||||||
|
{reply, {ok, Node}, State#{nodes := touch_node(Nodes, Node)}};
|
||||||
error ->
|
error ->
|
||||||
case start_archive_node(Day) of
|
case is_node_alive(Node) of
|
||||||
{ok, Node} ->
|
|
||||||
Ref = erlang:send_after(?TIMEOUT, self(), {release, Node}),
|
|
||||||
NewState = State#{Node => #{timer => Ref, last_access => erlang:monotonic_time()}},
|
|
||||||
{reply, {ok, Node}, NewState};
|
|
||||||
{error, Reason} ->
|
|
||||||
{reply, {error, Reason}, State}
|
|
||||||
end
|
|
||||||
end.
|
|
||||||
|
|
||||||
handle_cast(_, State) -> {noreply, State}.
|
|
||||||
|
|
||||||
handle_info({release, Node}, State) ->
|
|
||||||
case maps:find(Node, State) of
|
|
||||||
{ok, #{last_access := Last}} ->
|
|
||||||
Now = erlang:monotonic_time(),
|
|
||||||
if Now - Last >= (?TIMEOUT * 1000) ->
|
|
||||||
stop_archive_node(Node),
|
|
||||||
{noreply, maps:remove(Node, State)};
|
|
||||||
true ->
|
true ->
|
||||||
Ref = erlang:send_after(?TIMEOUT, self(), {release, Node}),
|
{reply, {ok, Node}, State#{nodes := register_node(Nodes, Node)}};
|
||||||
{noreply, State#{Node => #{timer => Ref, last_access => Last}}}
|
false ->
|
||||||
|
case maps:find(Day, Starting) of
|
||||||
|
{ok, Waiters} ->
|
||||||
|
{noreply, State#{starting := Starting#{Day => [From | Waiters]}}};
|
||||||
|
error ->
|
||||||
|
spawn_starter(Day),
|
||||||
|
{noreply, State#{starting := Starting#{Day => [From]}}}
|
||||||
|
end
|
||||||
|
end
|
||||||
|
end;
|
||||||
|
handle_call(_Req, _From, State) ->
|
||||||
|
{reply, {error, unknown_call}, State}.
|
||||||
|
|
||||||
|
handle_cast(_, State) ->
|
||||||
|
{noreply, State}.
|
||||||
|
|
||||||
|
handle_info({start_done, Day, Result}, State) ->
|
||||||
|
Starting = maps:get(starting, State),
|
||||||
|
Waiters = maps:get(Day, Starting, []),
|
||||||
|
NewStarting = maps:remove(Day, Starting),
|
||||||
|
Nodes = maps:get(nodes, State),
|
||||||
|
case Result of
|
||||||
|
{ok, Node} ->
|
||||||
|
reply_all(Waiters, {ok, Node}),
|
||||||
|
{noreply, State#{
|
||||||
|
nodes := register_node(Nodes, Node),
|
||||||
|
starting := NewStarting
|
||||||
|
}};
|
||||||
|
{error, Reason} ->
|
||||||
|
reply_all(Waiters, {error, Reason}),
|
||||||
|
{noreply, State#{starting := NewStarting}}
|
||||||
|
end;
|
||||||
|
handle_info({release, Node}, State) ->
|
||||||
|
Nodes = maps:get(nodes, State),
|
||||||
|
case maps:find(Node, Nodes) of
|
||||||
|
{ok, #{last_access := Last, timer := _Old}} ->
|
||||||
|
Idle = erlang:convert_time_unit(
|
||||||
|
erlang:monotonic_time() - Last, native, millisecond),
|
||||||
|
if Idle >= ?IDLE_MS ->
|
||||||
|
stop_archive_node(Node),
|
||||||
|
{noreply, State#{nodes := maps:remove(Node, Nodes)}};
|
||||||
|
true ->
|
||||||
|
{noreply, State#{nodes := touch_node(Nodes, Node)}}
|
||||||
end;
|
end;
|
||||||
error ->
|
error ->
|
||||||
{noreply, State}
|
{noreply, State}
|
||||||
end;
|
end;
|
||||||
handle_info(_, State) -> {noreply, State}.
|
handle_info({'EXIT', _Pid, _Reason}, State) ->
|
||||||
|
{noreply, State};
|
||||||
|
handle_info(_, State) ->
|
||||||
|
{noreply, State}.
|
||||||
|
|
||||||
update_access(State, Node) ->
|
terminate(_Reason, _State) ->
|
||||||
Info = maps:get(Node, State),
|
ok.
|
||||||
Ref = erlang:send_after(?TIMEOUT, self(), {release, Node}),
|
|
||||||
State#{Node => Info#{timer => Ref, last_access => erlang:monotonic_time()}}.
|
code_change(_OldVsn, State, _Extra) ->
|
||||||
|
{ok, State}.
|
||||||
|
|
||||||
|
%%%-------------------------------------------------------------------
|
||||||
|
%%% Internal
|
||||||
|
%%%-------------------------------------------------------------------
|
||||||
|
|
||||||
|
spawn_starter(Day) ->
|
||||||
|
Parent = self(),
|
||||||
|
spawn(fun() ->
|
||||||
|
Result =
|
||||||
|
try start_archive_node(Day) of
|
||||||
|
Ok -> Ok
|
||||||
|
catch
|
||||||
|
exit:{timeout, _} -> {error, peer_timeout};
|
||||||
|
exit:Reason -> {error, {exit, Reason}};
|
||||||
|
error:Reason -> {error, {error, Reason}};
|
||||||
|
throw:Reason -> {error, {throw, Reason}}
|
||||||
|
end,
|
||||||
|
Parent ! {start_done, Day, Result}
|
||||||
|
end).
|
||||||
|
|
||||||
|
reply_all(Waiters, Reply) ->
|
||||||
|
lists:foreach(fun(From) -> gen_server:reply(From, Reply) end, Waiters).
|
||||||
|
|
||||||
|
register_node(Nodes, Node) ->
|
||||||
|
case maps:find(Node, Nodes) of
|
||||||
|
{ok, #{timer := OldRef}} ->
|
||||||
|
_ = erlang:cancel_timer(OldRef),
|
||||||
|
ok;
|
||||||
|
error ->
|
||||||
|
ok
|
||||||
|
end,
|
||||||
|
Ref = erlang:send_after(?IDLE_MS, self(), {release, Node}),
|
||||||
|
Nodes#{Node => #{timer => Ref, last_access => erlang:monotonic_time()}}.
|
||||||
|
|
||||||
|
touch_node(Nodes, Node) ->
|
||||||
|
register_node(Nodes, Node).
|
||||||
|
|
||||||
|
is_node_alive(Node) ->
|
||||||
|
case net_adm:ping(Node) of
|
||||||
|
pong -> true;
|
||||||
|
pang -> false
|
||||||
|
end.
|
||||||
|
|
||||||
start_archive_node(Day) ->
|
start_archive_node(Day) ->
|
||||||
Node = archive_node_name(Day),
|
Node = archive_node_name(Day),
|
||||||
@@ -69,7 +149,11 @@ start_archive_node(Day) ->
|
|||||||
start_archive_peer(Node) ->
|
start_archive_peer(Node) ->
|
||||||
case os:getenv("CLUSTER_MODE") of
|
case os:getenv("CLUSTER_MODE") of
|
||||||
"true" ->
|
"true" ->
|
||||||
case peer:start_link(#{name => Node, host => host()}) of
|
case peer:start_link(#{
|
||||||
|
name => Node,
|
||||||
|
host => host(),
|
||||||
|
connection_timeout => ?PEER_CONN_MS
|
||||||
|
}) of
|
||||||
{ok, _} -> {ok, Node};
|
{ok, _} -> {ok, Node};
|
||||||
{error, {already_started, _}} -> {ok, Node};
|
{error, {already_started, _}} -> {ok, Node};
|
||||||
Error -> Error
|
Error -> Error
|
||||||
@@ -78,6 +162,7 @@ start_archive_peer(Node) ->
|
|||||||
CookieStr = atom_to_list(erlang:get_cookie()),
|
CookieStr = atom_to_list(erlang:get_cookie()),
|
||||||
case slave:start(host(), Node, "-setcookie " ++ CookieStr) of
|
case slave:start(host(), Node, "-setcookie " ++ CookieStr) of
|
||||||
{ok, _} -> {ok, Node};
|
{ok, _} -> {ok, Node};
|
||||||
|
{error, {already_running, _}} -> {ok, Node};
|
||||||
Error -> Error
|
Error -> Error
|
||||||
end
|
end
|
||||||
end.
|
end.
|
||||||
@@ -102,4 +187,4 @@ stop_archive_node(Node) ->
|
|||||||
|
|
||||||
host() ->
|
host() ->
|
||||||
{ok, Name} = inet:gethostname(),
|
{ok, Name} = inet:gethostname(),
|
||||||
Name.
|
Name.
|
||||||
|
|||||||
@@ -14,7 +14,7 @@
|
|||||||
|
|
||||||
-include("records.hrl").
|
-include("records.hrl").
|
||||||
|
|
||||||
-define(ARCHIVE_CALL_TIMEOUT, 3000).
|
-define(ARCHIVE_CALL_TIMEOUT, 8000).
|
||||||
|
|
||||||
%%% cowboy_handler callback
|
%%% cowboy_handler callback
|
||||||
-spec init(cowboy_req:req(), any()) -> {ok, cowboy_req:req(), any()}.
|
-spec init(cowboy_req:req(), any()) -> {ok, cowboy_req:req(), any()}.
|
||||||
|
|||||||
Reference in New Issue
Block a user