diff --git a/fluxer_gateway/src/gateway/gateway_dispatch_relay.erl b/fluxer_gateway/src/gateway/gateway_dispatch_relay.erl index 013b0f576..59099bcfe 100644 --- a/fluxer_gateway/src/gateway/gateway_dispatch_relay.erl +++ b/fluxer_gateway/src/gateway/gateway_dispatch_relay.erl @@ -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) -> diff --git a/fluxer_gateway/src/gateway/gateway_dispatch_relay_batch.erl b/fluxer_gateway/src/gateway/gateway_dispatch_relay_batch.erl index 595ff5615..c6e3c132f 100644 --- a/fluxer_gateway/src/gateway/gateway_dispatch_relay_batch.erl +++ b/fluxer_gateway/src/gateway/gateway_dispatch_relay_batch.erl @@ -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). diff --git a/fluxer_gateway/src/gateway/gateway_rollout_config_validate.erl b/fluxer_gateway/src/gateway/gateway_rollout_config_validate.erl index ebb03d7fb..3c7b2f800 100644 --- a/fluxer_gateway/src/gateway/gateway_rollout_config_validate.erl +++ b/fluxer_gateway/src/gateway/gateway_rollout_config_validate.erl @@ -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(). diff --git a/fluxer_gateway/test/gateway_dispatch_relay_bound_tests.erl b/fluxer_gateway/test/gateway_dispatch_relay_bound_tests.erl index fb1c8d321..cf984ddf9 100644 --- a/fluxer_gateway/test/gateway_dispatch_relay_bound_tests.erl +++ b/fluxer_gateway/test/gateway_dispatch_relay_bound_tests.erl @@ -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), diff --git a/fluxer_gateway/test/gateway_rollout_config_tests.erl b/fluxer_gateway/test/gateway_rollout_config_tests.erl index ecce4cd33..c5030e7d3 100644 --- a/fluxer_gateway/test/gateway_rollout_config_tests.erl +++ b/fluxer_gateway/test/gateway_rollout_config_tests.erl @@ -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() + ) + ).