fix(gateway): share the relay dispatch bound across producers (#2315)

This commit is contained in:
Hampus
2026-09-01 04:34:41 +02:00
committed by GitHub
parent cf3af50464
commit 90aa810ce4
5 changed files with 252 additions and 71 deletions
@@ -21,9 +21,13 @@
-define(STATE_KEY, {gateway_dispatch_relay, state}).
-define(WORKER_KEY(Index), {gateway_dispatch_relay_worker, Index}).
-type state() ::
coordinator_state()
| #{role := worker, index := non_neg_integer(), delivered := non_neg_integer()}.
-type state() :: coordinator_state() | worker_state().
-type worker_state() :: #{
role := worker,
index := non_neg_integer(),
inflight := atomics:atomics_ref() | undefined,
delivered := non_neg_integer()
}.
-type coordinator_state() :: #{
role := coordinator,
workers := tuple(),
@@ -163,23 +167,14 @@ init(coordinator) ->
init({worker, Index}) ->
erlang:process_flag(fullsweep_after, 50),
persistent_term:put(?WORKER_KEY(Index), self()),
{ok, #{role => worker, index => Index, delivered => 0}}.
{ok, #{
role => worker,
index => Index,
inflight => gateway_dispatch_relay_batch:inflight_ref(Index),
delivered => 0
}}.
-spec handle_call(term(), gen_server:from(), state()) -> {reply, term(), state()}.
handle_call(
{deliver, SessionPid, Event, Payload},
_From,
#{role := worker, delivered := Delivered} = State
) ->
dispatch_direct(SessionPid, Event, Payload),
{reply, ok, State#{delivered := Delivered + 1}};
handle_call(
{deliver_many, SessionPids, Event, Payload},
_From,
#{role := worker, delivered := Delivered} = State
) when is_list(SessionPids) ->
deliver_many_direct(SessionPids, Event, Payload),
{reply, ok, State#{delivered := Delivered + length(SessionPids)}};
handle_call(diagnostic_info, _From, State) ->
{reply, diagnostic_info(), State};
handle_call(_Request, _From, State) ->
@@ -188,14 +183,16 @@ handle_call(_Request, _From, State) ->
-spec handle_cast(term(), state()) -> {noreply, state()}.
handle_cast(
{deliver, SessionPid, Event, Payload},
#{role := worker, delivered := Delivered} = State
#{role := worker, inflight := Inflight, delivered := Delivered} = State
) ->
ok = gateway_dispatch_relay_batch:release_queue_slot(Inflight),
dispatch_direct(SessionPid, Event, Payload),
{noreply, State#{delivered := Delivered + 1}};
handle_cast(
{deliver_many, SessionPids, Event, Payload},
#{role := worker, delivered := Delivered} = State
#{role := worker, inflight := Inflight, delivered := Delivered} = State
) when is_list(SessionPids) ->
ok = gateway_dispatch_relay_batch:release_queue_slot(Inflight),
deliver_many_direct(SessionPids, Event, Payload),
{noreply, State#{delivered := Delivered + length(SessionPids)}};
handle_cast(_Msg, State) ->
@@ -13,14 +13,19 @@
normalize_workers_tuple/1,
message_queue_len/1,
max_queue/0,
inflight_ref/1,
release_queue_slot/1,
start_workers/1,
start_worker/1,
worker_index/3
]).
-define(STATE_KEY, {gateway_dispatch_relay, state}).
-define(SYNC_TIMEOUT_MS, 5000).
-define(SLOT_KEY(Worker), {?MODULE, Worker}).
-define(INFLIGHT_KEY(Slot), {?MODULE, inflight, Slot}).
-define(INFLIGHT_INDEX, 1).
-define(INFLIGHT_UNSAMPLED, -1).
-define(CLAIM_DEADLINE_MS, 5000).
-define(CLAIM_BACKOFF_MS, 1).
-spec relay_or_direct_many([pid()], atom(), term()) -> ok.
relay_or_direct_many(SessionPids, Event, Payload) ->
@@ -42,7 +47,7 @@ relay_or_direct_many(SessionPids, Event, Payload) ->
relay_many_to_shards(SessionPids, Event, Payload, Workers, Count) ->
ShardBuckets = build_shard_buckets(SessionPids, Count),
Deferred = deliver_shard_buckets(1, Count, ShardBuckets, Event, Payload, Workers, []),
deliver_deferred_shards(Deferred, Event, Payload).
deliver_deferred_shards(Deferred, Event, Payload, Workers).
-spec build_shard_buckets([pid()], pos_integer()) -> tuple().
build_shard_buckets(SessionPids, Count) ->
@@ -59,8 +64,8 @@ build_shard_buckets(SessionPids, Count) ->
).
-spec deliver_shard_buckets(
pos_integer(), pos_integer(), tuple(), atom(), term(), tuple(), [{pid(), [pid()]}]
) -> [{pid(), [pid()]}].
pos_integer(), pos_integer(), tuple(), atom(), term(), tuple(), [{pos_integer(), [pid()]}]
) -> [{pos_integer(), [pid()]}].
deliver_shard_buckets(Index, Count, _Buckets, _Event, _Payload, _Workers, Deferred) when
Index > Count
->
@@ -73,74 +78,122 @@ deliver_shard_buckets(Index, Count, Buckets, Event, Payload, Workers, Deferred)
end,
deliver_shard_buckets(Index + 1, Count, Buckets, Event, Payload, Workers, Next).
-spec deliver_shard(pos_integer(), [pid()], atom(), term(), tuple(), [{pid(), [pid()]}]) ->
[{pid(), [pid()]}].
-spec deliver_shard(
pos_integer(), [pid()], atom(), term(), tuple(), [{pos_integer(), [pid()]}]
) -> [{pos_integer(), [pid()]}].
deliver_shard(Index, Pids, Event, Payload, Workers, Deferred) ->
Worker = element(Index, Workers),
case enqueue_async(Worker, {deliver_many, Pids, Event, Payload}) of
Msg = {deliver_many, Pids, Event, Payload},
case enqueue_async(Index - 1, element(Index, Workers), Msg) of
ok -> Deferred;
full -> [{Worker, Pids} | Deferred]
full -> [{Index, Pids} | Deferred]
end.
-spec deliver_deferred_shards([{pid(), [pid()]}], atom(), term()) -> ok.
deliver_deferred_shards([], _Event, _Payload) ->
-spec deliver_deferred_shards([{pos_integer(), [pid()]}], atom(), term(), tuple()) -> ok.
deliver_deferred_shards([], _Event, _Payload, _Workers) ->
ok;
deliver_deferred_shards([{Worker, Pids} | Rest], Event, Payload) ->
enqueue_sync(Worker, {deliver_many, Pids, Event, Payload}),
deliver_deferred_shards(Rest, Event, Payload).
deliver_deferred_shards([{Index, Pids} | Rest], Event, Payload, Workers) ->
enqueue(Index - 1, element(Index, Workers), {deliver_many, Pids, Event, Payload}),
deliver_deferred_shards(Rest, Event, Payload, Workers).
-spec relay_or_direct(pid(), atom(), term()) -> ok.
relay_or_direct(SessionPid, Event, Payload) ->
case select_worker(SessionPid) of
case select_worker_slot(SessionPid) of
undefined ->
gateway_dispatch_relay:dispatch_direct(SessionPid, Event, Payload);
Worker ->
enqueue(Worker, {deliver, SessionPid, Event, Payload})
{Slot, Worker} ->
enqueue(Slot, Worker, {deliver, SessionPid, Event, Payload})
end.
-spec enqueue(pid(), term()) -> ok.
enqueue(Worker, Msg) ->
case enqueue_async(Worker, Msg) of
-spec enqueue(non_neg_integer(), pid(), term()) -> ok.
enqueue(Slot, Worker, Msg) ->
case enqueue_async(Slot, Worker, Msg) of
ok -> ok;
full -> enqueue_sync(Worker, Msg)
full -> enqueue_blocking(Slot, Worker, Msg, claim_deadline())
end.
-spec enqueue_async(pid(), term()) -> ok | full.
enqueue_async(Worker, Msg) ->
case claim_queue_slot(Worker, max_queue()) of
-spec enqueue_blocking(non_neg_integer(), pid(), term(), integer()) -> ok.
enqueue_blocking(Slot, Worker, Msg, Deadline) ->
case gateway_retry_timer:wait_until(?CLAIM_BACKOFF_MS, Deadline) of
ok -> enqueue_retry(Slot, Worker, Msg, Deadline);
_ -> enqueue_forced(Slot, Worker, Msg)
end.
-spec enqueue_retry(non_neg_integer(), pid(), term(), integer()) -> ok.
enqueue_retry(Slot, Worker, Msg, Deadline) ->
case enqueue_async(Slot, Worker, Msg) of
ok -> ok;
full -> enqueue_blocking(Slot, Worker, Msg, Deadline)
end.
-spec enqueue_forced(non_neg_integer(), pid(), term()) -> ok.
enqueue_forced(Slot, Worker, Msg) ->
ok = reserve_queue_slot(Slot),
gen_server:cast(Worker, Msg).
-spec claim_deadline() -> integer().
claim_deadline() ->
erlang:monotonic_time(millisecond) + ?CLAIM_DEADLINE_MS.
-spec enqueue_async(non_neg_integer(), pid(), term()) -> ok | full.
enqueue_async(Slot, Worker, Msg) ->
case claim_queue_slot(Slot, Worker, max_queue()) of
ok -> gen_server:cast(Worker, Msg);
full -> full
end.
-spec claim_queue_slot(pid(), pos_integer()) -> ok | full.
claim_queue_slot(Worker, MaxQueue) ->
case erlang:get(?SLOT_KEY(Worker)) of
Claimed when is_integer(Claimed), Claimed < MaxQueue ->
_ = erlang:put(?SLOT_KEY(Worker), Claimed + 1),
ok;
_ ->
sample_queue_slot(Worker, MaxQueue)
-spec claim_queue_slot(non_neg_integer(), pid(), pos_integer()) -> ok | full.
claim_queue_slot(Slot, Worker, MaxQueue) ->
case inflight_ref(Slot) of
undefined -> sample_queue_slot(Worker, MaxQueue);
Ref -> claim_inflight_slot(Ref, Worker, MaxQueue)
end.
-spec claim_inflight_slot(atomics:atomics_ref(), pid(), pos_integer()) -> ok | full.
claim_inflight_slot(Ref, Worker, MaxQueue) ->
case atomics:add_get(Ref, ?INFLIGHT_INDEX, 1) of
Claimed when Claimed > 0, Claimed =< MaxQueue -> ok;
_ -> resample_inflight_slot(Ref, Worker, MaxQueue)
end.
-spec resample_inflight_slot(atomics:atomics_ref(), pid(), pos_integer()) -> ok | full.
resample_inflight_slot(Ref, Worker, MaxQueue) ->
Current = atomics:sub_get(Ref, ?INFLIGHT_INDEX, 1),
case message_queue_len(Worker) of
Sampled when Sampled < MaxQueue -> exchange_inflight_slot(Ref, Current, Sampled + 1);
_ -> full
end.
-spec exchange_inflight_slot(atomics:atomics_ref(), integer(), pos_integer()) -> ok | full.
exchange_inflight_slot(Ref, Current, Desired) ->
case atomics:compare_exchange(Ref, ?INFLIGHT_INDEX, Current, Desired) of
ok -> ok;
_ -> full
end.
-spec sample_queue_slot(pid(), pos_integer()) -> ok | full.
sample_queue_slot(Worker, MaxQueue) ->
case message_queue_len(Worker) of
Sampled when Sampled < MaxQueue ->
_ = erlang:put(?SLOT_KEY(Worker), Sampled + 1),
ok;
_ ->
_ = erlang:erase(?SLOT_KEY(Worker)),
full
case message_queue_len(Worker) < MaxQueue of
true -> ok;
false -> full
end.
-spec enqueue_sync(pid(), term()) -> ok.
enqueue_sync(Worker, Msg) ->
try gen_server:call(Worker, Msg, ?SYNC_TIMEOUT_MS) of
_ -> ok
catch
exit:_Reason -> ok
-spec reserve_queue_slot(non_neg_integer()) -> ok.
reserve_queue_slot(Slot) ->
case inflight_ref(Slot) of
undefined -> ok;
Ref -> atomics:add(Ref, ?INFLIGHT_INDEX, 1)
end.
-spec release_queue_slot(atomics:atomics_ref() | undefined) -> ok.
release_queue_slot(undefined) ->
ok;
release_queue_slot(Ref) ->
atomics:sub(Ref, ?INFLIGHT_INDEX, 1).
-spec inflight_ref(non_neg_integer()) -> atomics:atomics_ref() | undefined.
inflight_ref(Slot) ->
persistent_term:get(?INFLIGHT_KEY(Slot), undefined).
-spec max_queue() -> pos_integer().
max_queue() ->
Ceiling = process_health_watchdog:kill_threshold(),
@@ -172,13 +225,20 @@ current_workers_tuple_normalized() ->
-spec select_worker(pid()) -> pid() | undefined.
select_worker(SessionPid) ->
case select_worker_slot(SessionPid) of
undefined -> undefined;
{_Slot, Worker} -> Worker
end.
-spec select_worker_slot(pid()) -> {non_neg_integer(), pid()} | undefined.
select_worker_slot(SessionPid) ->
Workers = current_workers_tuple_normalized(),
case tuple_size(Workers) of
0 ->
undefined;
Count ->
Index = erlang:phash2(SessionPid, Count) + 1,
element(Index, Workers)
{Index - 1, element(Index, Workers)}
end.
-spec message_queue_len(pid()) -> non_neg_integer().
@@ -194,6 +254,7 @@ start_workers(Count) ->
-spec start_worker(non_neg_integer()) -> pid().
start_worker(Index) ->
ok = reset_inflight(Index),
{ok, Pid} = gen_server:start_link(
gateway_dispatch_relay,
{worker, Index},
@@ -201,6 +262,19 @@ start_worker(Index) ->
),
Pid.
-spec reset_inflight(non_neg_integer()) -> ok.
reset_inflight(Index) ->
case inflight_ref(Index) of
undefined -> persistent_term:put(?INFLIGHT_KEY(Index), new_inflight());
Ref -> atomics:put(Ref, ?INFLIGHT_INDEX, ?INFLIGHT_UNSAMPLED)
end.
-spec new_inflight() -> atomics:atomics_ref().
new_inflight() ->
Ref = atomics:new(1, []),
atomics:put(Ref, ?INFLIGHT_INDEX, ?INFLIGHT_UNSAMPLED),
Ref.
-spec worker_index(pid(), tuple(), non_neg_integer()) -> non_neg_integer() | undefined.
worker_index(Pid, Workers, Index) ->
worker_index(Pid, tuple_size(Workers), Workers, Index).
@@ -8,7 +8,7 @@
-type field_kind() ::
percentage
| positive_integer
| non_negative_integer
| relay_max_queue
| rpc_timeout
| reconcile_interval
| concurrency
@@ -59,8 +59,8 @@ valid_config_field(Key, Value) ->
valid_percentage_value(Value);
positive_integer ->
is_integer(Value) andalso Value > 0;
non_negative_integer ->
is_integer(Value) andalso Value >= 0;
relay_max_queue ->
valid_relay_max_queue_value(Value);
rpc_timeout ->
is_integer(Value) andalso Value >= 1000 andalso Value =< 60000;
reconcile_interval ->
@@ -81,7 +81,7 @@ config_field_kind(<<"session_rollout_percentage">>) -> percentage;
config_field_kind(<<"guild_rollout_percentage">>) -> percentage;
config_field_kind(<<"voice_reconciliation_v3_percentage">>) -> percentage;
config_field_kind(<<"gateway_dispatch_relay_shards">>) -> positive_integer;
config_field_kind(<<"gateway_dispatch_relay_max_queue">>) -> non_negative_integer;
config_field_kind(<<"gateway_dispatch_relay_max_queue">>) -> relay_max_queue;
config_field_kind(<<"rpc_request_timeout_ms">>) -> rpc_timeout;
config_field_kind(<<"voice_reconciliation_v3_interval_ms">>) -> reconcile_interval;
config_field_kind(<<"max_concurrent_session_starts">>) -> concurrency;
@@ -93,3 +93,8 @@ config_field_kind(_) -> invalid.
-spec valid_percentage_value(term()) -> boolean().
valid_percentage_value(Value) ->
is_number(Value) andalso Value >= 0 andalso Value =< 100.
-spec valid_relay_max_queue_value(term()) -> boolean().
valid_relay_max_queue_value(Value) ->
is_integer(Value) andalso Value >= 1 andalso
Value =< process_health_watchdog:kill_threshold().
@@ -15,6 +15,9 @@
-define(PROBE_BUDGET, 2).
-define(SHARD_COUNT, 2).
-define(FAST_SHARD_TIMEOUT_MS, 1500).
-define(PRODUCER_COUNTS, [1, 10, 64]).
-define(PEAK_SAMPLES, 100).
-define(PEAK_INTERVAL_MS, 5).
dispatch_is_bounded_and_stays_ordered_test_() ->
{timeout, 30, fun dispatch_is_bounded_and_stays_ordered/0}.
@@ -37,6 +40,12 @@ dispatch_does_not_probe_the_worker_per_event_test_() ->
saturated_shard_does_not_delay_other_shards_test_() ->
{timeout, 60, fun saturated_shard_does_not_delay_other_shards/0}.
concurrent_producers_share_the_worker_bound_test_() ->
{timeout, 120, fun concurrent_producers_share_the_worker_bound/0}.
restarted_worker_starts_from_an_empty_bound_test_() ->
{timeout, 30, fun restarted_worker_starts_from_an_empty_bound/0}.
dispatch_is_bounded_and_stays_ordered() ->
assert_bounded_and_ordered(fun send_dispatch/2).
@@ -105,6 +114,66 @@ saturated_shard_does_not_delay_other_shards() ->
end)
end).
restarted_worker_starts_from_an_empty_bound() ->
with_max_queue(?BOUND, fun() ->
Saturated = saturate_worker(),
Ref = make_ref(),
Session = spawn_session(self(), Ref),
with_worker(fun(_Restarted) ->
Producer = spawn_producer(fun send_dispatch/2, Session),
?assertEqual(finished, producer_status(Producer, 5000)),
?assertEqual(lists:seq(1, ?EVENT_COUNT), collect_observed(Ref, ?EVENT_COUNT, []))
end),
Session ! stop,
stop_worker(Saturated)
end).
saturate_worker() ->
Worker = gateway_dispatch_relay_batch:start_worker(0),
Session = spawn_session(self(), make_ref()),
with_relay_workers([Worker], fun() ->
ok = sys:suspend(Worker),
fill_queue(Worker, ?BOUND),
Producer = spawn_producer(fun send_dispatch/2, Session),
?assertEqual(blocked, producer_status(Producer, 500))
end),
Session ! stop,
Worker.
concurrent_producers_share_the_worker_bound() ->
lists:foreach(fun assert_peak_within_bound/1, ?PRODUCER_COUNTS).
assert_peak_within_bound(Producers) ->
Verdict = bound_verdict(peak_queue_len(Producers)),
?assertEqual({Producers, within_bound}, {Producers, Verdict}).
bound_verdict(Peak) when Peak =< ?BOUND + 1 -> within_bound;
bound_verdict(Peak) -> {over_bound, Peak}.
peak_queue_len(Producers) ->
with_max_queue(?BOUND, fun() ->
with_worker(fun(Worker) -> concurrent_peak(Worker, Producers) end)
end).
concurrent_peak(Worker, Producers) ->
Ref = make_ref(),
Session = spawn_session(self(), Ref),
ok = sys:suspend(Worker),
Pids = [spawn_producer(fun send_dispatch/2, Session) || _ <- lists:seq(1, Producers)],
Peak = sample_peak(Worker, ?PEAK_SAMPLES, 0),
ok = sys:resume(Worker),
lists:foreach(fun(Pid) -> ok = await_producer(Pid, blocked) end, Pids),
_ = collect_observed(Ref, Producers * ?EVENT_COUNT, []),
Session ! stop,
Peak.
sample_peak(_Worker, 0, Peak) ->
Peak;
sample_peak(Worker, Remaining, Peak) ->
ok = gateway_retry_timer:wait(?PEAK_INTERVAL_MS),
Sampled = gateway_dispatch_relay_batch:message_queue_len(Worker),
sample_peak(Worker, Remaining - 1, max(Peak, Sampled)).
assert_bounded_and_ordered(SendFun) ->
{Status, QueueLen, Observed} = run_over_bound(SendFun),
?assertEqual(lists:seq(1, ?EVENT_COUNT), Observed),
@@ -251,3 +251,39 @@ validate_config_rejects_voice_reconciliation_v3_percentage_above_maximum_test()
default_config()
)
).
validate_config_rejects_relay_max_queue_zero_test() ->
?assertMatch(
{error, {invalid_field, <<"gateway_dispatch_relay_max_queue">>, 0}},
gateway_rollout_config_validate:validate(
#{<<"gateway_dispatch_relay_max_queue">> => 0},
default_config()
)
).
validate_config_rejects_relay_max_queue_above_kill_threshold_test() ->
Above = process_health_watchdog:kill_threshold() + 1,
?assertMatch(
{error, {invalid_field, <<"gateway_dispatch_relay_max_queue">>, Above}},
gateway_rollout_config_validate:validate(
#{<<"gateway_dispatch_relay_max_queue">> => Above},
default_config()
)
).
validate_config_accepts_relay_max_queue_range_bounds_test() ->
Ceiling = process_health_watchdog:kill_threshold(),
?assertMatch(
{ok, #{<<"gateway_dispatch_relay_max_queue">> := 1}},
gateway_rollout_config_validate:validate(
#{<<"gateway_dispatch_relay_max_queue">> => 1},
default_config()
)
),
?assertMatch(
{ok, #{<<"gateway_dispatch_relay_max_queue">> := Ceiling}},
gateway_rollout_config_validate:validate(
#{<<"gateway_dispatch_relay_max_queue">> => Ceiling},
default_config()
)
).