refactor(gateway): remove voice reconciliation v3 (#2601)

This commit is contained in:
Hampus
2026-09-08 19:17:48 +02:00
committed by GitHub
parent ef067f36c6
commit ce08f82a92
17 changed files with 79 additions and 792 deletions
+2 -16
View File
@@ -12287,14 +12287,7 @@
"max_concurrent_guild_starts": {"type": "integer", "minimum": 1, "maximum": 10000, "format": "int32"},
"gateway_dispatch_relay_shards": {"type": "integer", "minimum": 1, "maximum": 10000, "format": "int32"},
"gateway_dispatch_relay_max_queue": {"type": "integer", "minimum": 0, "maximum": 1000000, "format": "int32"},
"voice_e2ee_scope": {"enum": ["guild_feature_only", "platform_wide"], "type": "string"},
"voice_reconciliation_v3_percentage": {"type": "number", "minimum": 0, "maximum": 100},
"voice_reconciliation_v3_interval_ms": {
"type": "integer",
"minimum": 500,
"maximum": 60000,
"format": "int32"
}
"voice_e2ee_scope": {"enum": ["guild_feature_only", "platform_wide"], "type": "string"}
}
},
"InstanceConfigUpdateRequest": {
@@ -12534,14 +12527,7 @@
"max_concurrent_guild_starts": {"type": "integer", "minimum": 1, "maximum": 10000, "format": "int32"},
"gateway_dispatch_relay_shards": {"type": "integer", "minimum": 1, "maximum": 10000, "format": "int32"},
"gateway_dispatch_relay_max_queue": {"type": "integer", "minimum": 0, "maximum": 1000000, "format": "int32"},
"voice_e2ee_scope": {"enum": ["guild_feature_only", "platform_wide"], "type": "string"},
"voice_reconciliation_v3_percentage": {"type": "number", "minimum": 0, "maximum": 100},
"voice_reconciliation_v3_interval_ms": {
"type": "integer",
"minimum": 500,
"maximum": 60000,
"format": "int32"
}
"voice_e2ee_scope": {"enum": ["guild_feature_only", "platform_wide"], "type": "string"}
}
},
"BrandingAssetUploadRequest": {
@@ -63,8 +63,6 @@ describe('GatewayRolloutConfigPublisher', () => {
gateway_dispatch_relay_shards: 32,
gateway_dispatch_relay_max_queue: 50000,
voice_e2ee_scope: 'guild_feature_only',
voice_reconciliation_v3_percentage: 100,
voice_reconciliation_v3_interval_ms: 2000,
};
await publisher.publish(config);
@@ -42,8 +42,6 @@ const DEFAULT_GATEWAY_ROLLOUT_CONFIG: GatewayRolloutConfig = {
gateway_dispatch_relay_shards: 32,
gateway_dispatch_relay_max_queue: 50000,
voice_e2ee_scope: 'guild_feature_only',
voice_reconciliation_v3_percentage: 100,
voice_reconciliation_v3_interval_ms: 2000,
};
export type InstanceRegistrationMode = 'open' | 'approval' | 'closed';
export interface InstanceRegistrationConfig {
@@ -85,8 +85,6 @@ Admission and dispatch tuning for the Gateway cluster.
| gateway_dispatch_relay_shards | integer | Dispatch relay shard count (1-10000, default 32) |
| gateway_dispatch_relay_max_queue | integer | Dispatch relay queue ceiling (0-1000000, default 50000) |
| voice_e2ee_scope | string | `guild_feature_only` or `platform_wide` (default `guild_feature_only`) |
| voice_reconciliation_v3_percentage | number | Percentage of voice states on the newer reconciliation path (0-100, default 100) |
| voice_reconciliation_v3_interval_ms | integer | Voice reconciliation interval (500-60000, default 2000) |
Every field is present on read. A deployment that has stored nothing reports the defaults above.
-20
View File
@@ -33,7 +33,6 @@
{'DOWN', reference(), process, pid(), term()}
| {ring_timeout, integer()}
| {pending_connection_timeout, binary()}
| voice_reconcile_v3_tick
| idle_timeout.
-type start_result() :: {ok, pid()} | {error, term()}.
@@ -49,14 +48,12 @@ start_link_from_state(State) ->
init({transferred, TransferState}) ->
erlang:process_flag(fullsweep_after, 10),
State = call_handoff:restore_state(TransferState),
voice_reconciliation_v3:schedule_tick(voice_reconcile_v3_tick),
erlang:garbage_collect(),
{ok, State};
init(CallData) ->
erlang:process_flag(fullsweep_after, 10),
State = build_initial_state(CallData),
FinalState = run_init_pipeline(State),
voice_reconciliation_v3:schedule_tick(voice_reconcile_v3_tick),
erlang:garbage_collect(),
{ok, FinalState}.
@@ -176,9 +173,6 @@ handle_info_message({ring_timeout, UserId}, State) ->
call_ringing:handle_ring_timeout(UserId, State);
handle_info_message({pending_connection_timeout, ConnectionId}, State) ->
handle_pending_timeout(ConnectionId, State);
handle_info_message(voice_reconcile_v3_tick, State) ->
voice_reconciliation_v3:schedule_tick(voice_reconcile_v3_tick),
maybe_reconcile_voice_v3(State);
handle_info_message(idle_timeout, State) ->
call_ringing:handle_idle_timeout(State).
@@ -324,8 +318,6 @@ decode_info_message({ring_timeout, UserId}) when is_integer(UserId) ->
{ok, {ring_timeout, UserId}};
decode_info_message({pending_connection_timeout, ConnectionId}) when is_binary(ConnectionId) ->
{ok, {pending_connection_timeout, ConnectionId}};
decode_info_message(voice_reconcile_v3_tick) ->
{ok, voice_reconcile_v3_tick};
decode_info_message(idle_timeout) ->
{ok, idle_timeout};
decode_info_message(_) ->
@@ -451,15 +443,3 @@ check_pending_session_alive(ConnectionId, UserId, SessionId, SessionPid, State)
)
end.
-spec maybe_reconcile_voice_v3(map()) -> {noreply, map()} | {stop, normal, map()}.
maybe_reconcile_voice_v3(#{channel_id := ChannelId, voice_states := VoiceStates} = State) ->
case
maps:size(VoiceStates) > 0 andalso
voice_reconciliation_v3:enabled_for(call, ChannelId)
of
true ->
AbsentEntries = voice_reconciliation_v3:find_absent_call_entries(State),
call_voice:reconcile_absent_connections(AbsentEntries, State);
false ->
{noreply, State}
end.
+1 -89
View File
@@ -10,7 +10,6 @@
handle_disconnect_user/4,
handle_leave/2,
disconnect_user_after_pending_timeout/4,
reconcile_absent_connections/2,
maybe_notify_session_force_disconnect/4,
maybe_spawn_region_switch/3,
is_session_pid_alive/1
@@ -155,7 +154,7 @@ handle_disconnect_user(
end,
case
voice_disconnect_common:disconnect_user_if_in_channel(
UserId, ExpectedChannelId, VoiceStates, Sessions, CleanupFun
UserId, ExpectedChannelId, ConnectionId, VoiceStates, Sessions, CleanupFun
)
of
{not_found, _, _} ->
@@ -238,93 +237,6 @@ disconnect_user_after_pending_timeout(
}}
end.
-spec reconcile_absent_connections([voice_reconciliation_v3:participant_entry()], map()) ->
{noreply, map()} | {stop, normal, map()}.
reconcile_absent_connections([], State) ->
{noreply, State};
reconcile_absent_connections(
AbsentEntries,
#{
voice_states := VoiceStates,
sessions := Sessions,
pending_connections := PendingConns
} = State
) ->
ActiveAbsentEntries = active_absent_entries(AbsentEntries, VoiceStates),
case ActiveAbsentEntries of
[] ->
{noreply, State};
_ ->
do_reconcile_absent_connections(
ActiveAbsentEntries, VoiceStates, Sessions, PendingConns, State
)
end.
-spec active_absent_entries([voice_reconciliation_v3:participant_entry()], map()) ->
[voice_reconciliation_v3:participant_entry()].
active_absent_entries(AbsentEntries, VoiceStates) ->
lists:filter(
fun(Entry) -> absent_entry_active(Entry, VoiceStates) end,
AbsentEntries
).
-spec absent_entry_active(voice_reconciliation_v3:participant_entry(), map()) -> boolean().
absent_entry_active(#{user_id := UserId, connection_id := ConnectionId}, VoiceStates) ->
case maps:get(UserId, VoiceStates, undefined) of
VoiceState when is_map(VoiceState) ->
maps:get(<<"connection_id">>, VoiceState, undefined) =:= ConnectionId;
_ ->
false
end.
-spec do_reconcile_absent_connections(
[voice_reconciliation_v3:participant_entry()], map(), map(), map(), map()
) -> {noreply, map()} | {stop, normal, map()}.
do_reconcile_absent_connections(
ActiveAbsentEntries, VoiceStates, Sessions, PendingConns, State
) ->
RemovedUsers = [maps:get(user_id, Entry) || Entry <- ActiveAbsentEntries],
NewVoiceStates = lists:foldl(fun maps:remove/2, VoiceStates, RemovedUsers),
{NewSessions, NewPending} = remove_absent_sessions(
ActiveAbsentEntries, Sessions, PendingConns, State
),
BaseState = State#{
voice_states => NewVoiceStates,
sessions => NewSessions,
pending_connections => NewPending
},
CleanState = call_ringing:cancel_ringing_timers(RemovedUsers, BaseState),
RingState = call_ringing:remove_users_from_ringing(RemovedUsers, CleanState),
{UpdatedState, Dispatched} = call_ringing:maybe_dispatch_state_update(State, RingState),
call_ringing:maybe_stop_or_noreply(UpdatedState, Dispatched).
-spec remove_absent_sessions(
[voice_reconciliation_v3:participant_entry()], map(), map(), map()
) -> {map(), map()}.
remove_absent_sessions(ActiveAbsentEntries, Sessions, PendingConns, State) ->
lists:foldl(
fun(Entry, {SessionsAcc, PendingAcc}) ->
remove_absent_session(Entry, SessionsAcc, PendingAcc, State)
end,
{Sessions, PendingConns},
ActiveAbsentEntries
).
-spec remove_absent_session(voice_reconciliation_v3:participant_entry(), map(), map(), map()) ->
{map(), map()}.
remove_absent_session(
#{user_id := UserId, connection_id := ConnectionId}, Sessions, Pending, State
) ->
NewPending = voice_pending_common:remove_pending_connection(ConnectionId, Pending),
case voice_disconnect_common:find_session_by_user_id(UserId, Sessions) of
{ok, SessionId, _Pid, Ref} ->
demonitor(Ref, [flush]),
maybe_notify_session_force_disconnect(UserId, SessionId, ConnectionId, State),
{maps:remove(SessionId, Sessions), NewPending};
not_found ->
{Sessions, NewPending}
end.
-spec maybe_notify_session_force_disconnect(
integer(), binary(), binary() | undefined, map()
) -> ok.
@@ -16,8 +16,6 @@
gateway_dispatch_relay_shards/0,
gateway_dispatch_relay_max_queue/0,
rpc_request_timeout_ms/0,
voice_reconciliation_v3_percentage/0,
voice_reconciliation_v3_interval_ms/0,
max_concurrent_session_starts/0,
max_concurrent_guild_starts/0,
is_session_eligible/1,
@@ -99,14 +97,6 @@ gateway_dispatch_relay_max_queue() ->
rpc_request_timeout_ms() ->
maps:get(<<"rpc_request_timeout_ms">>, get(), 10000).
-spec voice_reconciliation_v3_percentage() -> number().
voice_reconciliation_v3_percentage() ->
maps:get(<<"voice_reconciliation_v3_percentage">>, get(), 100).
-spec voice_reconciliation_v3_interval_ms() -> integer().
voice_reconciliation_v3_interval_ms() ->
maps:get(<<"voice_reconciliation_v3_interval_ms">>, get(), 2000).
-spec max_concurrent_session_starts() -> integer().
max_concurrent_session_starts() ->
maps:get(<<"max_concurrent_session_starts">>, get(), 512).
@@ -237,9 +227,7 @@ default_config() ->
<<"max_concurrent_guild_starts">> => 256,
<<"gateway_dispatch_relay_shards">> => 32,
<<"gateway_dispatch_relay_max_queue">> => 50000,
<<"voice_e2ee_scope">> => <<"guild_feature_only">>,
<<"voice_reconciliation_v3_percentage">> => 100,
<<"voice_reconciliation_v3_interval_ms">> => 2000
<<"voice_e2ee_scope">> => <<"guild_feature_only">>
}.
-spec initial_config() -> map().
@@ -357,9 +345,7 @@ log_config_transitions(OldConfig, NewConfig) ->
WatchKeys = [
<<"session_rollout_percentage">>,
<<"guild_rollout_percentage">>,
<<"rpc_request_timeout_ms">>,
<<"voice_reconciliation_v3_percentage">>,
<<"voice_reconciliation_v3_interval_ms">>
<<"rpc_request_timeout_ms">>
],
lists:foreach(
fun(Key) -> log_key_transition(Key, OldConfig, NewConfig) end,
@@ -10,7 +10,6 @@
| positive_integer
| relay_max_queue
| rpc_timeout
| reconcile_interval
| concurrency
| rollout_mode
| voice_e2ee_scope
@@ -37,9 +36,7 @@ config_fields() ->
<<"max_concurrent_guild_starts">>,
<<"gateway_dispatch_relay_shards">>,
<<"gateway_dispatch_relay_max_queue">>,
<<"voice_e2ee_scope">>,
<<"voice_reconciliation_v3_percentage">>,
<<"voice_reconciliation_v3_interval_ms">>
<<"voice_e2ee_scope">>
].
-spec validate_fields([binary()], map()) -> ok | {error, term()}.
@@ -63,8 +60,6 @@ valid_config_field(Key, Value) ->
valid_relay_max_queue_value(Value);
rpc_timeout ->
is_integer(Value) andalso Value >= 1000 andalso Value =< 60000;
reconcile_interval ->
is_integer(Value) andalso Value >= 500 andalso Value =< 60000;
concurrency ->
is_integer(Value) andalso Value >= 1 andalso Value =< 10000;
rollout_mode ->
@@ -79,11 +74,9 @@ valid_config_field(Key, Value) ->
-spec config_field_kind(binary()) -> field_kind().
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">>) -> 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;
config_field_kind(<<"max_concurrent_guild_starts">>) -> concurrency;
config_field_kind(<<"session_rollout_mode">>) -> rollout_mode;
@@ -8,7 +8,6 @@
-export([disconnect_voice_user/2]).
-export([disconnect_voice_user_if_in_channel/2]).
-export([disconnect_all_voice_users_in_channel/2]).
-export([reconcile_absent_voice_connections/2]).
-export([cleanup_virtual_channel_access_for_user/2]).
-export([recently_disconnected_voice_states/1]).
-export([clear_recently_disconnected/2]).
@@ -52,10 +51,6 @@ disconnect_voice_user_if_in_channel(Request, State) ->
disconnect_all_voice_users_in_channel(Request, State) ->
guild_voice_disconnect_channel:disconnect_all_voice_users_in_channel(Request, State).
-spec reconcile_absent_voice_connections([binary()], guild_state()) -> guild_state().
reconcile_absent_voice_connections(ConnectionIds, State) ->
guild_voice_disconnect_user:reconcile_absent_voice_connections(ConnectionIds, State).
-spec cleanup_virtual_channel_access_for_user(integer(), guild_state()) -> guild_state().
cleanup_virtual_channel_access_for_user(UserId, State) ->
guild_voice_disconnect_user:cleanup_virtual_channel_access_for_user(UserId, State).
@@ -208,23 +203,6 @@ disconnect_all_voice_users_in_channel_test() ->
?assert(maps:is_key(<<"c">>, Remaining)),
?assertNot(maps:is_key(<<"a">>, Remaining)).
reconcile_absent_voice_connections_removes_without_force_disconnect_test() ->
TestFun = test_force_disconnect_fun(),
VoiceStates = #{
<<"a">> => voice_state_fixture(5, 10, 20),
<<"b">> => voice_state_fixture(6, 10, 20)
},
State = #{
id => 10,
voice_states => VoiceStates,
test_force_disconnect_fun => TestFun
},
NewState = reconcile_absent_voice_connections([<<"a">>], State),
Remaining = maps:get(voice_states, NewState),
?assertNot(maps:is_key(<<"a">>, Remaining)),
?assert(maps:is_key(<<"b">>, Remaining)),
?assertEqual([], collect_force_disconnect_messages(1)).
pending_connection_tests_test() ->
Pending = #{
<<"conn1">> => #{user_id => 5, channel_id => 100},
@@ -7,7 +7,6 @@
handle_voice_disconnect/5,
disconnect_voice_user/2,
disconnect_voice_user_if_in_channel/2,
reconcile_absent_voice_connections/2,
force_disconnect_participant/4,
cleanup_virtual_channel_access_for_user/2
]).
@@ -197,51 +196,6 @@ disconnect_voice_user_if_in_channel(
)
end.
-spec reconcile_absent_voice_connections([binary()], guild_state()) -> guild_state().
reconcile_absent_voice_connections([], State) ->
State;
reconcile_absent_voice_connections(ConnectionIds, State) when is_list(ConnectionIds) ->
VoiceStates = voice_state_utils:voice_states(State),
RemovedVoiceStates = maps:with(ConnectionIds, VoiceStates),
case maps:size(RemovedVoiceStates) of
0 ->
State;
_ ->
do_reconcile_absent_voice_connections(RemovedVoiceStates, VoiceStates, State)
end.
-spec do_reconcile_absent_voice_connections(
voice_state_map(), voice_state_map(), guild_state()
) ->
guild_state().
do_reconcile_absent_voice_connections(RemovedVoiceStates, VoiceStates, State) ->
ok = guild_voice_disconnect_broadcast:purge_count_cache(maps:keys(RemovedVoiceStates)),
NewVoiceStates = voice_state_utils:drop_voice_states(RemovedVoiceStates, VoiceStates),
NewState0 = State#{voice_states => NewVoiceStates},
NewState1 = guild_voice_disconnect_broadcast:clear_e2ee_room_keys_for_removed(
RemovedVoiceStates, NewVoiceStates, NewState0
),
NewState2 = clear_recently_disconnected_connections(RemovedVoiceStates, NewState1),
voice_state_utils:broadcast_disconnects(RemovedVoiceStates, NewState2),
cleanup_absent_users(RemovedVoiceStates, NewVoiceStates, NewState2).
-spec cleanup_absent_users(voice_state_map(), voice_state_map(), guild_state()) ->
guild_state().
cleanup_absent_users(RemovedVoiceStates, RemainingVoiceStates, State) ->
UserIds = lists:usort([
UserId
|| VoiceState <- maps:values(RemovedVoiceStates),
UserId <- [voice_state_utils:voice_state_user_id(VoiceState)],
is_integer(UserId)
]),
lists:foldl(
fun(UserId, AccState) ->
maybe_cleanup_virtual_channel_access(UserId, RemainingVoiceStates, AccState)
end,
State,
UserIds
).
-spec force_disconnect_participant(integer(), integer(), integer(), binary()) ->
{ok, map()} | {error, term()}.
force_disconnect_participant(GuildId, ChannelId, UserId, ConnectionId) ->
@@ -131,7 +131,6 @@ init(#{guild_id := GuildId, guild_pid := GuildPid} = Args) ->
guild_voice_server_sync:ensure_registry(),
ets:insert(?REGISTRY_TABLE, {GuildId, self()}),
erlang:send_after(?SWEEP_INTERVAL_MS, self(), sweep_pending_joins),
voice_reconciliation_v3:schedule_tick(voice_reconcile_v3_tick),
InitialVoiceStates = voice_state_utils:ensure_voice_states(
maps:get(initial_voice_states, Args, #{})
),
@@ -261,9 +260,6 @@ handle_info(sweep_pending_joins, State) ->
NewGuildState = guild_voice_connection:sweep_expired_pending_joins(GuildState),
{noreply, guild_voice_server_state:apply_guild_state(NewGuildState, State2)}
end;
handle_info(voice_reconcile_v3_tick, State) ->
voice_reconciliation_v3:schedule_tick(voice_reconcile_v3_tick),
{noreply, maybe_reconcile_voice_v3(State)};
handle_info({'EXIT', Pid, Reason}, #{guild_pid := GuildPid} = State) when Pid =:= GuildPid ->
logger:info(
"Voice server shutting down because guild process exited",
@@ -399,20 +395,6 @@ bounded_put_new(Key, Value, Map, MaxSize) ->
Map#{Key => Value}
end.
-spec maybe_reconcile_voice_v3(server_state()) -> server_state().
maybe_reconcile_voice_v3(#{guild_id := GuildId, voice_states := VoiceStates} = State) ->
case
maps:size(VoiceStates) > 0 andalso
voice_reconciliation_v3:enabled_for(guild, GuildId)
of
true ->
AbsentConnectionIds = voice_reconciliation_v3:find_absent_guild_connections(State),
guild_voice_disconnect:reconcile_absent_voice_connections(
AbsentConnectionIds, State
);
false ->
State
end.
-spec sweep_recently_disconnected(map()) -> map().
sweep_recently_disconnected(State) ->
@@ -7,6 +7,7 @@
find_session_by_user_id/2,
disconnect_user/4,
disconnect_user_if_in_channel/5,
disconnect_user_if_in_channel/6,
channel_has_capacity/3
]).
@@ -93,18 +94,46 @@ sessions_for_user(UserId, Sessions) ->
| {not_found, voice_states_map(), sessions_map()}
| {channel_mismatch, voice_states_map(), sessions_map()}.
disconnect_user_if_in_channel(UserId, ExpectedChannelId, VoiceStates, Sessions, CleanupFun) ->
disconnect_user_if_in_channel(
UserId, ExpectedChannelId, undefined, VoiceStates, Sessions, CleanupFun
).
-spec disconnect_user_if_in_channel(
user_id(), integer(), binary() | undefined, voice_states_map(), sessions_map(), cleanup_fun()
) ->
{ok, voice_states_map(), sessions_map()}
| {not_found, voice_states_map(), sessions_map()}
| {channel_mismatch, voice_states_map(), sessions_map()}.
disconnect_user_if_in_channel(
UserId, ExpectedChannelId, ExpectedConnectionId, VoiceStates, Sessions, CleanupFun
) ->
case maps:get(UserId, VoiceStates, undefined) of
undefined ->
{not_found, VoiceStates, Sessions};
VoiceState ->
disconnect_matching_channel(
UserId,
ExpectedChannelId,
VoiceState,
VoiceStates,
Sessions,
CleanupFun
)
case connection_matches(VoiceState, ExpectedConnectionId) of
false ->
{not_found, VoiceStates, Sessions};
true ->
disconnect_matching_channel(
UserId,
ExpectedChannelId,
VoiceState,
VoiceStates,
Sessions,
CleanupFun
)
end
end.
-spec connection_matches(map(), binary() | undefined) -> boolean().
connection_matches(_VoiceState, undefined) ->
true;
connection_matches(VoiceState, ExpectedConnectionId) ->
case maps:get(<<"connection_id">>, VoiceState, undefined) of
undefined -> true;
ExpectedConnectionId -> true;
_ -> false
end.
-spec disconnect_matching_channel(
@@ -1,461 +0,0 @@
%% SPDX-License-Identifier: AGPL-3.0-or-later
-module(voice_reconciliation_v3).
-typing([eqwalizer]).
-export([
schedule_tick/1,
interval_ms/0,
enabled_for/2,
find_absent_guild_connections/1,
find_absent_call_entries/1,
find_absent_entries/2
]).
-ifdef(TEST).
-include_lib("eunit/include/eunit.hrl").
-endif.
-export_type([
owner_kind/0,
participant_entry/0,
room_key/0,
snapshot_fun/0
]).
-type owner_kind() :: guild | call.
-type room_key() :: {integer() | null, integer(), binary(), binary()}.
-type participant_entry() :: #{
connection_id := binary(),
user_id := integer(),
channel_id := integer(),
guild_id := integer() | null,
region_id := binary(),
server_id := binary(),
pending := boolean()
}.
-type snapshot_fun() ::
fun((integer() | null, integer(), binary(), binary()) -> {ok, term()} | {error, term()}).
-define(DEFAULT_INTERVAL_MS, 2000).
-spec schedule_tick(term()) -> reference().
schedule_tick(Message) ->
erlang:send_after(jittered_interval_ms(), self(), Message).
-spec interval_ms() -> pos_integer().
interval_ms() ->
normalize_interval(gateway_rollout_config:voice_reconciliation_v3_interval_ms()).
-spec enabled_for(owner_kind(), integer()) -> boolean().
enabled_for(_Kind, OwnerId) when is_integer(OwnerId), OwnerId > 0 ->
Percentage = gateway_rollout_config:voice_reconciliation_v3_percentage(),
Percentage > 0 andalso
(Percentage >= 100 orelse erlang:phash2(integer_to_binary(OwnerId), 100) < Percentage);
enabled_for(_Kind, _OwnerId) ->
false.
-spec find_absent_guild_connections(map()) -> [binary()].
find_absent_guild_connections(State) ->
GuildId = maps:get(guild_id, State, maps:get(id, State, undefined)),
VoiceStates = voice_state_utils:voice_states(State),
Pending = ensure_map(maps:get(pending_voice_connections, State, #{})),
Entries = guild_entries(GuildId, VoiceStates, Pending),
Absent = find_absent_entries(Entries, snapshot_fun(State)),
[maps:get(connection_id, Entry) || Entry <- Absent].
-spec find_absent_call_entries(map()) -> [participant_entry()].
find_absent_call_entries(#{channel_id := ChannelId} = State) when
is_integer(ChannelId), ChannelId > 0
->
VoiceStates = ensure_map(maps:get(voice_states, State, #{})),
Pending = ensure_map(maps:get(pending_connections, State, #{})),
Entries = call_entries(ChannelId, VoiceStates, Pending),
find_absent_entries(Entries, snapshot_fun(State));
find_absent_call_entries(_State) ->
[].
-spec find_absent_entries([participant_entry()], snapshot_fun()) -> [participant_entry()].
find_absent_entries(Entries, SnapshotFun) ->
RoomGroups = group_entries_by_room([E || E <- Entries, not maps:get(pending, E)]),
maps:fold(
fun(RoomKey, RoomEntries, Acc) ->
find_absent_entries_for_room(RoomKey, RoomEntries, SnapshotFun, Acc)
end,
[],
RoomGroups
).
-spec find_absent_entries_for_room(room_key(), [participant_entry()], snapshot_fun(), [
participant_entry()
]) ->
[participant_entry()].
find_absent_entries_for_room(
{GuildId, ChannelId, RegionId, ServerId} = RoomKey, RoomEntries, SnapshotFun, Acc
) ->
case fetch_present_connections(SnapshotFun, GuildId, ChannelId, RegionId, ServerId) of
{ok, Present} ->
absent_from_present(RoomEntries, Present, Acc);
{error, Reason} ->
log_snapshot_error(RoomKey, Reason),
Acc
end.
-spec absent_from_present([participant_entry()], map(), [participant_entry()]) ->
[participant_entry()].
absent_from_present(RoomEntries, Present, Acc) ->
lists:foldl(
fun(Entry, EntryAcc) -> prepend_if_absent(Entry, Present, EntryAcc) end,
Acc,
RoomEntries
).
-spec prepend_if_absent(participant_entry(), map(), [participant_entry()]) ->
[participant_entry()].
prepend_if_absent(Entry, Present, Acc) ->
case maps:is_key(maps:get(connection_id, Entry), Present) of
true -> Acc;
false -> [Entry | Acc]
end.
-spec guild_entries(term(), map(), map()) -> [participant_entry()].
guild_entries(GuildId, VoiceStates, Pending) ->
case normalize_guild_id(GuildId) of
undefined ->
[];
NormalizedGuildId ->
guild_entries_for_id(NormalizedGuildId, VoiceStates, Pending)
end.
-spec guild_entries_for_id(integer(), map(), map()) -> [participant_entry()].
guild_entries_for_id(GuildId, VoiceStates, Pending) ->
maps:fold(
fun(ConnId, VoiceState, Acc) ->
prepend_entry(
entry_from_voice_state(GuildId, ConnId, undefined, VoiceState, Pending),
Acc
)
end,
[],
VoiceStates
).
-spec call_entries(integer(), map(), map()) -> [participant_entry()].
call_entries(ChannelId, VoiceStates, Pending) ->
maps:fold(
fun(UserId, VoiceState, Acc) ->
ConnId = normalize_connection_id(
maps:get(
<<"connection_id">>,
VoiceState,
maps:get(connection_id, VoiceState, undefined)
)
),
prepend_entry(
entry_from_voice_state(null, ConnId, UserId, VoiceState, Pending, ChannelId),
Acc
)
end,
[],
VoiceStates
).
-spec entry_from_voice_state(
integer() | null, term(), term(), term(), map()
) -> {ok, participant_entry()} | skip.
entry_from_voice_state(GuildId, ConnId, FallbackUserId, VoiceState, Pending) ->
entry_from_voice_state(GuildId, ConnId, FallbackUserId, VoiceState, Pending, undefined).
-spec entry_from_voice_state(
integer() | null, term(), term(), term(), map(), integer() | undefined
) -> {ok, participant_entry()} | skip.
entry_from_voice_state(
GuildId, ConnId0, FallbackUserId, VoiceState, Pending, FallbackChannelId
) when
is_map(VoiceState)
->
ConnId = normalize_connection_id(ConnId0),
UserId = normalize_user_id(VoiceState, FallbackUserId),
ChannelId = normalize_channel_id(VoiceState, FallbackChannelId),
RegionId = normalize_optional_binary(
maps:get(<<"region_id">>, VoiceState, maps:get(region_id, VoiceState, undefined))
),
ServerId = normalize_optional_binary(
maps:get(<<"server_id">>, VoiceState, maps:get(server_id, VoiceState, undefined))
),
case {ConnId, UserId, ChannelId, RegionId, ServerId} of
{C, U, Ch, R, S} when
is_binary(C),
is_integer(U),
U > 0,
is_integer(Ch),
Ch > 0,
is_binary(R),
is_binary(S)
->
{ok, #{
connection_id => C,
user_id => U,
channel_id => Ch,
guild_id => GuildId,
region_id => R,
server_id => S,
pending => maps:is_key(C, Pending)
}};
_ ->
skip
end;
entry_from_voice_state(_GuildId, _ConnId, _FallbackUserId, _VoiceState, _Pending, _ChannelId) ->
skip.
-spec prepend_entry({ok, participant_entry()} | skip, [participant_entry()]) ->
[participant_entry()].
prepend_entry({ok, Entry}, Acc) -> [Entry | Acc];
prepend_entry(skip, Acc) -> Acc.
-spec group_entries_by_room([participant_entry()]) -> #{room_key() => [participant_entry()]}.
group_entries_by_room(Entries) ->
lists:foldl(
fun(Entry, Acc) ->
Key = room_key(Entry),
Acc#{Key => [Entry | maps:get(Key, Acc, [])]}
end,
#{},
Entries
).
-spec room_key(participant_entry()) -> room_key().
room_key(Entry) ->
{
maps:get(guild_id, Entry),
maps:get(channel_id, Entry),
maps:get(region_id, Entry),
maps:get(server_id, Entry)
}.
-spec fetch_present_connections(
snapshot_fun(), integer() | null, integer(), binary(), binary()
) -> {ok, map()} | {error, term()}.
fetch_present_connections(SnapshotFun, GuildId, ChannelId, RegionId, ServerId) ->
try SnapshotFun(GuildId, ChannelId, RegionId, ServerId) of
{ok, Snapshot} -> {ok, connection_set(Snapshot)};
{error, Reason} -> {error, Reason};
Other -> {error, {unexpected_snapshot_result, Other}}
catch
Class:Reason -> {error, {Class, Reason}}
end.
-spec snapshot_fun(map()) -> snapshot_fun().
snapshot_fun(State) ->
case maps:get(test_voice_reconciliation_v3_fun, State, undefined) of
Fun when is_function(Fun, 4) -> Fun;
_ -> fun fetch_snapshot_from_rpc/4
end.
-spec fetch_snapshot_from_rpc(integer() | null, integer(), binary(), binary()) ->
{ok, term()} | {error, term()}.
fetch_snapshot_from_rpc(GuildId, ChannelId, RegionId, ServerId) ->
Req = voice_utils:build_list_participants_rpc_request(
GuildId, ChannelId, RegionId, ServerId
),
case rpc_client:call(Req) of
{ok, #{<<"status">> := <<"ok">>, <<"participants">> := Participants}} ->
{ok, Participants};
{ok, #{<<"status">> := <<"error">>} = ErrorData} ->
{error, ErrorData};
{ok, Other} ->
{error, {unexpected_rpc_response, Other}};
{error, Reason} ->
{error, Reason}
end.
-spec connection_set(term()) -> map().
connection_set(#{<<"participants">> := Participants}) ->
connection_set(Participants);
connection_set(#{participants := Participants}) ->
connection_set(Participants);
connection_set(Participants) when is_list(Participants) ->
lists:foldl(fun add_snapshot_connection/2, #{}, Participants);
connection_set(_) ->
#{}.
-spec add_snapshot_connection(term(), map()) -> map().
add_snapshot_connection(ConnectionId, Acc) when
is_binary(ConnectionId), byte_size(ConnectionId) > 0
->
Acc#{ConnectionId => true};
add_snapshot_connection(Participant, Acc) when is_map(Participant) ->
case participant_connection_id(Participant) of
undefined -> Acc;
ConnectionId -> Acc#{ConnectionId => true}
end;
add_snapshot_connection(_Participant, Acc) ->
Acc.
-spec participant_connection_id(map()) -> binary() | undefined.
participant_connection_id(Participant) ->
case
normalize_connection_id(
maps:get(
<<"connection_id">>,
Participant,
maps:get(connection_id, Participant, undefined)
)
)
of
undefined ->
parse_identity_connection_id(
maps:get(
<<"identity">>, Participant, maps:get(identity, Participant, undefined)
)
);
ConnectionId ->
ConnectionId
end.
-spec parse_identity_connection_id(term()) -> binary() | undefined.
parse_identity_connection_id(Identity) when is_binary(Identity) ->
parse_identity_parts(binary:split(Identity, <<"_">>, [global]));
parse_identity_connection_id(_) ->
undefined.
-spec parse_identity_parts([binary()]) -> binary() | undefined.
parse_identity_parts([<<"user">>, UserId, ConnId]) when
byte_size(UserId) > 0, byte_size(ConnId) > 0
->
ConnId;
parse_identity_parts([<<"user">>, UserId | Rest]) when byte_size(UserId) > 0, Rest =/= [] ->
nonempty_joined_connection_id(Rest);
parse_identity_parts(_) ->
undefined.
-spec nonempty_joined_connection_id([binary()]) -> binary() | undefined.
nonempty_joined_connection_id(Parts) ->
case join_binary(Parts, <<"_">>) of
<<>> -> undefined;
ConnId -> ConnId
end.
-spec join_binary([binary()], binary()) -> binary().
join_binary([Part], _Separator) ->
Part;
join_binary([Part | Rest], Separator) ->
<<Part/binary, Separator/binary, (join_binary(Rest, Separator))/binary>>.
-spec normalize_user_id(map(), term()) -> integer() | undefined.
normalize_user_id(VoiceState, FallbackUserId) ->
case voice_state_utils:voice_state_user_id(VoiceState) of
undefined -> normalize_positive_snowflake(FallbackUserId);
UserId -> UserId
end.
-spec normalize_channel_id(map(), integer() | undefined) -> integer() | undefined.
normalize_channel_id(VoiceState, FallbackChannelId) ->
case voice_state_utils:voice_state_channel_id(VoiceState) of
undefined -> FallbackChannelId;
ChannelId -> ChannelId
end.
-spec normalize_guild_id(term()) -> integer() | undefined.
normalize_guild_id(GuildId) ->
normalize_positive_snowflake(GuildId).
-spec normalize_positive_snowflake(term()) -> integer() | undefined.
normalize_positive_snowflake(Value) ->
guild_voice_connection_normalize:normalize_positive_snowflake(Value).
-spec normalize_connection_id(term()) -> binary() | undefined.
normalize_connection_id(Value) ->
normalize_optional_binary(Value).
-spec normalize_optional_binary(term()) -> binary() | undefined.
normalize_optional_binary(Value) ->
guild_voice_connection_normalize:normalize_optional_binary(Value).
-spec normalize_interval(term()) -> pos_integer().
normalize_interval(Value) when is_integer(Value), Value >= 500, Value =< 60000 ->
Value;
normalize_interval(_) ->
?DEFAULT_INTERVAL_MS.
-spec jittered_interval_ms() -> pos_integer().
jittered_interval_ms() ->
Interval = interval_ms(),
Jitter = min(250, max(0, Interval div 10)),
Interval + erlang:phash2(self(), Jitter + 1).
-spec ensure_map(term()) -> map().
ensure_map(Map) when is_map(Map) -> Map;
ensure_map(_) -> #{}.
-spec log_snapshot_error(room_key(), term()) -> ok.
log_snapshot_error({GuildId, ChannelId, RegionId, ServerId}, Reason) ->
logger:debug(
"voice_reconciliation_v3_snapshot_error:"
" guild_id=~p channel_id=~p region_id=~p server_id=~p reason=~p",
[GuildId, ChannelId, RegionId, ServerId, Reason]
),
ok.
-ifdef(TEST).
voice_state(UserId, ChannelId, ConnId, RegionId, ServerId) ->
#{
<<"user_id">> => integer_to_binary(UserId),
<<"channel_id">> => integer_to_binary(ChannelId),
<<"connection_id">> => ConnId,
<<"region_id">> => RegionId,
<<"server_id">> => ServerId
}.
find_absent_guild_connections_skips_pending_test() ->
State = #{
guild_id => 10,
voice_states => #{
<<"conn-present">> => voice_state(1, 20, <<"conn-present">>, <<"local">>, <<"s1">>),
<<"conn-absent">> => voice_state(2, 20, <<"conn-absent">>, <<"local">>, <<"s1">>),
<<"conn-pending">> => voice_state(3, 20, <<"conn-pending">>, <<"local">>, <<"s1">>)
},
pending_voice_connections => #{<<"conn-pending">> => #{}},
test_voice_reconciliation_v3_fun =>
fun(10, 20, <<"local">>, <<"s1">>) -> {ok, [<<"conn-present">>]} end
},
?assertEqual([<<"conn-absent">>], find_absent_guild_connections(State)).
find_absent_entries_keeps_room_on_snapshot_error_test() ->
Entries = [
#{
connection_id => <<"conn">>,
user_id => 1,
channel_id => 20,
guild_id => 10,
region_id => <<"local">>,
server_id => <<"s1">>,
pending => false
}
],
Fun = fun(_GuildId, _ChannelId, _RegionId, _ServerId) -> {error, unavailable} end,
?assertEqual([], find_absent_entries(Entries, Fun)).
find_absent_call_entries_uses_connection_id_from_voice_state_test() ->
State = #{
channel_id => 30,
voice_states => #{5 => voice_state(5, 30, <<"conn-a">>, <<"local">>, <<"s1">>)},
pending_connections => #{},
test_voice_reconciliation_v3_fun =>
fun(null, 30, <<"local">>, <<"s1">>) -> {ok, []} end
},
?assertMatch(
[#{connection_id := <<"conn-a">>, user_id := 5}], find_absent_call_entries(State)
).
connection_set_parses_identity_fallback_test() ->
Snapshot = [
#{<<"identity">> => <<"user_42_conn_with_underscores">>},
#{<<"connection_id">> => <<"conn-direct">>}
],
Set = connection_set(Snapshot),
?assert(maps:is_key(<<"conn_with_underscores">>, Set)),
?assert(maps:is_key(<<"conn-direct">>, Set)).
-endif.
+32 -50
View File
@@ -123,56 +123,6 @@ pending_connection_timeout_preserves_active_voice_state_test() ->
)
).
reconcile_absent_connections_removes_matching_voice_state_test() ->
UserId = 42,
OtherUserId = 84,
SessionId = <<"session-42">>,
OtherSessionId = <<"session-84">>,
ConnectionId = <<"conn-42">>,
OtherConnectionId = <<"conn-84">>,
Ref = monitor(process, self()),
OtherRef = monitor(process, self()),
VoiceState = #{
<<"user_id">> => <<"42">>,
<<"channel_id">> => <<"1">>,
<<"connection_id">> => ConnectionId
},
OtherVoiceState = #{
<<"user_id">> => <<"84">>,
<<"channel_id">> => <<"1">>,
<<"connection_id">> => OtherConnectionId
},
State = new_call_test_state(#{
voice_states => #{UserId => VoiceState, OtherUserId => OtherVoiceState},
sessions => #{
SessionId => {UserId, self(), Ref},
OtherSessionId => {OtherUserId, self(), OtherRef}
},
pending_connections => #{ConnectionId => #{user_id => UserId}},
initiator_ready => true
}),
Absent = [
#{
connection_id => ConnectionId,
user_id => UserId,
channel_id => 1,
guild_id => null,
region_id => <<"local">>,
server_id => <<"s1">>,
pending => false
}
],
{noreply, NewState} = call_voice:reconcile_absent_connections(Absent, State),
?assertNot(maps:is_key(UserId, maps:get(voice_states, NewState))),
?assert(maps:is_key(OtherUserId, maps:get(voice_states, NewState))),
?assertNot(maps:is_key(SessionId, maps:get(sessions, NewState))),
?assertNot(maps:is_key(ConnectionId, maps:get(pending_connections, NewState))),
receive
{'$gen_cast', {call_force_disconnect, 1, ConnectionId}} -> ok
after 1000 ->
?assert(false, force_disconnect_not_sent)
end.
leave_then_rejoin_keeps_voice_state_test() ->
UserId = 42,
OtherUserId = 84,
@@ -264,6 +214,38 @@ disconnect_user_if_in_channel_notifies_session_cleanup_test() ->
?assert(false)
end.
disconnect_user_if_in_channel_ignores_a_stale_connection_id_test() ->
UserId = 42,
SessionId = <<"session-42">>,
LiveConnectionId = <<"conn-live">>,
StaleConnectionId = <<"conn-stale">>,
VoiceState = #{
<<"user_id">> => <<"42">>,
<<"channel_id">> => <<"1">>,
<<"connection_id">> => LiveConnectionId
},
State = new_call_test_state(#{
voice_states => #{UserId => VoiceState},
sessions => #{SessionId => {UserId, self(), make_ref()}},
initiator_ready => true
}),
{reply, Reply, NewState} = call:handle_call(
{disconnect_user_if_in_channel, UserId, 1, StaleConnectionId},
{self(), make_ref()},
State
),
?assertMatch(#{ignored := true, reason := <<"not_in_call">>}, Reply),
?assertEqual(
VoiceState,
maps:get(UserId, maps:get(voice_states, NewState), undefined)
),
receive
{'$gen_cast', {call_force_disconnect, _, _}} ->
?assert(false, force_disconnected_a_live_connection)
after 200 ->
ok
end.
leave_removes_voice_state_count_without_full_rebuild_test() ->
UserId = 42,
SessionId = <<"session-count-leave">>,
@@ -16,9 +16,7 @@ default_config() ->
<<"max_concurrent_guild_starts">> => 256,
<<"gateway_dispatch_relay_shards">> => 32,
<<"gateway_dispatch_relay_max_queue">> => 50000,
<<"voice_e2ee_scope">> => <<"guild_feature_only">>,
<<"voice_reconciliation_v3_percentage">> => 100,
<<"voice_reconciliation_v3_interval_ms">> => 2000
<<"voice_e2ee_scope">> => <<"guild_feature_only">>
}.
default_config_has_expected_keys_test() ->
@@ -29,9 +27,7 @@ default_config_has_expected_keys_test() ->
?assertEqual(100, maps:get(<<"guild_rollout_percentage">>, Config)),
?assertEqual(10000, maps:get(<<"rpc_request_timeout_ms">>, Config)),
?assertEqual(512, maps:get(<<"max_concurrent_session_starts">>, Config)),
?assertEqual(256, maps:get(<<"max_concurrent_guild_starts">>, Config)),
?assertEqual(100, maps:get(<<"voice_reconciliation_v3_percentage">>, Config)),
?assertEqual(2000, maps:get(<<"voice_reconciliation_v3_interval_ms">>, Config)).
?assertEqual(256, maps:get(<<"max_concurrent_guild_starts">>, Config)).
is_session_eligible_full_rollout_test() ->
persistent_term:put(?PERSISTENT_TERM_KEY, default_config()),
@@ -234,24 +230,6 @@ validate_config_rejects_rpc_timeout_above_maximum_test() ->
)
).
validate_config_rejects_voice_reconciliation_v3_interval_below_minimum_test() ->
?assertMatch(
{error, {invalid_field, <<"voice_reconciliation_v3_interval_ms">>, 499}},
gateway_rollout_config_validate:validate(
#{<<"voice_reconciliation_v3_interval_ms">> => 499},
default_config()
)
).
validate_config_rejects_voice_reconciliation_v3_percentage_above_maximum_test() ->
?assertMatch(
{error, {invalid_field, <<"voice_reconciliation_v3_percentage">>, 101}},
gateway_rollout_config_validate:validate(
#{<<"voice_reconciliation_v3_percentage">> => 101},
default_config()
)
).
validate_config_rejects_relay_max_queue_zero_test() ->
?assertMatch(
{error, {invalid_field, <<"gateway_dispatch_relay_max_queue">>, 0}},
@@ -16,8 +16,6 @@ export const GatewayRolloutConfigSchema = z.object({
gateway_dispatch_relay_shards: z.number().int().min(1).max(10000).default(32),
gateway_dispatch_relay_max_queue: z.number().int().min(0).max(1000000).default(50000),
voice_e2ee_scope: VoiceE2EEScopeEnum.default('guild_feature_only'),
voice_reconciliation_v3_percentage: z.number().min(0).max(100).default(100),
voice_reconciliation_v3_interval_ms: z.number().int().min(500).max(60000).default(2000),
});
export type GatewayRolloutConfig = z.infer<typeof GatewayRolloutConfigSchema>;
@@ -32,8 +30,6 @@ export const GatewayRolloutConfigUpdateRequest = z.object({
gateway_dispatch_relay_shards: z.number().int().min(1).max(10000).optional(),
gateway_dispatch_relay_max_queue: z.number().int().min(0).max(1000000).optional(),
voice_e2ee_scope: VoiceE2EEScopeEnum.optional(),
voice_reconciliation_v3_percentage: z.number().min(0).max(100).optional(),
voice_reconciliation_v3_interval_ms: z.number().int().min(500).max(60000).optional(),
});
export type GatewayRolloutConfigUpdateRequest = z.infer<typeof GatewayRolloutConfigUpdateRequest>;
@@ -12,19 +12,17 @@ describe('gateway rollout schemas', () => {
gateway_dispatch_relay_shards: 32,
gateway_dispatch_relay_max_queue: 50000,
voice_e2ee_scope: 'guild_feature_only',
voice_reconciliation_v3_percentage: 100,
voice_reconciliation_v3_interval_ms: 2000,
});
});
test('update request remains partial and does not inject defaults', () => {
expect(
GatewayRolloutConfigUpdateRequest.parse({
rpc_request_timeout_ms: 5000,
voice_reconciliation_v3_interval_ms: 1500,
gateway_dispatch_relay_shards: 16,
}),
).toEqual({
rpc_request_timeout_ms: 5000,
voice_reconciliation_v3_interval_ms: 1500,
gateway_dispatch_relay_shards: 16,
});
});
test('voice e2ee scope accepts only known modes', () => {