From e1bc6c2f7ec92622518fb475495c0ae64b2b3bc0 Mon Sep 17 00:00:00 2001 From: Hampus Date: Sat, 12 Sep 2026 01:37:45 +0200 Subject: [PATCH] refactor(voice): remove the unused voice state ack (#2716) --- .../features/gateway/events/EventRouter.ts | 2 - .../voice/engine/MediaEngineFacade.ts | 39 ------ .../MediaEngineFacadeStateMachine.test.ts | 25 ---- .../engine/MediaEngineFacadeStateMachine.ts | 12 -- .../features/voice/events/VoiceStateAck.ts | 22 --- .../src/content/docs/gateway/commands.md | 13 +- .../src/content/docs/gateway/events.md | 66 --------- .../docs/gateway/limits-and-rate-limits.md | 2 +- fluxer_docs/src/content/docs/voice/index.md | 4 +- .../voice/guild_voice_connection_apply.erl | 4 +- .../voice/guild_voice_connection_update.erl | 33 ++--- .../voice/guild_voice_connection_util.erl | 91 +----------- .../src/guild/voice/guild_voice_mutation.erl | 129 ------------------ .../src/session/session_voice_connect.erl | 14 +- .../src/session/session_voice_dispatch.erl | 32 +---- .../test/guild_voice_connection_tests.erl | 114 ---------------- 16 files changed, 21 insertions(+), 581 deletions(-) delete mode 100644 fluxer_app/src/features/voice/events/VoiceStateAck.ts delete mode 100644 fluxer_gateway/src/guild/voice/guild_voice_mutation.erl diff --git a/fluxer_app/src/features/gateway/events/EventRouter.ts b/fluxer_app/src/features/gateway/events/EventRouter.ts index 36ce99a97..115defcfa 100644 --- a/fluxer_app/src/features/gateway/events/EventRouter.ts +++ b/fluxer_app/src/features/gateway/events/EventRouter.ts @@ -67,7 +67,6 @@ import {handleCallDelete} from '@app/features/voice/events/CallDelete'; import {handleCallUpdate} from '@app/features/voice/events/CallUpdate'; import {handleEntranceSoundPlay} from '@app/features/voice/events/EntranceSoundPlay'; import {handleVoiceServerUpdate} from '@app/features/voice/events/VoiceServerUpdate'; -import {handleVoiceStateAck} from '@app/features/voice/events/VoiceStateAck'; import {handleVoiceStateUpdate} from '@app/features/voice/events/VoiceStateUpdate'; export interface GatewayGeoipPayload { @@ -143,7 +142,6 @@ export function createHandlerRegistry(): GatewayHandlerRegistry { registry.set('SAVED_MESSAGE_DELETE', handleSavedMessageDelete as GatewayEventHandler); registry.set('PRESENCE_UPDATE', handlePresenceUpdate as GatewayEventHandler); registry.set('PRESENCE_UPDATE_BULK', handlePresenceUpdateBulk as GatewayEventHandler); - registry.set('VOICE_STATE_ACK', handleVoiceStateAck as GatewayEventHandler); registry.set('VOICE_STATE_UPDATE', handleVoiceStateUpdate as GatewayEventHandler); registry.set('VOICE_SERVER_UPDATE', handleVoiceServerUpdate as GatewayEventHandler); registry.set('CALL_CREATE', handleCallCreate as GatewayEventHandler); diff --git a/fluxer_app/src/features/voice/engine/MediaEngineFacade.ts b/fluxer_app/src/features/voice/engine/MediaEngineFacade.ts index 26cbd4bb9..c10759b6d 100644 --- a/fluxer_app/src/features/voice/engine/MediaEngineFacade.ts +++ b/fluxer_app/src/features/voice/engine/MediaEngineFacade.ts @@ -36,7 +36,6 @@ import { selectMediaEngineGatewayErrorDecision, shouldCancelMediaEngineReconnectForServerVoiceStateRemoval, shouldImmediatelyDisconnectMediaEngineForServerVoiceStateRemoval, - shouldNotifyCameraUserLimitRejection, shouldRunMediaEngineDeferredDisconnect, transitionMediaEngineFacadeSnapshot, } from '@app/features/voice/engine/MediaEngineFacadeStateMachine'; @@ -47,7 +46,6 @@ import { CLAIM_YOUR_ACCOUNT_TO_START_OR_JOIN_1_DESCRIPTOR, DEFERRED_DISCONNECT_TIMEOUT_MS, RECONNECT_SUCCEEDED_PICK_A_SCREEN_AGAIN_IF_YOU_DESCRIPTOR, - VOICE_CAMERA_USER_LIMIT_REACHED_DESCRIPTOR, VOICE_CHANNEL_NO_LONGER_AVAILABLE_DESCRIPTOR, VOICE_CONNECTION_FAILED_DESCRIPTOR, VOICE_CONNECTION_LIMIT_REACHED_DESCRIPTOR, @@ -136,7 +134,6 @@ import {VoiceEngineV2AppStatsHostAdapter} from '@app/features/voice/engine/v2/Vo import VoiceEngineV2AppSubscriptionAdapter from '@app/features/voice/engine/v2/VoiceEngineV2AppSubscriptionAdapter'; import voiceEngineV2AppVoiceStateAdapter from '@app/features/voice/engine/v2/VoiceEngineV2AppVoiceStateAdapter'; import type {DisplayScreenShareCaptureContext} from '@app/features/voice/engine/voice_screen_share_manager/shared'; -import type {VoiceStateAckPayload} from '@app/features/voice/events/VoiceStateAck'; import CallMediaPrefs from '@app/features/voice/state/CallMediaPrefs'; import {type ChannelE2EEStatus, computeChannelE2EEStatus} from '@app/features/voice/state/ChannelE2EEStatus'; import LocalVoiceState from '@app/features/voice/state/LocalVoiceState'; @@ -156,7 +153,6 @@ import {getActiveVoiceProcessingMode} from '@app/features/voice/utils/VoiceProce import type {NativeAudioStartOptions} from '@app/types/electron.d'; import {ME} from '@fluxer/constants/src/AppConstants'; import {ChannelTypes} from '@fluxer/constants/src/ChannelConstants'; -import {VOICE_CHANNEL_CAMERA_USER_LIMIT} from '@fluxer/constants/src/LimitConstants'; import type { VoiceEngineV2AudioMode, VoiceEngineV2CameraOptions, @@ -1925,41 +1921,6 @@ class MediaEngineFacade extends Store { return voiceEngineV2AppVoiceStateAdapter.getAllVoiceStates(); } - handleGatewayVoiceStateAck(data: VoiceStateAckPayload): void { - if (shouldNotifyCameraUserLimitRejection({status: data.status, errorCode: data.error_code})) { - this.handleCameraUserLimitRejection(); - } - const canonicalState = data.canonical_state; - if (!canonicalState?.channel_id || !canonicalState.connection_id) { - logger.debug('Ignoring voice state ack without canonical connection state', { - mutationId: data.mutation_id, - status: data.status, - guildId: data.guild_id, - channelId: data.channel_id, - connectionId: data.connection_id, - }); - return; - } - this.applyEngineGatewayEcho( - this.toEngineGatewayVoiceState(canonicalState, canonicalState.guild_id ?? data.guild_id ?? null), - ); - } - - private handleCameraUserLimitRejection(): void { - logger.warn('Camera enable rejected by the gateway camera user limit', { - limit: VOICE_CHANNEL_CAMERA_USER_LIMIT, - }); - void this.setCameraEnabled(false, {sendUpdate: false}).catch((error) => { - logger.warn('Failed to turn the camera back off after a camera user limit rejection', {error}); - }); - if (!this.i18n) return; - this.showVoiceErrorModal( - VOICE_CAMERA_USER_LIMIT_REACHED_DESCRIPTOR, - 'voice.media-engine-facade.camera-user-limit-error-modal', - {voiceChannelCameraUserLimit: VOICE_CHANNEL_CAMERA_USER_LIMIT}, - ); - } - private toEngineGatewayVoiceState( source: { guild_id?: string | null; diff --git a/fluxer_app/src/features/voice/engine/MediaEngineFacadeStateMachine.test.ts b/fluxer_app/src/features/voice/engine/MediaEngineFacadeStateMachine.test.ts index 1a3b60768..417989721 100644 --- a/fluxer_app/src/features/voice/engine/MediaEngineFacadeStateMachine.test.ts +++ b/fluxer_app/src/features/voice/engine/MediaEngineFacadeStateMachine.test.ts @@ -10,10 +10,8 @@ import { selectMediaEngineGatewayErrorDecision, shouldCancelMediaEngineReconnectForServerVoiceStateRemoval, shouldImmediatelyDisconnectMediaEngineForServerVoiceStateRemoval, - shouldNotifyCameraUserLimitRejection, shouldRunMediaEngineDeferredDisconnect, transitionMediaEngineFacadeSnapshot, - VOICE_CAMERA_USER_LIMIT_ERROR_CODE, } from './MediaEngineFacadeStateMachine'; function connected(guildId: string | null = 'guild-1', channelId = 'channel-1') { @@ -434,27 +432,4 @@ describe('MediaEngineFacadeStateMachine', () => { }), ).toEqual({type: 'navigate-channel-gate'}); }); - - it('notifies only for rejected voice state acks carrying the camera user limit error code', () => { - expect( - shouldNotifyCameraUserLimitRejection({ - status: 'rejected', - errorCode: VOICE_CAMERA_USER_LIMIT_ERROR_CODE, - }), - ).toBe(true); - expect( - shouldNotifyCameraUserLimitRejection({ - status: 'applied', - errorCode: VOICE_CAMERA_USER_LIMIT_ERROR_CODE, - }), - ).toBe(false); - expect( - shouldNotifyCameraUserLimitRejection({ - status: 'rejected', - errorCode: 'VOICE_PERMISSION_DENIED', - }), - ).toBe(false); - expect(shouldNotifyCameraUserLimitRejection({})).toBe(false); - expect(shouldNotifyCameraUserLimitRejection({status: 'rejected'})).toBe(false); - }); }); diff --git a/fluxer_app/src/features/voice/engine/MediaEngineFacadeStateMachine.ts b/fluxer_app/src/features/voice/engine/MediaEngineFacadeStateMachine.ts index 74745c12b..3afd7ed8b 100644 --- a/fluxer_app/src/features/voice/engine/MediaEngineFacadeStateMachine.ts +++ b/fluxer_app/src/features/voice/engine/MediaEngineFacadeStateMachine.ts @@ -84,18 +84,6 @@ export interface MediaEngineFacadeGatewayErrorInput { channelId: string | null; } -export const VOICE_CAMERA_USER_LIMIT_ERROR_CODE = 'VOICE_CAMERA_USER_LIMIT'; - -export interface MediaEngineFacadeVoiceStateAckRejectionInput { - status?: string; - errorCode?: string; -} - -export function shouldNotifyCameraUserLimitRejection(input: MediaEngineFacadeVoiceStateAckRejectionInput): boolean { - if (input.status !== 'rejected') return false; - return input.errorCode === VOICE_CAMERA_USER_LIMIT_ERROR_CODE; -} - export type MediaEngineFacadeGatewayErrorDecision = | {type: 'ignore'} | { diff --git a/fluxer_app/src/features/voice/events/VoiceStateAck.ts b/fluxer_app/src/features/voice/events/VoiceStateAck.ts deleted file mode 100644 index 48114b61e..000000000 --- a/fluxer_app/src/features/voice/events/VoiceStateAck.ts +++ /dev/null @@ -1,22 +0,0 @@ -// SPDX-License-Identifier: AGPL-3.0-or-later - -import type {GatewayHandlerContext} from '@app/features/gateway/events/EventRouter'; -import type {VoiceState} from '@app/features/gateway/types/GatewayVoiceTypes'; -import MediaEngine from '@app/features/voice/engine/MediaEngineFacade'; - -export interface VoiceStateAckPayload { - mutation_id?: string; - runtime_epoch?: string | null; - connection_id?: string | null; - guild_id?: string | null; - channel_id?: string | null; - status?: string; - server_version?: number; - canonical_state?: VoiceState | null; - error_code?: string; - error_message?: string; -} - -export function handleVoiceStateAck(data: VoiceStateAckPayload, _context: GatewayHandlerContext): void { - MediaEngine.handleGatewayVoiceStateAck(data); -} diff --git a/fluxer_docs/src/content/docs/gateway/commands.md b/fluxer_docs/src/content/docs/gateway/commands.md index f6d9c517a..66aa877ae 100644 --- a/fluxer_docs/src/content/docs/gateway/commands.md +++ b/fluxer_docs/src/content/docs/gateway/commands.md @@ -15,7 +15,7 @@ Except for [Heartbeat](#heartbeat), [Identify](#identify), and [Resume](#resume) | 1 | Heartbeat | Opcode `11` Heartbeat ACK | | 2 | Identify | [Ready](/gateway/events/#ready), a close frame, or silence when the payload is held or discarded | | 3 | Presence Update | No direct response | -| 4 | Voice State Update | [Voice State Ack](/gateway/events/#voice-state-ack), [Voice State Update](/gateway/events/#voice-state-update), and [Voice Server Update](/gateway/events/#voice-server-update) when state changes | +| 4 | Voice State Update | [Voice State Update](/gateway/events/#voice-state-update) and [Voice Server Update](/gateway/events/#voice-server-update) when state changes | | 6 | Resume | Replayed Dispatches followed by [Resumed](/gateway/events/#resumed), Invalid Session, or a close frame | | 8 | Request Guild Members | One or more [Guild Members Chunk](/gateway/events/#guild-members-chunk) events | | 14 | Lazy Request | [Guild Sync](/gateway/events/#guild-sync) and [Guild Member List Update](/gateway/events/#guild-member-list-update) | @@ -265,9 +265,6 @@ Opcode `4` joins, moves, updates, or leaves the voice membership associated with | viewer_stream_keys?3 | ?array[string] | The stream keys this connection is watching, where an omitted key keeps the current list and null clears it | | latitude? | number or string | The client latitude, used to pick a voice region | | longitude? | number or string | The client longitude, used to pick a voice region | -| mutation_id? | string | A client-generated identity echoed in [Voice State Ack](/gateway/events/#voice-state-ack) | -| runtime_epoch? | string | The client runtime generation echoed in Voice State Ack | -| base_version?4 | integer | The voice state version this update was computed against | 1 A `channel_id` with no `connection_id` opens a new connection, and a `channel_id` with one updates or moves that connection. An update that leaves one guild, meaning a non-null `guild_id` with `channel_id: null`, requires a `connection_id`, and one that omits it is refused with `VOICE_MISSING_CONNECTION_ID` @@ -288,9 +285,7 @@ Opcode `4` joins, moves, updates, or leaves the voice membership associated with "self_mute": false, "self_deaf": false, "self_video": false, - "self_stream": false, - "mutation_id": "e0a1c2", - "base_version": 7 + "self_stream": false } } ``` @@ -305,10 +300,6 @@ The command has no `session_id` field. The current Gateway session is the member Joining or replacing a grant produces [Voice Server Update](/gateway/events/#voice-server-update) with the token and endpoint for the media connection, and [Voice State Update](/gateway/events/#voice-state-update) for every session that can see the channel. -When a guild update has `mutation_id`, Fluxer also reports the outcome to the requesting session as [Voice State Ack](/gateway/events/#voice-state-ack), whose `status` is `applied` or `rejected` and whose `error_code` names the exact refusal. Without `mutation_id` a refusal produces no event at all, and the DM and group DM call context never acks. - -`base_version` is a staleness check for an update that names an existing guild connection. An update whose `base_version` is more than one behind that connection's current voice state version is rejected with `stale_base_version`. The check runs after the member, channel, and connection lookups and before the permission checks. It does not apply to opening a new connection, to leaving a channel, or to the DM and group DM call context. - The first two updates in a rolling one-second window are processed immediately. Later updates enter a [per-session queue](/gateway/limits-and-rate-limits/#connection-and-command-rate-limits) that holds at most 64 commands and drains one command every 500 ms. A newer update replaces an older queued update for the same `guild_id` and `connection_id` pair, and a full queue discards its oldest entry before accepting the new one. ## Request Guild Members diff --git a/fluxer_docs/src/content/docs/gateway/events.md b/fluxer_docs/src/content/docs/gateway/events.md index e33c553d9..a76bef283 100644 --- a/fluxer_docs/src/content/docs/gateway/events.md +++ b/fluxer_docs/src/content/docs/gateway/events.md @@ -105,7 +105,6 @@ A Dispatch is buffered for [Resume](/gateway/commands/#resume) replay unless it | [Channel Pins Update](#channel-pins-update) | A channel's most recent pin time changes | Channel visibility | | [Channel Pins ACK](#channel-pins-ack) | The current user acknowledges a channel's pins | Current user | | [Voice State Update](#voice-state-update) | A guild or call participant's voice state changes | Channel visibility | -| [Voice State Ack](#voice-state-ack) | The session's own voice mutation is applied or rejected | Current session | | [Voice Server Update](#voice-server-update) | The session receives or replaces its own voice grant | Current session | | [Entrance Sound Play](#entrance-sound-play) | A participant's entrance sound plays in a voice channel | Voice channel | | [Call Create](#call-create) | A private channel call begins or becomes visible | Call recipient | @@ -942,71 +941,6 @@ A `channel_id` of null means the participant left. The broadcast form has no `region_id`, `server_id`, `latitude`, or `longitude`. -### VOICE_STATE_ACK - -Reports the outcome of the session's own [Voice State Update](/gateway/commands/#voice-state-update). Sent only when that command supplied `mutation_id`, and delivered to the requesting session alone. - -| Field | Type | Description | -| --- | --- | --- | -| mutation_id | string | The `mutation_id` the command supplied | -| runtime_epoch1 | string | The `runtime_epoch` the command supplied | -| connection_id | ?string | Voice connection the mutation applied to | -| guild_id | ?snowflake | Guild the mutation applied to | -| channel_id | ?snowflake | Channel the mutation applied to | -| status | string | `applied` or `rejected` | -| server_version | integer | The voice state version after the mutation | -| canonical_state2 | [voice state object](#voice-state-object) | The authoritative voice state, an empty object when none exists | -| error_code? | string | Stable rejection code, present only when `status` is `rejected` | -| error_message? | string | Human-readable rejection message | - -1 Echoed back unchanged. A command that omitted the field produces the literal string `undefined` here, so a client MUST compare the value against the epoch it sent - -2 Also has the `region_id` and `server_id` fields a [Voice State Update](#voice-state-update) omits, and never has coordinates. Both extra fields are internal routing identity, and a client MUST NOT depend on them - -A mutation whose `base_version` is more than one behind the server's current version is rejected after the connection lookup and before the permission checks, with `error_code` and `error_message` both set to `stale_base_version`. - -Every other rejection has one of these codes, with `error_message` set to the registry text for the same code. - -#### `VOICE_CONNECTION_NOT_FOUND` - -The named voice connection does not exist and no matching pending connection could be restored. - -#### `VOICE_PENDING_EXPIRED` - -The pending voice connection expired before the mutation arrived. - -#### `VOICE_INVALID_STATE` - -A `viewer_stream_keys` entry is malformed, names another channel, or names a connection that is not in the channel. - -#### `VOICE_MEMBER_TIMED_OUT`1 - -The member is timed out. - -#### `VOICE_PERMISSION_DENIED`1 - -The user lacks `VIEW_CHANNEL` or `CONNECT` on the target channel. - -#### `VOICE_CHANNEL_FULL`1 - -The channel is at its `user_limit`. - -#### `VOICE_CONNECTION_LIMIT_REACHED`1 - -The user holds too many voice connections. - -#### `VOICE_CAMERA_USER_LIMIT` - -The channel already has 25 users with cameras enabled. - -#### `VOICE_E2EE_REQUIRED`1 - -The channel is end-to-end encrypted and the client does not support it. - -1 Checked only when the mutation moves the connection to a different channel. A mutation that keeps the connection in its current channel skips these checks and is applied - -Only a command that supplies `connection_id` can produce an ack. A command that omits it opens a new connection, and a `connection_id` that belongs to another user is rejected before any other check. A refusal of either kind produces no Dispatch. Without `mutation_id`, a refused voice state update also produces no Dispatch at all. - ### VOICE_SERVER_UPDATE The session received or replaced its own voice grant. Delivered to the requesting session alone. diff --git a/fluxer_docs/src/content/docs/gateway/limits-and-rate-limits.md b/fluxer_docs/src/content/docs/gateway/limits-and-rate-limits.md index 6b80a0084..ae13aeb24 100644 --- a/fluxer_docs/src/content/docs/gateway/limits-and-rate-limits.md +++ b/fluxer_docs/src/content/docs/gateway/limits-and-rate-limits.md @@ -110,7 +110,7 @@ Identify accepts at most 256 `ignored_events` entries. A longer array closes wit ## Voice admission -Voice admission follows the enclosing guild, DM, group DM, channel, and permission rules. A refusal is reported as [Voice State Ack](/gateway/events/#voice-state-ack) with an `error_code` and closes nothing. A Voice State Update that has no `mutation_id` produces no event when it is refused. +Voice admission follows the enclosing guild, DM, group DM, channel, and permission rules. A refusal produces no event and closes nothing. A guild voice channel admits at most `user_limit` users, where `0` means unlimited. A channel in which any participant has a camera enabled also admits at most 25 users in total, whatever its `user_limit`. A channel that already holds 25 users with cameras enabled refuses a further camera with `VOICE_CAMERA_USER_LIMIT`. [Capacity](/voice/#capacity) states these bounds in full. diff --git a/fluxer_docs/src/content/docs/voice/index.md b/fluxer_docs/src/content/docs/voice/index.md index c11a370c4..ece0db22c 100644 --- a/fluxer_docs/src/content/docs/voice/index.md +++ b/fluxer_docs/src/content/docs/voice/index.md @@ -71,9 +71,7 @@ A grant is issued when a connection opens, when it moves to another channel, and The grant `token` is issued for the media server and consumed by the media connection alone. No route on this API accepts it. -Fluxer reports a guild refusal as [Voice State Ack](/gateway/events/#voice-state-ack) with a `status` of `rejected` and an `error_code` naming the exact reason, and only a command with a `mutation_id` receives one. A refused placement without it produces no Dispatch, so a client that needs to observe a guild refusal MUST send `mutation_id`. - -A call never acks. A refused placement into a direct message or group direct message call produces no Dispatch whether or not the command had `mutation_id`. A client observes that refusal only as the absence of a grant. +Fluxer reports a refusal by sending no Dispatch. A client observes a refused placement only as the absence of a grant. ## Guild voice channels diff --git a/fluxer_gateway/src/guild/voice/guild_voice_connection_apply.erl b/fluxer_gateway/src/guild/voice/guild_voice_connection_apply.erl index 606872b90..cc1ba2c7c 100644 --- a/fluxer_gateway/src/guild/voice/guild_voice_connection_apply.erl +++ b/fluxer_gateway/src/guild/voice/guild_voice_connection_apply.erl @@ -140,6 +140,4 @@ apply_same_channel_update(Update, ParsedViewerKey) -> needs_token => false, viewer_stream_keys => ParsedViewerKey }), - guild_voice_connection_util:applied_mutation_reply( - UpdateResult, Context, ChannelIdValue - ). + UpdateResult. diff --git a/fluxer_gateway/src/guild/voice/guild_voice_connection_update.erl b/fluxer_gateway/src/guild/voice/guild_voice_connection_update.erl index 2e97b40d9..527ec23e5 100644 --- a/fluxer_gateway/src/guild/voice/guild_voice_connection_update.erl +++ b/fluxer_gateway/src/guild/voice/guild_voice_connection_update.erl @@ -181,28 +181,17 @@ normal_update( ) -> ConnectionId = maps:get(connection_id, Context), ExistingVS = maps:get(ConnectionId, VoiceStates, #{}), - CurrentVersion = voice_state_utils:voice_state_version(ExistingVS), - MutationDecision = guild_voice_mutation:evaluate( - maps:get(base_version, Context, undefined), CurrentVersion, valid - ), - case MutationDecision of - {reject, Reason} -> - guild_voice_connection_util:rejected_mutation_reply( - Context, ExistingVS, State, ChannelIdValue, <<"rejected">>, Reason - ); - apply -> - check_permissions( - Context, - ChannelIdValue, - Member, - Channel, - VoiceStates, - State, - IsChannelChange, - ViewerKeyResult, - ExistingVS - ) - end. + check_permissions( + Context, + ChannelIdValue, + Member, + Channel, + VoiceStates, + State, + IsChannelChange, + ViewerKeyResult, + ExistingVS + ). -spec check_permissions( context(), diff --git a/fluxer_gateway/src/guild/voice/guild_voice_connection_util.erl b/fluxer_gateway/src/guild/voice/guild_voice_connection_util.erl index 795982fc5..ed14f0f03 100644 --- a/fluxer_gateway/src/guild/voice/guild_voice_connection_util.erl +++ b/fluxer_gateway/src/guild/voice/guild_voice_connection_util.erl @@ -19,9 +19,7 @@ maybe_attach_e2ee_key_to_reply/2, normalize_session_id/1, normalize_optional_binary/1, - applied_mutation_reply/3, maybe_error_reply/5, - rejected_mutation_reply/6, clear_virtual_access_flags/2 ]). @@ -62,9 +60,6 @@ build_context(Request0) -> viewer_stream_keys => maps:get(viewer_stream_keys, Request, undefined), latitude => Coord(maps:get(latitude, Request, undefined)), longitude => Coord(maps:get(longitude, Request, undefined)), - mutation_id => maps:get(mutation_id, Request, undefined), - runtime_epoch => maps:get(runtime_epoch, Request, undefined), - base_version => maps:get(base_version, Request, undefined), e2ee_capable => Norm(maps:get(e2ee_capable, Request, false)), bot => Norm(maps:get(bot, Request, false)) }. @@ -166,72 +161,10 @@ normalize_session_id(Value) -> normalize_optional_binary(Value) -> guild_voice_connection_normalize:normalize_optional_binary(Value). --spec applied_mutation_reply({reply, map(), guild_state()}, context(), integer()) -> - {reply, map(), guild_state()}. -applied_mutation_reply({reply, BaseReply, NewState}, Context, ChannelIdValue) -> - GuildId = guild_voice_connection_normalize:normalize_positive_snowflake( - maps:get(id, NewState, undefined) - ), - NewVoiceState = maps:get(voice_state, BaseReply, #{}), - NewVersion = voice_state_utils:voice_state_version(NewVoiceState), - NormalizedConnId = normalize_conn_id_for_ack(Context), - Ack = guild_voice_mutation:build_ack( - maps:get(mutation_id, Context, undefined), - maps:get(runtime_epoch, Context, undefined), - NormalizedConnId, - GuildId, - ChannelIdValue, - #{ - status => <<"applied">>, - server_version => NewVersion, - canonical_state => voice_state_utils:external_voice_state(NewVoiceState) - } - ), - Reply = merge_ack(BaseReply, Ack), - {reply, Reply, NewState}. - -spec maybe_error_reply(context(), voice_state(), guild_state(), integer(), atom()) -> {reply, map(), guild_state()} | {reply, {error, atom(), atom()}, guild_state()}. -maybe_error_reply(Context, ExistingVoiceState, State, ChannelIdValue, ErrorAtom) -> - case maps:get(mutation_id, Context, undefined) of - undefined -> - {reply, gateway_errors:error(ErrorAtom), State}; - _ -> - rejected_mutation_reply( - Context, ExistingVoiceState, State, ChannelIdValue, <<"rejected">>, ErrorAtom - ) - end. - --spec rejected_mutation_reply( - context(), voice_state(), guild_state(), integer(), binary(), atom() | binary() -) -> {reply, map(), guild_state()}. -rejected_mutation_reply(Context, ExistingVS, State, ChannelIdValue, Status, Error) -> - GuildId = guild_voice_connection_normalize:normalize_positive_snowflake( - maps:get(id, State, undefined) - ), - CurrentVersion = voice_state_utils:voice_state_version(ExistingVS), - NormalizedConnId = normalize_conn_id_for_ack(Context), - Ack = guild_voice_mutation:build_ack( - maps:get(mutation_id, Context, undefined), - maps:get(runtime_epoch, Context, undefined), - NormalizedConnId, - GuildId, - ChannelIdValue, - #{ - status => Status, - server_version => CurrentVersion, - canonical_state => canonical_state_for_ack(ExistingVS), - error_code => rejection_error_code(Error), - error_message => rejection_error_message(Error) - } - ), - {reply, #{success => false, ack => Ack}, State}. - --spec canonical_state_for_ack(voice_state()) -> map(). -canonical_state_for_ack(VoiceState) when is_map(VoiceState), map_size(VoiceState) > 0 -> - voice_state_utils:external_voice_state(VoiceState); -canonical_state_for_ack(_) -> - #{}. +maybe_error_reply(_Context, _ExistingVoiceState, State, _ChannelIdValue, ErrorAtom) -> + {reply, gateway_errors:error(ErrorAtom), State}. -spec clear_virtual_access_flags(voice_state(), guild_state()) -> guild_state(). clear_virtual_access_flags(VoiceState, State) when is_map(VoiceState) -> @@ -244,26 +177,6 @@ clear_virtual_access_flags(VoiceState, State) when is_map(VoiceState) -> State end. --spec normalize_conn_id_for_ack(context()) -> binary() | null. -normalize_conn_id_for_ack(Context) -> - case maps:get(connection_id, Context, undefined) of - undefined -> null; - C -> C - end. - --spec merge_ack(map(), term()) -> map(). -merge_ack(BaseReply, undefined) -> BaseReply; -merge_ack(BaseReply, Ack) -> BaseReply#{ack => Ack}. - --spec rejection_error_code(atom() | binary()) -> binary(). -rejection_error_code(ErrorAtom) when is_atom(ErrorAtom) -> gateway_errors:error_code(ErrorAtom); -rejection_error_code(ErrorCode) when is_binary(ErrorCode) -> ErrorCode. - --spec rejection_error_message(atom() | binary()) -> binary(). -rejection_error_message(ErrorAtom) when is_atom(ErrorAtom) -> - gateway_errors:error_message(ErrorAtom); -rejection_error_message(ErrorMessage) when is_binary(ErrorMessage) -> ErrorMessage. - -spec guild_data(guild_state()) -> map(). guild_data(State) -> map_utils:ensure_map(maps:get(data, State, #{})). diff --git a/fluxer_gateway/src/guild/voice/guild_voice_mutation.erl b/fluxer_gateway/src/guild/voice/guild_voice_mutation.erl deleted file mode 100644 index 2ed9e40c6..000000000 --- a/fluxer_gateway/src/guild/voice/guild_voice_mutation.erl +++ /dev/null @@ -1,129 +0,0 @@ -%% SPDX-License-Identifier: AGPL-3.0-or-later - --module(guild_voice_mutation). --typing([eqwalizer]). - --export([evaluate/3, build_ack/6]). - --ifdef(TEST). --include_lib("eunit/include/eunit.hrl"). --endif. - --spec evaluate( - BaseVersion :: integer() | undefined, - CurrentVersion :: integer(), - ValidationResult :: valid | invalid -) -> apply | {reject, binary()}. -evaluate(_BaseVersion, _CurrentVersion, invalid) -> - {reject, <<"invalid_payload">>}; -evaluate(BaseVersion, CurrentVersion, valid) when - is_integer(BaseVersion), BaseVersion < CurrentVersion - 1 --> - {reject, <<"stale_base_version">>}; -evaluate(_BaseVersion, _CurrentVersion, valid) -> - apply. - --spec build_ack( - MutationId :: binary() | undefined, - RuntimeEpoch :: binary() | undefined, - ConnectionId :: binary() | null, - GuildId :: integer() | null | undefined, - ChannelId :: integer() | null | undefined, - Outcome :: #{ - status := binary(), - server_version := integer(), - canonical_state := map(), - error_code => binary(), - error_message => binary() - } -) -> map() | undefined. -build_ack(undefined, _RuntimeEpoch, _ConnectionId, _GuildId, _ChannelId, _Outcome) -> - undefined; -build_ack(MutationId, RuntimeEpoch, ConnectionId, GuildId, ChannelId, Outcome) when - is_binary(MutationId) --> - BaseAck = #{ - <<"mutation_id">> => MutationId, - <<"runtime_epoch">> => RuntimeEpoch, - <<"connection_id">> => ConnectionId, - <<"guild_id">> => maybe_int_to_binary(GuildId), - <<"channel_id">> => maybe_int_to_binary(ChannelId), - <<"status">> => maps:get(status, Outcome), - <<"server_version">> => maps:get(server_version, Outcome), - <<"canonical_state">> => maps:get(canonical_state, Outcome) - }, - put_optional_ack_fields(BaseAck, Outcome). - --spec maybe_int_to_binary(integer() | null | undefined) -> binary() | null. -maybe_int_to_binary(null) -> null; -maybe_int_to_binary(undefined) -> null; -maybe_int_to_binary(N) when is_integer(N) -> integer_to_binary(N). - --spec maybe_put_optional_ack_field(map(), binary(), binary() | undefined) -> map(). -maybe_put_optional_ack_field(Ack, _Field, undefined) -> - Ack; -maybe_put_optional_ack_field(Ack, Field, Value) -> - Ack#{Field => Value}. - --spec put_optional_ack_fields(map(), map()) -> map(). -put_optional_ack_fields(BaseAck, Outcome) -> - lists:foldl( - fun({Field, Key}, Ack) -> - maybe_put_optional_ack_field(Ack, Field, maps:get(Key, Outcome, undefined)) - end, - BaseAck, - [ - {<<"error_code">>, error_code}, - {<<"error_message">>, error_message} - ] - ). - --ifdef(TEST). - -evaluate_invalid_payload_test() -> - ?assertEqual({reject, <<"invalid_payload">>}, evaluate(0, 0, invalid)), - ?assertEqual({reject, <<"invalid_payload">>}, evaluate(undefined, 5, invalid)). - -evaluate_apply_no_version_test() -> - ?assertEqual(apply, evaluate(undefined, 0, valid)), - ?assertEqual(apply, evaluate(undefined, 100, valid)). - -evaluate_apply_current_version_test() -> - ?assertEqual(apply, evaluate(5, 5, valid)). - -evaluate_apply_one_behind_test() -> - ?assertEqual(apply, evaluate(4, 5, valid)). - -evaluate_stale_base_test() -> - ?assertEqual({reject, <<"stale_base_version">>}, evaluate(3, 5, valid)), - ?assertEqual({reject, <<"stale_base_version">>}, evaluate(0, 10, valid)). - -build_ack_no_mutation_id_test() -> - Outcome = #{status => <<"ignored">>, server_version => 0, canonical_state => #{}}, - ?assertEqual(undefined, build_ack(undefined, undefined, <<"conn1">>, 123, 456, Outcome)). - -test_build_ack(Outcome) -> - build_ack(<<"mut1">>, <<"epoch1">>, <<"conn1">>, 123, 456, Outcome). - -build_ack_full_test() -> - Outcome = #{status => <<"applied">>, server_version => 7, canonical_state => #{}}, - #{ - <<"mutation_id">> := <<"mut1">>, - <<"status">> := <<"applied">>, - <<"server_version">> := 7 - } = test_build_ack(Outcome). - -build_ack_with_error_fields_test() -> - Outcome = #{ - status => <<"rejected">>, - server_version => 9, - canonical_state => #{}, - error_code => <<"stale_base_version">>, - error_message => <<"stale_base_version">> - }, - #{ - <<"error_code">> := <<"stale_base_version">>, - <<"error_message">> := <<"stale_base_version">> - } = test_build_ack(Outcome). - --endif. diff --git a/fluxer_gateway/src/session/session_voice_connect.erl b/fluxer_gateway/src/session/session_voice_connect.erl index e15cc56ab..a1b77beb2 100644 --- a/fluxer_gateway/src/session/session_voice_connect.erl +++ b/fluxer_gateway/src/session/session_voice_connect.erl @@ -53,12 +53,6 @@ handle_voice_state_update(Data, State) -> -spec extract_voice_params(map()) -> map(). extract_voice_params(Data) -> - BaseVersionRaw = maps:get(<<"base_version">>, Data, undefined), - BaseVersion = - case BaseVersionRaw of - V when is_integer(V), V >= 0 -> V; - _ -> undefined - end, #{ guild_id_raw => maps:get(<<"guild_id">>, Data, null), channel_id_raw => maps:get(<<"channel_id">>, Data, null), @@ -70,10 +64,7 @@ extract_voice_params(Data) -> viewer_stream_keys => maps:get(<<"viewer_stream_keys">>, Data, undefined), is_mobile => maps:get(<<"is_mobile">>, Data, false), latitude => maps:get(<<"latitude">>, Data, null), - longitude => maps:get(<<"longitude">>, Data, null), - mutation_id => maps:get(<<"mutation_id">>, Data, undefined), - runtime_epoch => maps:get(<<"runtime_epoch">>, Data, undefined), - base_version => BaseVersion + longitude => maps:get(<<"longitude">>, Data, null) }. -spec dispatch_validated( @@ -510,9 +501,6 @@ build_guild_request(ChId, Params, UserId, SId, E2EE, Bot) -> is_mobile => maps:get(is_mobile, Params), latitude => maps:get(latitude, Params), longitude => maps:get(longitude, Params), - mutation_id => maps:get(mutation_id, Params), - runtime_epoch => maps:get(runtime_epoch, Params), - base_version => maps:get(base_version, Params), e2ee_capable => E2EE, bot => Bot }. diff --git a/fluxer_gateway/src/session/session_voice_dispatch.erl b/fluxer_gateway/src/session/session_voice_dispatch.erl index 73d1cc3bc..90a267c8e 100644 --- a/fluxer_gateway/src/session/session_voice_dispatch.erl +++ b/fluxer_gateway/src/session/session_voice_dispatch.erl @@ -102,14 +102,6 @@ handle_guild_reply_ok(Reply, Ctx, SessionPid) -> maybe_dispatch_voice_server_update( Reply, GId, ChId, SessionPid ), - case maps:get(ack, Reply, undefined) of - Ack when is_map(Ack) -> - dispatch_to_session( - SessionPid, voice_state_ack, Ack, GId - ); - _ -> - ok - end, ok. -spec maybe_dispatch_voice_server_update( @@ -305,9 +297,9 @@ dispatch_to_session_converts_wire_payload_test() -> <<"permissions">> => 8, <<"roles">> => [456] }, - ok = dispatch_to_session(self(), voice_state_ack, Payload, 123), + ok = dispatch_to_session(self(), voice_state_update, Payload, 123), receive - {'$gen_cast', {dispatch, voice_state_ack, WirePayload}} -> + {'$gen_cast', {dispatch, voice_state_update, WirePayload}} -> ?assertEqual( #{ <<"id">> => <<"123">>, @@ -362,26 +354,6 @@ plain_success_reply_does_not_dispatch_test() -> ok = handle_guild_reply_ok(#{success => true}, test_voice_ctx(456), self()), assert_no_dispatch(). -rejected_mutation_reply_dispatches_ack_only_test() -> - Ack = #{<<"status">> => <<"rejected">>, <<"mutation_id">> => <<"m1">>}, - Reply = #{success => false, ack => Ack}, - ok = handle_guild_reply_ok(Reply, test_voice_ctx(456), self()), - Payload = receive_dispatch(voice_state_ack), - ?assertEqual(<<"rejected">>, maps:get(<<"status">>, Payload)), - assert_no_dispatch(). - -in_channel_update_reply_dispatches_ack_test() -> - Ack = #{<<"status">> => <<"applied">>, <<"mutation_id">> => <<"m2">>}, - Reply = #{ - success => true, - voice_state => #{<<"channel_id">> => <<"456">>}, - ack => Ack - }, - ok = handle_guild_reply_ok(Reply, test_voice_ctx(456), self()), - Payload = receive_dispatch(voice_state_ack), - ?assertEqual(<<"applied">>, maps:get(<<"status">>, Payload)), - assert_no_dispatch(). - join_reply_dispatches_voice_server_update_test() -> Reply = #{ success => true, diff --git a/fluxer_gateway/test/guild_voice_connection_tests.erl b/fluxer_gateway/test/guild_voice_connection_tests.erl index 2809ec761..7bc5e0e9a 100644 --- a/fluxer_gateway/test/guild_voice_connection_tests.erl +++ b/fluxer_gateway/test/guild_voice_connection_tests.erl @@ -91,22 +91,6 @@ voice_state_update_connection_not_found_test() -> Request, State, {error, not_found, voice_connection_not_found} ). -voice_state_update_connection_not_found_returns_rejected_ack_test() -> - State = base_test_state(), - Request = #{ - user_id => 10, - channel_id => 100, - connection_id => <<"missing-conn">>, - mutation_id => <<"m-missing-connection">>, - runtime_epoch => <<"epoch-1">>, - base_version => 0 - }, - Ack = rejected_voice_state_ack(Request, State), - ?assertEqual(<<"rejected">>, maps:get(<<"status">>, Ack)), - ?assertEqual(0, maps:get(<<"server_version">>, Ack)), - ?assertEqual(#{}, maps:get(<<"canonical_state">>, Ack)), - ?assertEqual(<<"VOICE_CONNECTION_NOT_FOUND">>, maps:get(<<"error_code">>, Ack)). - voice_state_update_invalid_viewer_stream_keys_test() -> State = connected_state(<<"10">>), Request = #{ @@ -147,98 +131,6 @@ voice_state_update_rejects_dm_scope_viewer_stream_key_test() -> Request, State, {error, validation_error, voice_invalid_state} ). -voice_state_update_viewer_stream_keys_missing_connection_returns_rejected_ack_test() -> - VoiceStates = #{ - <<"conn-1">> => #{ - <<"channel_id">> => <<"100">>, - <<"connection_id">> => <<"conn-1">>, - <<"user_id">> => <<"10">>, - <<"version">> => 2 - } - }, - State = maps:put(voice_states, VoiceStates, base_test_state()), - Request = #{ - user_id => 10, - channel_id => 100, - connection_id => <<"conn-1">>, - mutation_id => <<"m-missing-watched">>, - runtime_epoch => <<"epoch-1">>, - base_version => 2, - viewer_stream_keys => [<<"999:100:missing-conn">>] - }, - Ack = rejected_voice_state_ack(Request, State), - ?assertEqual(<<"rejected">>, maps:get(<<"status">>, Ack)), - ?assertEqual(<<"VOICE_CONNECTION_NOT_FOUND">>, maps:get(<<"error_code">>, Ack)). - -voice_state_update_stale_base_version_returns_rejected_ack_test() -> - VoiceStates = #{ - <<"conn-1">> => #{ - <<"channel_id">> => <<"100">>, - <<"connection_id">> => <<"conn-1">>, - <<"user_id">> => <<"10">>, - <<"version">> => 5 - } - }, - State = maps:put(voice_states, VoiceStates, base_test_state()), - Request = #{ - user_id => 10, - channel_id => 100, - connection_id => <<"conn-1">>, - mutation_id => <<"m1">>, - runtime_epoch => <<"epoch-1">>, - base_version => 3 - }, - Ack = rejected_voice_state_ack(Request, State), - ?assertEqual(<<"rejected">>, maps:get(<<"status">>, Ack)), - ?assertEqual(5, maps:get(<<"server_version">>, Ack)), - ?assertEqual(<<"stale_base_version">>, maps:get(<<"error_code">>, Ack)). - -voice_state_update_stale_base_version_no_superseded_arm_test() -> - VoiceStates = #{ - <<"conn-2">> => #{ - <<"channel_id">> => <<"100">>, - <<"connection_id">> => <<"conn-2">>, - <<"user_id">> => <<"10">>, - <<"version">> => 10 - } - }, - State = maps:put(voice_states, VoiceStates, base_test_state()), - Request = #{ - user_id => 10, - channel_id => 100, - connection_id => <<"conn-2">>, - mutation_id => <<"m-reg">>, - runtime_epoch => <<"epoch-reg">>, - base_version => 7 - }, - Ack = rejected_voice_state_ack(Request, State), - ?assertEqual(<<"rejected">>, maps:get(<<"status">>, Ack)), - ?assertEqual(10, maps:get(<<"server_version">>, Ack)), - ?assertEqual(<<"stale_base_version">>, maps:get(<<"error_code">>, Ack)). - -voice_state_update_invalid_viewer_stream_keys_returns_rejected_ack_test() -> - VoiceStates = #{ - <<"conn-1">> => #{ - <<"channel_id">> => <<"100">>, - <<"connection_id">> => <<"conn-1">>, - <<"user_id">> => <<"10">>, - <<"version">> => 2 - } - }, - State = maps:put(voice_states, VoiceStates, base_test_state()), - Request = #{ - user_id => 10, - channel_id => 100, - connection_id => <<"conn-1">>, - mutation_id => <<"m2">>, - runtime_epoch => <<"epoch-1">>, - base_version => 2, - viewer_stream_keys => 123 - }, - Ack = rejected_voice_state_ack(Request, State), - ?assertEqual(<<"rejected">>, maps:get(<<"status">>, Ack)), - ?assertEqual(<<"VOICE_INVALID_STATE">>, maps:get(<<"error_code">>, Ack)). - voice_state_update_guild_id_missing_test() -> State0 = base_test_state(), State1 = replace_guild_id(State0, undefined), @@ -487,9 +379,3 @@ viewer_keys(State) -> assert_voice_state_update_error(Request, State, Error) -> {reply, Error, _} = guild_voice_connection:voice_state_update(Request, State). - --spec rejected_voice_state_ack(map(), map()) -> map(). -rejected_voice_state_ack(Request, State) -> - {reply, #{ack := #{} = Ack, success := false}, _} = - guild_voice_connection:voice_state_update(Request, State), - Ack.