191 lines
5.3 KiB
Erlang
191 lines
5.3 KiB
Erlang
-module(archive_manager).
|
|
-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([init/1, handle_call/3, handle_cast/2, handle_info/2,
|
|
terminate/2, code_change/3]).
|
|
|
|
-define(IDLE_MS, 30000).
|
|
-define(PEER_CONN_MS, 15000).
|
|
|
|
start_link() ->
|
|
gen_server:start_link({local, ?MODULE}, ?MODULE, [], []).
|
|
|
|
get_archive_node(Day) ->
|
|
gen_server:call(?MODULE, {get_node, Day}).
|
|
|
|
init([]) ->
|
|
process_flag(trap_exit, true),
|
|
{ok, #{nodes => #{}, starting => #{}}}.
|
|
|
|
handle_call({get_node, Day}, From, State) ->
|
|
Node = archive_node_name(Day),
|
|
Nodes = maps:get(nodes, State),
|
|
Starting = maps:get(starting, State),
|
|
case maps:find(Node, Nodes) of
|
|
{ok, _} ->
|
|
{reply, {ok, Node}, State#{nodes := touch_node(Nodes, Node)}};
|
|
error ->
|
|
case is_node_alive(Node) of
|
|
true ->
|
|
{reply, {ok, Node}, State#{nodes := register_node(Nodes, Node)}};
|
|
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;
|
|
error ->
|
|
{noreply, State}
|
|
end;
|
|
handle_info({'EXIT', _Pid, _Reason}, State) ->
|
|
{noreply, State};
|
|
handle_info(_, State) ->
|
|
{noreply, State}.
|
|
|
|
terminate(_Reason, _State) ->
|
|
ok.
|
|
|
|
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) ->
|
|
Node = archive_node_name(Day),
|
|
case start_archive_peer(Node) of
|
|
{ok, _} ->
|
|
case ensure_archive_node_ready(Node) of
|
|
ok -> {ok, Node};
|
|
{error, Reason} -> {error, Reason}
|
|
end;
|
|
{error, Reason} ->
|
|
{error, Reason}
|
|
end.
|
|
|
|
start_archive_peer(Node) ->
|
|
case os:getenv("CLUSTER_MODE") of
|
|
"true" ->
|
|
case peer:start_link(#{
|
|
name => Node,
|
|
host => host(),
|
|
connection_timeout => ?PEER_CONN_MS
|
|
}) of
|
|
{ok, _} -> {ok, Node};
|
|
{error, {already_started, _}} -> {ok, Node};
|
|
Error -> Error
|
|
end;
|
|
_ ->
|
|
CookieStr = atom_to_list(erlang:get_cookie()),
|
|
case slave:start(host(), Node, "-setcookie " ++ CookieStr) of
|
|
{ok, _} -> {ok, Node};
|
|
{error, {already_running, _}} -> {ok, Node};
|
|
Error -> Error
|
|
end
|
|
end.
|
|
|
|
ensure_archive_node_ready(Node) ->
|
|
case rpc:call(Node, mnesia, start, [], 5000) of
|
|
ok ->
|
|
case rpc:call(Node, code, ensure_loaded, [archive_fetcher], 5000) of
|
|
{module, archive_fetcher} -> ok;
|
|
{error, Reason} -> {error, Reason};
|
|
{badrpc, Reason} -> {error, Reason}
|
|
end;
|
|
{badrpc, Reason} -> {error, Reason};
|
|
Other -> {error, Other}
|
|
end.
|
|
|
|
archive_node_name(Day) ->
|
|
list_to_atom("eventhub_archive_" ++ Day ++ "@" ++ host()).
|
|
|
|
stop_archive_node(Node) ->
|
|
rpc:cast(Node, init, stop, []).
|
|
|
|
host() ->
|
|
{ok, Name} = inet:gethostname(),
|
|
Name.
|