feat(gateway): add an undrain endpoint and sweep orphan tables (#2493)

This commit is contained in:
Hampus
2026-09-06 15:25:43 +02:00
committed by GitHub
parent e828398e06
commit d8f2aa3184
4 changed files with 104 additions and 8 deletions
@@ -45,6 +45,7 @@ start_cowboy() ->
{<<"/_health">>, health_handler, liveness},
{<<"/_health/ready">>, health_handler, readiness},
{<<"/_health/drain">>, health_handler, drain},
{<<"/_health/undrain">>, health_handler, undrain},
{<<"/_metrics">>, metrics_handler, []},
{<<"/">>, gateway_handler, []}
]}
@@ -17,6 +17,7 @@
-define(DEBOUNCE_MS, 2000).
-define(RECONCILE_MS, 60000).
-define(HANDOFF_WORKER_TIMEOUT_MS, 30000).
-define(UNDRAIN_CALL_TIMEOUT_MS, 5000).
-type timer_state() :: undefined | {reference(), reference()}.
-type reconcile_timer() :: undefined | reference().
@@ -50,8 +51,17 @@ drain_async() ->
ok = drain_notify_role(),
drain_dispatch().
-spec undrain() -> ok.
-spec undrain() -> ok | {error, handoff_in_flight | unavailable}.
undrain() ->
case whereis(?MODULE) of
undefined ->
clear_draining();
_Pid ->
shard_utils:safe_gen_call(?MODULE, undrain, ?UNDRAIN_CALL_TIMEOUT_MS)
end.
-spec clear_draining() -> ok.
clear_draining() ->
persistent_term:erase({fluxer_gateway, draining}),
logger:info("Gateway un-cordoned: draining flag cleared"),
ok.
@@ -89,6 +99,10 @@ init([]) ->
{reply, term(), state()}.
handle_call(diagnostic_info, _From, State) ->
{reply, info(State), State};
handle_call(undrain, _From, #{handoff := undefined} = State) ->
{reply, clear_draining(), State};
handle_call(undrain, _From, State) ->
{reply, {error, handoff_in_flight}, State};
handle_call(_Request, _From, State) ->
{reply, ok, State}.
@@ -615,6 +629,18 @@ undrain_clears_draining_flag_test() ->
?assertEqual(false, persistent_term:get({fluxer_gateway, draining}, false)),
restore_persistent_term({fluxer_gateway, draining}, Previous).
undrain_is_refused_while_handoff_in_flight_test() ->
Previous = persistent_term:get({fluxer_gateway, draining}, undefined),
persistent_term:put({fluxer_gateway, draining}, true),
Busy = (idle_state([node()]))#{handoff := {self(), make_ref(), [node()], #{}}},
{reply, Refused, Busy} = handle_call(undrain, {self(), make_ref()}, Busy),
?assertEqual({error, handoff_in_flight}, Refused),
?assertEqual(true, persistent_term:get({fluxer_gateway, draining}, false)),
{reply, Cleared, _Idle} = handle_call(undrain, {self(), make_ref()}, idle_state([node()])),
?assertEqual(ok, Cleared),
?assertEqual(false, persistent_term:get({fluxer_gateway, draining}, false)),
restore_persistent_term({fluxer_gateway, draining}, Previous).
handoff_wait_loop() ->
receive
stop -> ok
+43 -4
View File
@@ -5,7 +5,7 @@
-export([init/2]).
-type mode() :: liveness | readiness | drain.
-type mode() :: liveness | readiness | drain | undrain.
-spec init(cowboy_req:req(), term()) -> {ok, cowboy_req:req(), term()}.
init(Req0, Mode0) ->
@@ -21,16 +21,19 @@ init(Req0, Mode0) ->
-spec normalize_mode(term()) -> mode().
normalize_mode(drain) -> drain;
normalize_mode(undrain) -> undrain;
normalize_mode(readiness) -> readiness;
normalize_mode(_) -> liveness.
-spec response_for_mode(mode(), cowboy_req:req()) -> {200 | 403 | 503, binary()}.
-spec response_for_mode(mode(), cowboy_req:req()) -> {200 | 403 | 409 | 503, binary()}.
response_for_mode(liveness, _Req) ->
{200, <<"OK">>};
response_for_mode(readiness, Req) ->
readiness_response(Req);
response_for_mode(drain, Req) ->
drain_response(Req).
drain_response(Req);
response_for_mode(undrain, Req) ->
undrain_response(Req).
-spec readiness_response(cowboy_req:req()) -> {200 | 403 | 503, binary()}.
readiness_response(Req) ->
@@ -57,10 +60,31 @@ drain_response(Req) ->
{200, <<"DRAINING">>}
end.
-spec undrain_response(cowboy_req:req()) -> {200 | 403 | 409, binary()}.
undrain_response(Req) ->
case is_loopback_request(Req) of
false ->
{403, <<"FORBIDDEN">>};
true ->
undrain_status(deactivate_drain())
end.
-spec undrain_status(term()) -> {200 | 409, binary()}.
undrain_status(ok) ->
{200, <<"READY">>};
undrain_status({error, handoff_in_flight}) ->
{409, <<"HANDOFF_IN_FLIGHT">>};
undrain_status(_Error) ->
{409, <<"UNAVAILABLE">>}.
-spec activate_drain() -> ok.
activate_drain() ->
gateway_cluster_handoff:drain_async().
-spec deactivate_drain() -> ok | {error, term()}.
deactivate_drain() ->
gateway_cluster_handoff:undrain().
-spec is_loopback_request(cowboy_req:req()) -> boolean().
is_loopback_request(Req) ->
case cowboy_req:peer(Req) of
@@ -79,16 +103,31 @@ normalize_mode_test() ->
?assertEqual(liveness, normalize_mode(undefined)),
?assertEqual(liveness, normalize_mode([])),
?assertEqual(readiness, normalize_mode(readiness)),
?assertEqual(drain, normalize_mode(drain)).
?assertEqual(drain, normalize_mode(drain)),
?assertEqual(undrain, normalize_mode(undrain)).
readiness_status_test() ->
?assertEqual({200, <<"OK">>}, readiness_status(true)),
?assertEqual({503, <<"DRAINING">>}, readiness_status(false)).
undrain_status_test() ->
?assertEqual({200, <<"READY">>}, undrain_status(ok)),
?assertEqual({409, <<"HANDOFF_IN_FLIGHT">>}, undrain_status({error, handoff_in_flight})),
?assertEqual({409, <<"UNAVAILABLE">>}, undrain_status({error, unavailable})).
activate_drain_sets_draining_flag_test() ->
persistent_term:erase({fluxer_gateway, draining}),
?assertEqual(ok, activate_drain()),
?assert(gateway_node_router:is_draining()),
persistent_term:erase({fluxer_gateway, draining}).
deactivate_drain_clears_draining_flag_test() ->
persistent_term:erase({fluxer_gateway, draining}),
?assertEqual(ok, activate_drain()),
?assert(gateway_node_router:is_draining()),
?assertEqual(ok, deactivate_drain()),
?assertNot(gateway_node_router:is_draining()),
?assertEqual({200, <<"OK">>}, readiness_status(gateway_node_router:is_ready())),
persistent_term:erase({fluxer_gateway, draining}).
-endif.
+33 -3
View File
@@ -8,8 +8,9 @@
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2, code_change/3]).
-define(DEFAULT_CALL_TIMEOUT_MS, 5000).
-define(ORPHAN_SWEEP_INTERVAL_MS, 60000).
-type state() :: #{}.
-type state() :: #{sweep_timer => reference()}.
-spec start_link() -> {ok, pid()} | {error, term()}.
start_link() ->
@@ -32,7 +33,7 @@ ensure_table(TableName, Options) ->
init([]) ->
erlang:process_flag(fullsweep_after, 0),
ok = ensure_core_tables(),
{ok, #{}}.
{ok, #{sweep_timer => schedule_orphan_sweep()}}.
-spec handle_call(term(), gen_server:from(), state()) -> {reply, ok, state()}.
handle_call({ensure_table, TableName, Options}, _From, State) when
@@ -47,11 +48,26 @@ handle_cast(_Msg, State) ->
{noreply, State}.
-spec handle_info(term(), state()) -> {noreply, state()}.
handle_info(sweep_orphan_tables, State) ->
_ = sweep_orphan_tables(),
{noreply, State#{sweep_timer => schedule_orphan_sweep()}};
handle_info(_Info, State) ->
{noreply, State}.
-spec schedule_orphan_sweep() -> reference().
schedule_orphan_sweep() ->
erlang:send_after(?ORPHAN_SWEEP_INTERVAL_MS, self(), sweep_orphan_tables).
-spec terminate(term(), state()) -> ok.
terminate(_Reason, _State) ->
terminate(_Reason, State) ->
cancel_orphan_sweep(State),
ok.
-spec cancel_orphan_sweep(state()) -> ok.
cancel_orphan_sweep(#{sweep_timer := Ref}) when is_reference(Ref) ->
_ = erlang:cancel_timer(Ref),
ok;
cancel_orphan_sweep(_State) ->
ok.
-spec code_change(term(), state(), term()) -> {ok, state()}.
@@ -229,6 +245,20 @@ sweep_orphan_tables_detects_dead_owner_test() ->
?assert(Count >= 0),
ets:delete(Tab).
orphan_sweep_is_scheduled_and_rearmed_test() ->
{ok, Pid} = start_link(),
try
#{sweep_timer := Ref} = sys:get_state(Pid),
?assert(is_reference(Ref)),
Pid ! sweep_orphan_tables,
#{sweep_timer := NextRef} = sys:get_state(Pid),
?assert(is_reference(NextRef)),
?assertNotEqual(Ref, NextRef),
?assertEqual(Pid, ets:info(guild_voice_registry, owner))
after
gen_server:stop(?MODULE)
end.
sweep_orphan_tables_returns_zero_when_clean_test() ->
{ok, Count} = sweep_orphan_tables(),
?assert(is_integer(Count)),