feat(voice): soft connection limits for voice servers (#2694)

This commit is contained in:
Hampus
2026-09-11 20:05:14 +02:00
committed by GitHub
parent cadca2c18e
commit 1b1d48b05e
24 changed files with 654 additions and 28 deletions
@@ -235,6 +235,7 @@ export class AdminVoiceService {
serverId: data.server_id,
endpoint: data.endpoint,
isActive: data.is_active ?? true,
softConnectionLimit: data.soft_connection_limit ?? null,
apiKey: data.api_key ?? null,
apiSecret: data.api_secret ?? null,
latitude: data.latitude ?? null,
@@ -271,6 +272,7 @@ export class AdminVoiceService {
if (data.latitude !== undefined) updates.latitude = data.latitude;
if (data.longitude !== undefined) updates.longitude = data.longitude;
if (data.is_active !== undefined) updates.isActive = data.is_active;
if (data.soft_connection_limit !== undefined) updates.softConnectionLimit = data.soft_connection_limit;
updates.restrictions = patchVoiceRestrictions(existing.restrictions, data);
updates.updatedAt = new Date();
await voiceRepository.upsertServer(updates);
@@ -339,6 +341,7 @@ export class AdminVoiceService {
latitude: server.latitude ?? null,
longitude: server.longitude ?? null,
is_active: server.isActive,
soft_connection_limit: server.softConnectionLimit ?? null,
vip_only: server.restrictions.vipOnly,
required_guild_features: Array.from(server.restrictions.requiredGuildFeatures),
allowed_guild_ids: allowedGuildIds,
@@ -266,4 +266,77 @@ describe('VoiceAdminController', () => {
expect(persisted?.apiKey).toBe(fixture.initialApiKey);
expect(persisted?.apiSecret).toBe(fixture.initialApiSecret);
});
test('stores, keeps, and clears a voice server soft connection limit', async () => {
const admin = await createAdminWithAcls(harness, [
AdminACLs.VOICE_REGION_CREATE,
AdminACLs.VOICE_SERVER_CREATE,
AdminACLs.VOICE_SERVER_LIST,
AdminACLs.VOICE_SERVER_UPDATE,
]);
const regionId = 'voice-region-soft-limit';
const serverId = 'voice-server-soft-limit';
await createBuilder<CreateVoiceRegionResponse>(harness, `${admin.token}`)
.post('/admin/voice/regions')
.body({
id: regionId,
name: `Region ${regionId}`,
emoji: ':earth_americas:',
latitude: 1,
longitude: 2,
})
.expect(HTTP_STATUS.OK)
.execute();
const created = await createBuilder<CreateVoiceServerResponse>(harness, `${admin.token}`)
.post(`/admin/voice/regions/${regionId}/servers`)
.body({
server_id: serverId,
endpoint: 'https://voice-soft-limit.example.com/socket',
api_key: 'soft-limit-api-key',
api_secret: 'soft-limit-api-secret',
soft_connection_limit: 250,
})
.expect(HTTP_STATUS.OK)
.execute();
expect(created.server.soft_connection_limit).toBe(250);
expect((await voiceRepository.getServer(regionId, serverId))?.softConnectionLimit).toBe(250);
await createBuilder<UpdateVoiceServerResponse>(harness, `${admin.token}`)
.patch(`/admin/voice/regions/${regionId}/servers/${serverId}`)
.body({endpoint: 'https://voice-soft-limit-2.example.com/socket'})
.expect(HTTP_STATUS.OK)
.execute();
expect((await voiceRepository.getServer(regionId, serverId))?.softConnectionLimit).toBe(250);
const cleared = await createBuilder<UpdateVoiceServerResponse>(harness, `${admin.token}`)
.patch(`/admin/voice/regions/${regionId}/servers/${serverId}`)
.body({soft_connection_limit: null})
.expect(HTTP_STATUS.OK)
.execute();
expect(cleared.server.soft_connection_limit).toBeNull();
expect((await voiceRepository.getServer(regionId, serverId))?.softConnectionLimit).toBeNull();
});
test('rejects a voice server soft connection limit below one', async () => {
const admin = await createAdminWithAcls(harness, [AdminACLs.VOICE_REGION_CREATE, AdminACLs.VOICE_SERVER_CREATE]);
const regionId = 'voice-region-soft-limit-invalid';
await createBuilder<CreateVoiceRegionResponse>(harness, `${admin.token}`)
.post('/admin/voice/regions')
.body({
id: regionId,
name: `Region ${regionId}`,
emoji: ':earth_americas:',
latitude: 1,
longitude: 2,
})
.expect(HTTP_STATUS.OK)
.execute();
await createBuilder(harness, `${admin.token}`)
.post(`/admin/voice/regions/${regionId}/servers`)
.body({
server_id: 'voice-server-soft-limit-invalid',
endpoint: 'https://voice-soft-limit-invalid.example.com/socket',
api_key: 'soft-limit-invalid-api-key',
api_secret: 'soft-limit-invalid-api-secret',
soft_connection_limit: 0,
})
.expect(HTTP_STATUS.BAD_REQUEST, APIErrorCodes.INVALID_FORM_BODY)
.execute();
});
});
@@ -39,6 +39,7 @@ export interface VoiceServerRow {
latitude: number | null;
longitude: number | null;
is_active: boolean | null;
soft_connection_limit: number | null;
vip_only: boolean | null;
required_guild_features: Set<string> | null;
allowed_guild_ids: Set<bigint> | null;
@@ -56,6 +57,7 @@ export const VOICE_SERVER_COLUMNS = [
'latitude',
'longitude',
'is_active',
'soft_connection_limit',
'vip_only',
'required_guild_features',
'allowed_guild_ids',
@@ -28,6 +28,7 @@ import {setInjectedSearchProvider} from '../SearchFactory';
import type {ISearchProvider} from '../search/ISearchProvider';
import {VoiceAvailabilityService} from '../voice/VoiceAvailabilityService';
import {VoiceRepository} from '../voice/VoiceRepository';
import {VoiceServerLoadTracker} from '../voice/VoiceServerLoad';
import {VoiceTopology} from '../voice/VoiceTopology';
import type {WorkerTaskName} from '../worker/WorkerLaneConfig';
@@ -289,7 +290,10 @@ export async function ensureVoiceResourcesInitialized(): Promise<void> {
const topology = new VoiceTopology(voiceRepository, voiceConfigSubscriber);
await topology.initialize();
voiceTopology = topology;
voiceAvailabilityService = new VoiceAvailabilityService(topology);
voiceAvailabilityService = new VoiceAvailabilityService(
topology,
new VoiceServerLoadTracker({gatewayService: getGatewayService()}),
);
liveKitServiceInstance = new LiveKitService(topology);
voiceRoomStoreInstance = new VoiceRoomStore(getKVClient());
})().finally(() => {
@@ -3,6 +3,8 @@
import {GuildFeatures} from '@fluxer/constants/src/GuildConstants';
import type {GuildID, UserID} from '../BrandedTypes';
import type {VoiceRegionAvailability, VoiceRegionMetadata, VoiceRegionRecord, VoiceServerRecord} from './VoiceModel';
import {preferServersUnderSoftLimit} from './VoiceRegionSelection';
import type {VoiceServerLoadSource} from './VoiceServerLoad';
import type {VoiceTopology} from './VoiceTopology';
export interface VoiceAccessContext {
@@ -11,10 +13,19 @@ export interface VoiceAccessContext {
guildFeatures?: Set<string>;
}
const EMPTY_CONNECTION_COUNTS: ReadonlyMap<string, number> = new Map();
export class VoiceAvailabilityService {
private rotationIndex: Map<string, number> = new Map();
constructor(private topology: VoiceTopology) {}
constructor(
private topology: VoiceTopology,
private loadSource: VoiceServerLoadSource | null = null,
) {}
getServerConnectionCounts(): ReadonlyMap<string, number> {
return this.loadSource?.getConnectionCounts() ?? EMPTY_CONNECTION_COUNTS;
}
getRegionMetadata(): Array<VoiceRegionMetadata> {
return this.topology.getRegionMetadataList();
@@ -140,9 +151,10 @@ export class VoiceAvailabilityService {
if (accessibleServers.length === 0) {
return null;
}
const candidateServers = preferServersUnderSoftLimit(accessibleServers, this.getServerConnectionCounts());
const index = this.rotationIndex.get(regionId) ?? 0;
const server = accessibleServers[index % accessibleServers.length];
this.rotationIndex.set(regionId, (index + 1) % accessibleServers.length);
const server = candidateServers[index % candidateServers.length];
this.rotationIndex.set(regionId, (index + 1) % candidateServers.length);
return server;
}
@@ -56,6 +56,7 @@ export class VoiceDataInitializer {
latitude: null,
longitude: null,
isActive: true,
softConnectionLimit: null,
restrictions: {
vipOnly: false,
requiredGuildFeatures: new Set(),
+1
View File
@@ -30,6 +30,7 @@ export interface VoiceServerRecord {
latitude: number | null;
longitude: number | null;
isActive: boolean;
softConnectionLimit: number | null;
restrictions: VoiceRestriction;
createdAt: Date | null;
updatedAt: Date | null;
@@ -77,12 +77,14 @@ export function selectVoiceRegionId({
export function selectClosestPseudoRegionServer({
mode,
accessibleServers,
connectionCounts,
latitude,
longitude,
selectionKey,
}: {
mode: VoiceRegionPreference['mode'];
accessibleServers: Array<VoiceServerRecord>;
connectionCounts: ReadonlyMap<string, number>;
latitude?: string;
longitude?: string;
selectionKey: string;
@@ -95,10 +97,31 @@ export function selectClosestPseudoRegionServer({
if (userLat === null || userLon === null) {
return null;
}
const closestServers = findClosestServers(accessibleServers, userLat, userLon);
const preferredServers = preferServersUnderSoftLimit(accessibleServers, connectionCounts);
const closestServers = findClosestServers(preferredServers, userLat, userLon);
return selectBalancedServer(closestServers, selectionKey);
}
export function preferServersUnderSoftLimit(
servers: Array<VoiceServerRecord>,
connectionCounts: ReadonlyMap<string, number>,
): Array<VoiceServerRecord> {
const serversUnderLimit = servers.filter((server) => !isServerAtSoftLimit(server, connectionCounts));
return serversUnderLimit.length > 0 ? serversUnderLimit : servers;
}
function isServerAtSoftLimit(server: VoiceServerRecord, connectionCounts: ReadonlyMap<string, number>): boolean {
const limit = server.softConnectionLimit;
if (limit === null || limit <= 0) {
return false;
}
const connectionCount = connectionCounts.get(server.serverId);
if (connectionCount === undefined) {
return false;
}
return connectionCount >= limit;
}
function findClosestRegionIds(
latitude: string | undefined,
longitude: string | undefined,
@@ -147,6 +147,7 @@ export class VoiceRepository implements IVoiceRepository {
latitude: server.latitude ?? null,
longitude: server.longitude ?? null,
is_active: server.isActive,
soft_connection_limit: server.softConnectionLimit ?? null,
vip_only: server.restrictions.vipOnly,
required_guild_features: new Set(server.restrictions.requiredGuildFeatures),
allowed_guild_ids: new Set(Array.from(server.restrictions.allowedGuildIds).map((id) => BigInt(id))),
@@ -190,6 +191,7 @@ export class VoiceRepository implements IVoiceRepository {
latitude: row.latitude ?? null,
longitude: row.longitude ?? null,
isActive: row.is_active ?? true,
softConnectionLimit: row.soft_connection_limit ?? null,
restrictions: {
vipOnly: row.vip_only ?? false,
requiredGuildFeatures: new Set(toIterable<string>(row.required_guild_features)),
@@ -0,0 +1,69 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import type {IGatewayService} from '../infrastructure/IGatewayService';
import {Logger} from '../Logger';
const DEFAULT_REFRESH_INTERVAL_MS = 15000;
const STALE_REFRESH_INTERVALS = 4;
const EMPTY_CONNECTION_COUNTS: ReadonlyMap<string, number> = new Map();
export interface VoiceServerLoadSource {
getConnectionCounts(): ReadonlyMap<string, number>;
}
export class VoiceServerLoadTracker implements VoiceServerLoadSource {
private readonly gatewayService: IGatewayService;
private readonly refreshIntervalMs: number;
private readonly staleAfterMs: number;
private readonly now: () => number;
private connectionCounts: ReadonlyMap<string, number> = EMPTY_CONNECTION_COUNTS;
private lastAttemptAt = 0;
private lastSuccessAt = 0;
private refreshing: Promise<void> | null = null;
constructor(options: {
gatewayService: IGatewayService;
refreshIntervalMs?: number;
now?: () => number;
}) {
this.gatewayService = options.gatewayService;
this.refreshIntervalMs = options.refreshIntervalMs ?? DEFAULT_REFRESH_INTERVAL_MS;
this.staleAfterMs = this.refreshIntervalMs * STALE_REFRESH_INTERVALS;
this.now = options.now ?? Date.now;
}
getConnectionCounts(): ReadonlyMap<string, number> {
const now = this.now();
if (now - this.lastAttemptAt >= this.refreshIntervalMs) {
void this.refresh();
}
if (this.lastSuccessAt === 0 || now - this.lastSuccessAt > this.staleAfterMs) {
return EMPTY_CONNECTION_COUNTS;
}
return this.connectionCounts;
}
async refresh(): Promise<void> {
if (this.refreshing) {
return this.refreshing;
}
this.lastAttemptAt = this.now();
this.refreshing = this.gatewayService
.getVoiceStateCounts()
.then((counts) => {
const nextCounts = new Map<string, number>();
for (const server of counts.servers) {
nextCounts.set(server.server_id, server.voice_state_count);
}
this.connectionCounts = nextCounts;
this.lastSuccessAt = this.now();
})
.catch((error) => {
Logger.warn({error}, 'Failed to refresh voice server connection counts');
})
.finally(() => {
this.refreshing = null;
});
return this.refreshing;
}
}
+1
View File
@@ -161,6 +161,7 @@ export class VoiceService {
const pseudoRegionServer = selectClosestPseudoRegionServer({
mode: regionPreference.mode,
accessibleServers,
connectionCounts: this.voiceAvailabilityService.getServerConnectionCounts(),
latitude: params.latitude,
longitude: params.longitude,
selectionKey,
@@ -38,6 +38,7 @@ function createMockServer(overrides: Partial<VoiceServerRecord> = {}): VoiceServ
latitude: null,
longitude: null,
isActive: true,
softConnectionLimit: null,
restrictions: {
vipOnly: false,
requiredGuildFeatures: new Set(),
@@ -384,5 +385,37 @@ describe('VoiceAvailabilityService', () => {
expect(first!.serverId).toBe('server-1');
expect(second!.serverId).toBe('server-2');
});
it('rotates only between servers below their soft connection limit', () => {
const region = createMockRegion();
const server1 = createMockServer({serverId: 'server-1', softConnectionLimit: 50});
const server2 = createMockServer({serverId: 'server-2'});
const topology = createMockTopology([region], new Map([['us-default', [server1, server2]]]));
service = new VoiceAvailabilityService(topology, {
getConnectionCounts: () => new Map([['server-1', 50]]),
});
const context: VoiceAccessContext = {
requestingUserId: 123n as UserID,
};
expect(service.selectServer('us-default', context)!.serverId).toBe('server-2');
expect(service.selectServer('us-default', context)!.serverId).toBe('server-2');
});
it('rotates across every server when all of them are at their soft connection limit', () => {
const region = createMockRegion();
const server1 = createMockServer({serverId: 'server-1', softConnectionLimit: 50});
const server2 = createMockServer({serverId: 'server-2', softConnectionLimit: 50});
const topology = createMockTopology([region], new Map([['us-default', [server1, server2]]]));
service = new VoiceAvailabilityService(topology, {
getConnectionCounts: () =>
new Map([
['server-1', 90],
['server-2', 90],
]),
});
const context: VoiceAccessContext = {
requestingUserId: 123n as UserID,
};
expect(service.selectServer('us-default', context)!.serverId).toBe('server-1');
expect(service.selectServer('us-default', context)!.serverId).toBe('server-2');
});
});
});
@@ -3,6 +3,7 @@
import {describe, expect, it} from 'vitest';
import type {VoiceRegionAvailability, VoiceServerRecord} from '../VoiceModel';
import {
preferServersUnderSoftLimit,
resolveVoiceRegionPreference,
selectClosestPseudoRegionServer,
selectVoiceRegionId,
@@ -45,11 +46,13 @@ function createVoiceServer({
serverId,
latitude,
longitude,
softConnectionLimit = null,
}: {
regionId: string;
serverId: string;
latitude: number | null;
longitude: number | null;
softConnectionLimit?: number | null;
}): VoiceServerRecord {
return {
regionId,
@@ -60,6 +63,7 @@ function createVoiceServer({
latitude,
longitude,
isActive: true,
softConnectionLimit,
restrictions: {
vipOnly: false,
requiredGuildFeatures: new Set(),
@@ -123,6 +127,7 @@ describe('VoiceRegionSelection', () => {
const selectedServer = selectClosestPseudoRegionServer({
mode: 'automatic',
accessibleServers: [serverA, serverB],
connectionCounts: new Map(),
latitude: '50',
longitude: '50',
selectionKey: 'guild:1:channel:1',
@@ -136,6 +141,7 @@ describe('VoiceRegionSelection', () => {
const selectedFromForwardOrder = selectClosestPseudoRegionServer({
mode: 'automatic',
accessibleServers: [serverB, serverA],
connectionCounts: new Map(),
latitude: '50',
longitude: '50',
selectionKey: 'guild:1:channel:1',
@@ -143,6 +149,7 @@ describe('VoiceRegionSelection', () => {
const selectedFromReverseOrder = selectClosestPseudoRegionServer({
mode: 'automatic',
accessibleServers: [serverA, serverB],
connectionCounts: new Map(),
latitude: '50',
longitude: '50',
selectionKey: 'guild:1:channel:1',
@@ -150,6 +157,7 @@ describe('VoiceRegionSelection', () => {
const selectedForAnotherRoom = selectClosestPseudoRegionServer({
mode: 'automatic',
accessibleServers: [serverB, serverA],
connectionCounts: new Map(),
latitude: '50',
longitude: '50',
selectionKey: 'guild:1:channel:2',
@@ -164,6 +172,7 @@ describe('VoiceRegionSelection', () => {
const selectedServer = selectClosestPseudoRegionServer({
mode: 'explicit',
accessibleServers: [serverA, serverB],
connectionCounts: new Map(),
latitude: '50',
longitude: '50',
selectionKey: 'guild:1:channel:1',
@@ -206,4 +215,91 @@ describe('VoiceRegionSelection', () => {
expect(selectedFromReverseOrder).toBe('b');
expect(selectedForAnotherRoom).toBe('a');
});
it('skips a pseudo-region server that reached its soft connection limit', () => {
const nearServer = createVoiceServer({
regionId: 'a',
serverId: 'a1',
latitude: 51,
longitude: 51,
softConnectionLimit: 100,
});
const farServer = createVoiceServer({regionId: 'b', serverId: 'b1', latitude: 0, longitude: 0});
const selectedServer = selectClosestPseudoRegionServer({
mode: 'automatic',
accessibleServers: [nearServer, farServer],
connectionCounts: new Map([['a1', 100]]),
latitude: '50',
longitude: '50',
selectionKey: 'guild:1:channel:1',
});
expect(selectedServer?.serverId).toBe('b1');
});
it('keeps a pseudo-region server that is still below its soft connection limit', () => {
const nearServer = createVoiceServer({
regionId: 'a',
serverId: 'a1',
latitude: 51,
longitude: 51,
softConnectionLimit: 100,
});
const farServer = createVoiceServer({regionId: 'b', serverId: 'b1', latitude: 0, longitude: 0});
const selectedServer = selectClosestPseudoRegionServer({
mode: 'automatic',
accessibleServers: [nearServer, farServer],
connectionCounts: new Map([['a1', 99]]),
latitude: '50',
longitude: '50',
selectionKey: 'guild:1:channel:1',
});
expect(selectedServer?.serverId).toBe('a1');
});
it('falls back to a server over its soft connection limit when every candidate is over', () => {
const serverA = createVoiceServer({
regionId: 'a',
serverId: 'a1',
latitude: 51,
longitude: 51,
softConnectionLimit: 10,
});
const serverB = createVoiceServer({
regionId: 'b',
serverId: 'b1',
latitude: 0,
longitude: 0,
softConnectionLimit: 10,
});
const selectedServer = selectClosestPseudoRegionServer({
mode: 'automatic',
accessibleServers: [serverA, serverB],
connectionCounts: new Map([
['a1', 40],
['b1', 40],
]),
latitude: '50',
longitude: '50',
selectionKey: 'guild:1:channel:1',
});
expect(selectedServer?.serverId).toBe('a1');
});
it('ignores a soft connection limit when no count is known for the server', () => {
const serverA = createVoiceServer({
regionId: 'a',
serverId: 'a1',
latitude: null,
longitude: null,
softConnectionLimit: 1,
});
const serverB = createVoiceServer({regionId: 'b', serverId: 'b1', latitude: null, longitude: null});
expect(preferServersUnderSoftLimit([serverA, serverB], new Map())).toEqual([serverA, serverB]);
});
it('ignores a soft connection limit that is not positive', () => {
const serverA = createVoiceServer({
regionId: 'a',
serverId: 'a1',
latitude: null,
longitude: null,
softConnectionLimit: 0,
});
expect(preferServersUnderSoftLimit([serverA], new Map([['a1', 500]]))).toEqual([serverA]);
});
});
@@ -0,0 +1,87 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {describe, expect, it} from 'vitest';
import type {GatewayVoiceStateCounts, IGatewayService} from '../../infrastructure/IGatewayService';
import {VoiceServerLoadTracker} from '../VoiceServerLoad';
function createGatewayService(respond: () => Promise<GatewayVoiceStateCounts>): {
gatewayService: IGatewayService;
callCount: () => number;
} {
let calls = 0;
const gatewayService = {
getVoiceStateCounts: () => {
calls += 1;
return respond();
},
} as IGatewayService;
return {gatewayService, callCount: () => calls};
}
function counts(servers: Array<{server_id: string; voice_state_count: number}>): GatewayVoiceStateCounts {
return {
total_voice_states: servers.reduce((total, server) => total + server.voice_state_count, 0),
regions: [],
servers,
};
}
describe('VoiceServerLoadTracker', () => {
it('reports no counts until the first refresh resolves', () => {
const {gatewayService} = createGatewayService(async () => counts([{server_id: 'server-1', voice_state_count: 7}]));
const tracker = new VoiceServerLoadTracker({gatewayService});
expect(tracker.getConnectionCounts().size).toBe(0);
});
it('reports the counts the gateway returned', async () => {
const {gatewayService} = createGatewayService(async () => counts([{server_id: 'server-1', voice_state_count: 7}]));
const tracker = new VoiceServerLoadTracker({gatewayService});
await tracker.refresh();
expect(tracker.getConnectionCounts().get('server-1')).toBe(7);
});
it('keeps the last counts when a refresh fails', async () => {
let shouldFail = false;
const {gatewayService} = createGatewayService(async () => {
if (shouldFail) {
throw new Error('gateway unavailable');
}
return counts([{server_id: 'server-1', voice_state_count: 7}]);
});
const tracker = new VoiceServerLoadTracker({gatewayService});
await tracker.refresh();
shouldFail = true;
await tracker.refresh();
expect(tracker.getConnectionCounts().get('server-1')).toBe(7);
});
it('refreshes no more often than the refresh interval', async () => {
let currentTime = 1000;
const {gatewayService, callCount} = createGatewayService(async () =>
counts([{server_id: 'server-1', voice_state_count: 7}]),
);
const tracker = new VoiceServerLoadTracker({
gatewayService,
refreshIntervalMs: 5000,
now: () => currentTime,
});
await tracker.refresh();
tracker.getConnectionCounts();
currentTime += 4999;
tracker.getConnectionCounts();
expect(callCount()).toBe(1);
currentTime += 1;
tracker.getConnectionCounts();
expect(callCount()).toBe(2);
});
it('drops counts that are too old to place against', async () => {
let currentTime = 1000;
const {gatewayService} = createGatewayService(async () => counts([{server_id: 'server-1', voice_state_count: 7}]));
const tracker = new VoiceServerLoadTracker({
gatewayService,
refreshIntervalMs: 5000,
now: () => currentTime,
});
await tracker.refresh();
expect(tracker.getConnectionCounts().get('server-1')).toBe(7);
currentTime += 20001;
expect(tracker.getConnectionCounts().size).toBe(0);
});
});