mirror of
https://github.com/fluxerapp/fluxer
synced 2026-10-07 19:22:14 +09:00
fix(gateway): keep dispatch ordered under broadcaster load (#2627)
This commit is contained in:
@@ -11,6 +11,8 @@
|
|||||||
]).
|
]).
|
||||||
|
|
||||||
-define(MAX_BULK_ENCODE_GROUPS, 1024).
|
-define(MAX_BULK_ENCODE_GROUPS, 1024).
|
||||||
|
-define(BROADCASTER_CLAIM_DEADLINE_MS, 250).
|
||||||
|
-define(BROADCASTER_CLAIM_BACKOFF_MS, 1).
|
||||||
|
|
||||||
-type event() :: atom().
|
-type event() :: atom().
|
||||||
-type event_data() :: map().
|
-type event_data() :: map().
|
||||||
@@ -240,12 +242,41 @@ check_eligible_pid(SessionData, Event, FinalData, GuildId, State) ->
|
|||||||
dispatch_to_pids([], _Event, _EncodedData, _GuildId, _State) ->
|
dispatch_to_pids([], _Event, _EncodedData, _GuildId, _State) ->
|
||||||
ok;
|
ok;
|
||||||
dispatch_to_pids(Pids, Event, EncodedData, GuildId, State) ->
|
dispatch_to_pids(Pids, Event, EncodedData, GuildId, State) ->
|
||||||
BroadcasterPid = maps:get(broadcaster_pid, State, undefined),
|
case maps:get(broadcaster_pid, State, undefined) of
|
||||||
case guild_broadcaster:cast_event(BroadcasterPid, Event, EncodedData, Pids) of
|
BroadcasterPid when is_pid(BroadcasterPid) ->
|
||||||
true -> ok;
|
claim_broadcaster(BroadcasterPid, Pids, Event, EncodedData, claim_deadline());
|
||||||
false -> gateway_dispatch_relay:dispatch_many(Pids, Event, EncodedData, GuildId)
|
_ ->
|
||||||
|
gateway_dispatch_relay:dispatch_many(Pids, Event, EncodedData, GuildId)
|
||||||
end.
|
end.
|
||||||
|
|
||||||
|
-spec claim_broadcaster(pid(), [pid()], event(), term(), integer()) -> ok.
|
||||||
|
claim_broadcaster(BroadcasterPid, Pids, Event, EncodedData, Deadline) ->
|
||||||
|
case guild_broadcaster:cast_event(BroadcasterPid, Event, EncodedData, Pids) of
|
||||||
|
true ->
|
||||||
|
ok;
|
||||||
|
false ->
|
||||||
|
claim_broadcaster_after_wait(
|
||||||
|
BroadcasterPid,
|
||||||
|
Pids,
|
||||||
|
Event,
|
||||||
|
EncodedData,
|
||||||
|
Deadline,
|
||||||
|
gateway_retry_timer:wait_until(?BROADCASTER_CLAIM_BACKOFF_MS, Deadline)
|
||||||
|
)
|
||||||
|
end.
|
||||||
|
|
||||||
|
-spec claim_broadcaster_after_wait(
|
||||||
|
pid(), [pid()], event(), term(), integer(), term()
|
||||||
|
) -> ok.
|
||||||
|
claim_broadcaster_after_wait(BroadcasterPid, Pids, Event, EncodedData, Deadline, ok) ->
|
||||||
|
claim_broadcaster(BroadcasterPid, Pids, Event, EncodedData, Deadline);
|
||||||
|
claim_broadcaster_after_wait(BroadcasterPid, Pids, Event, EncodedData, _Deadline, _Expired) ->
|
||||||
|
gen_server:cast(BroadcasterPid, {event_broadcast, Event, EncodedData, Pids}).
|
||||||
|
|
||||||
|
-spec claim_deadline() -> integer().
|
||||||
|
claim_deadline() ->
|
||||||
|
erlang:monotonic_time(millisecond) + ?BROADCASTER_CLAIM_DEADLINE_MS.
|
||||||
|
|
||||||
-spec normalize_success(non_neg_integer()) -> non_neg_integer().
|
-spec normalize_success(non_neg_integer()) -> non_neg_integer().
|
||||||
normalize_success(Count) when Count > 0 -> 1;
|
normalize_success(Count) when Count > 0 -> 1;
|
||||||
normalize_success(_) -> 0.
|
normalize_success(_) -> 0.
|
||||||
|
|||||||
@@ -0,0 +1,132 @@
|
|||||||
|
%% SPDX-License-Identifier: AGPL-3.0-or-later
|
||||||
|
|
||||||
|
-module(guild_dispatch_send_ordering_tests).
|
||||||
|
-typing([eqwalizer]).
|
||||||
|
|
||||||
|
-include_lib("eunit/include/eunit.hrl").
|
||||||
|
|
||||||
|
-define(SATURATION_FILL, 600).
|
||||||
|
-define(RECEIVE_TIMEOUT_MS, 600).
|
||||||
|
-define(BROADCASTER_IDLE_TIMEOUT_MS, 30000).
|
||||||
|
|
||||||
|
saturated_broadcaster_keeps_dispatch_ordered_test_() ->
|
||||||
|
{timeout, 30, fun saturated_broadcaster_keeps_dispatch_ordered/0}.
|
||||||
|
|
||||||
|
healthy_broadcaster_keeps_dispatch_ordered_test_() ->
|
||||||
|
{timeout, 30, fun healthy_broadcaster_keeps_dispatch_ordered/0}.
|
||||||
|
|
||||||
|
absent_broadcaster_dispatches_through_relay_test_() ->
|
||||||
|
{timeout, 30, fun absent_broadcaster_dispatches_through_relay/0}.
|
||||||
|
|
||||||
|
saturated_broadcaster_keeps_dispatch_ordered() ->
|
||||||
|
Broadcaster = spawn_idle_broadcaster(),
|
||||||
|
try
|
||||||
|
saturate(Broadcaster),
|
||||||
|
?assertEqual(false, cast_event(Broadcaster)),
|
||||||
|
QueueBefore = queue_len(Broadcaster),
|
||||||
|
flush_dispatches(),
|
||||||
|
dispatch(Broadcaster),
|
||||||
|
?assertNot(received_relay_dispatch()),
|
||||||
|
?assertEqual(QueueBefore + 1, queue_len(Broadcaster))
|
||||||
|
after
|
||||||
|
stop_broadcaster(Broadcaster)
|
||||||
|
end.
|
||||||
|
|
||||||
|
healthy_broadcaster_keeps_dispatch_ordered() ->
|
||||||
|
Broadcaster = spawn_draining_broadcaster(),
|
||||||
|
try
|
||||||
|
flush_dispatches(),
|
||||||
|
dispatch(Broadcaster),
|
||||||
|
?assertNot(received_relay_dispatch())
|
||||||
|
after
|
||||||
|
stop_broadcaster(Broadcaster)
|
||||||
|
end.
|
||||||
|
|
||||||
|
absent_broadcaster_dispatches_through_relay() ->
|
||||||
|
flush_dispatches(),
|
||||||
|
dispatch(undefined),
|
||||||
|
?assert(received_relay_dispatch()).
|
||||||
|
|
||||||
|
dispatch(Broadcaster) ->
|
||||||
|
_ = guild_dispatch_send:dispatch_to_sessions(
|
||||||
|
[session()], message_update, event_data(), state(Broadcaster)
|
||||||
|
),
|
||||||
|
ok.
|
||||||
|
|
||||||
|
cast_event(Broadcaster) ->
|
||||||
|
guild_broadcaster:cast_event(Broadcaster, message_update, {pre_encoded, <<"{}">>}, [self()]).
|
||||||
|
|
||||||
|
saturate(Broadcaster) ->
|
||||||
|
lists:foreach(
|
||||||
|
fun(N) -> Broadcaster ! {filler, N} end,
|
||||||
|
lists:seq(1, ?SATURATION_FILL)
|
||||||
|
).
|
||||||
|
|
||||||
|
spawn_idle_broadcaster() ->
|
||||||
|
spawn(fun idle_broadcaster/0).
|
||||||
|
|
||||||
|
idle_broadcaster() ->
|
||||||
|
receive
|
||||||
|
stop -> ok
|
||||||
|
after ?BROADCASTER_IDLE_TIMEOUT_MS -> ok
|
||||||
|
end.
|
||||||
|
|
||||||
|
spawn_draining_broadcaster() ->
|
||||||
|
spawn(fun draining_broadcaster/0).
|
||||||
|
|
||||||
|
draining_broadcaster() ->
|
||||||
|
receive
|
||||||
|
stop -> ok;
|
||||||
|
_Other -> draining_broadcaster()
|
||||||
|
after ?BROADCASTER_IDLE_TIMEOUT_MS -> ok
|
||||||
|
end.
|
||||||
|
|
||||||
|
stop_broadcaster(Broadcaster) ->
|
||||||
|
Broadcaster ! stop,
|
||||||
|
ok.
|
||||||
|
|
||||||
|
queue_len(Pid) ->
|
||||||
|
{message_queue_len, Len} = erlang:process_info(Pid, message_queue_len),
|
||||||
|
Len.
|
||||||
|
|
||||||
|
received_relay_dispatch() ->
|
||||||
|
receive
|
||||||
|
{'$gen_cast', {dispatch, message_update, _Payload}} -> true
|
||||||
|
after ?RECEIVE_TIMEOUT_MS -> false
|
||||||
|
end.
|
||||||
|
|
||||||
|
flush_dispatches() ->
|
||||||
|
receive
|
||||||
|
{'$gen_cast', _Cast} -> flush_dispatches()
|
||||||
|
after 0 -> ok
|
||||||
|
end.
|
||||||
|
|
||||||
|
state(undefined) ->
|
||||||
|
base_state();
|
||||||
|
state(Broadcaster) ->
|
||||||
|
(base_state())#{broadcaster_pid => Broadcaster}.
|
||||||
|
|
||||||
|
base_state() ->
|
||||||
|
#{
|
||||||
|
id => 42,
|
||||||
|
member_count => 100,
|
||||||
|
data => #{
|
||||||
|
<<"guild">> => #{<<"owner_id">> => <<"999">>},
|
||||||
|
<<"roles">> => [],
|
||||||
|
<<"members">> => [],
|
||||||
|
<<"channels">> => []
|
||||||
|
}
|
||||||
|
}.
|
||||||
|
|
||||||
|
session() ->
|
||||||
|
{<<"ordered">>, #{
|
||||||
|
session_id => <<"ordered">>,
|
||||||
|
user_id => 10,
|
||||||
|
pid => self(),
|
||||||
|
active_guilds => sets:new(),
|
||||||
|
bot => false,
|
||||||
|
viewable_channels => #{100 => true}
|
||||||
|
}}.
|
||||||
|
|
||||||
|
event_data() ->
|
||||||
|
#{<<"guild_id">> => <<"42">>, <<"id">> => <<"7">>, <<"channel_id">> => <<"100">>}.
|
||||||
Reference in New Issue
Block a user