This repository has no description
Something went wrong. Try again.
22 kB · 695 lines
TypeScript
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696// SPDX-License-Identifier: AGPL-3.0-or-later
import type {WebhookEvent} from 'livekit-server-sdk';import {TrackSource, WebhookReceiver} from 'livekit-server-sdk';import type {ChannelID, GuildID, UserID} from '../BrandedTypes';import {Config} from '../Config';import {Logger} from '../Logger';import type {LimitConfigService} from '../limits/LimitConfigService';import {resolveLimitSafe} from '../limits/LimitConfigUtils';import {createLimitMatchContext} from '../limits/LimitMatchContextBuilder';import type {IUserRepository} from '../user/IUserRepository';import type {VoicePresenceHeartbeatStore} from '../voice/VoicePresenceHeartbeatStore';import type {VoiceTopology} from '../voice/VoiceTopology';import type {IGatewayService} from './IGatewayService';import type {ILiveKitService} from './ILiveKitService';import type {IVoiceRoomStore} from './IVoiceRoomStore';import {isDMRoom, parseParticipantIdentity, parseParticipantMetadataWithRaw, parseRoomName} from './VoiceRoomContext';
interface VoiceWebhookParticipantContext { readonly type: 'dm' | 'guild'; readonly channelId: ChannelID; readonly guildId?: GuildID;}
export class LiveKitWebhookService { private receivers: Map<string, WebhookReceiver>; private serverMap: Map< string, { regionId: string; serverId: string; } >;
constructor( private voiceRoomStore: IVoiceRoomStore, private gatewayService: IGatewayService, private userRepository: IUserRepository, private liveKitService: ILiveKitService, private voiceTopology: VoiceTopology, private limitConfigService: LimitConfigService, private voicePresenceHeartbeatStore: VoicePresenceHeartbeatStore | null = null, ) { this.receivers = new Map(); this.serverMap = new Map(); this.rebuildReceivers(); this.voiceTopology.registerSubscriber(() => this.rebuildReceivers()); }
async verifyAndParse( body: string, authHeader: string | undefined, ): Promise<{ event: WebhookEvent; apiKey: string; }> { if (!authHeader) { throw new Error('Missing authorization header'); } let lastError: Error | null = null; for (const [apiKey, receiver] of this.receivers.entries()) { try { const event = await receiver.receive(body, authHeader); return {event: event as WebhookEvent, apiKey}; } catch (error) { lastError = error as Error; } } throw lastError || new Error('No webhook receivers configured'); }
async handleWebhookRequest(params: {body: string; authHeader: string | undefined}): Promise<{ status: number; body: string | null; }> { const {body, authHeader} = params; Logger.debug( { bodySize: body.length, hasAuthHeader: Boolean(authHeader), }, 'Received LiveKit webhook request', ); try { const data = await this.verifyAndParse(body, authHeader); const eventName = data.event.event; Logger.debug( { apiKey: data.apiKey, event: eventName, roomName: data.event.room?.name ?? null, participantIdentity: data.event.participant?.identity ?? null, trackType: data.event.track?.type ?? null, }, 'Parsed LiveKit webhook event', ); if (data.event.numDropped != null && data.event.numDropped > 0) { Logger.warn( { numDropped: data.event.numDropped, roomName: data.event.room?.name ?? null, eventType: data.event.event, }, 'LiveKit webhook reports dropped events - reconciliation may be needed', ); } await this.processEvent(data); return {status: 200, body: null}; } catch (error) { Logger.debug({error}, 'Error processing LiveKit webhook'); return {status: 400, body: 'Invalid webhook'}; } }
private rebuildReceivers(): void { const newReceivers = new Map<string, WebhookReceiver>(); const newServerMap = new Map< string, { regionId: string; serverId: string; } >(); const regions = this.voiceTopology.getAllRegions(); for (const region of regions) { const servers = this.voiceTopology.getServersForRegion(region.id); for (const server of servers) { newReceivers.set(server.apiKey, new WebhookReceiver(server.apiKey, server.apiSecret)); newServerMap.set(server.apiKey, {regionId: region.id, serverId: server.serverId}); } } this.receivers = newReceivers; this.serverMap = newServerMap; Logger.debug( { regionCount: regions.length, serverCount: newReceivers.size, }, 'Rebuilt LiveKit webhook receivers', ); }
private async markVoicePresenceHeartbeatEnded(params: { channelId: ChannelID; userId: UserID; connectionId: string; }): Promise<void> { if (this.voicePresenceHeartbeatStore === null) { return; } try { await this.voicePresenceHeartbeatStore.markHeartbeatEnded(params); } catch (error) { Logger.warn( { error, channelId: params.channelId.toString(), userId: params.userId.toString(), connectionId: params.connectionId, }, 'Failed to mark v2 voice presence heartbeat ended', ); } }
private async markChannelVoicePresenceHeartbeatsEnded(channelId: ChannelID): Promise<void> { if (this.voicePresenceHeartbeatStore === null) { return; } try { await this.voicePresenceHeartbeatStore.markChannelHeartbeatsEnded(channelId); } catch (error) { Logger.warn( {error, channelId: channelId.toString()}, 'Failed to clear active v2 voice presence heartbeats for finished room', ); } }
async handleRoomFinished(event: WebhookEvent, apiKey: string): Promise<void> { if (event.event !== 'room_finished' || !event.room) { return; } const roomName = event.room.name; const context = parseRoomName(roomName); if (!context) { Logger.warn({roomName}, 'Unknown room name format'); return; } Logger.debug( { roomName, contextType: context.type, guildId: isDMRoom(context) ? undefined : context.guildId.toString(), channelId: context.channelId.toString(), }, 'Processing LiveKit room_finished event', ); const sourceServer = this.serverMap.get(apiKey); if (isDMRoom(context)) { const pinned = await this.voiceRoomStore.getPinnedRoomServer(undefined, context.channelId); if (pinned && sourceServer && pinned.serverId !== sourceServer.serverId) { Logger.debug( { channelId: context.channelId.toString(), finishedServer: sourceServer.serverId, pinnedServer: pinned.serverId, }, 'Ignoring room_finished from stale server — room has moved to a different server', ); return; } await this.markChannelVoicePresenceHeartbeatsEnded(context.channelId); await this.voiceRoomStore.deleteRoomServer(undefined, context.channelId); Logger.debug({channelId: context.channelId.toString()}, 'Cleared DM voice room server pinning'); } else { const pinned = await this.voiceRoomStore.getPinnedRoomServer(context.guildId, context.channelId); if (pinned && sourceServer && pinned.serverId !== sourceServer.serverId) { Logger.debug( { guildId: context.guildId.toString(), channelId: context.channelId.toString(), finishedServer: sourceServer.serverId, pinnedServer: pinned.serverId, }, 'Ignoring room_finished from stale server — room has moved to a different server', ); return; } await this.markChannelVoicePresenceHeartbeatsEnded(context.channelId); await this.voiceRoomStore.deleteRoomServer(context.guildId, context.channelId); Logger.debug( {guildId: context.guildId.toString(), channelId: context.channelId.toString()}, 'Cleared guild voice room server pinning', ); try { const result = await this.gatewayService.disconnectAllVoiceUsersInChannel({ guildId: context.guildId, channelId: context.channelId, }); Logger.info( { guildId: context.guildId.toString(), channelId: context.channelId.toString(), disconnectedCount: result.disconnectedCount, }, 'Cleaned up zombie voice connections for finished room', ); } catch (error) { Logger.error( {error, guildId: context.guildId.toString(), channelId: context.channelId.toString()}, 'Failed to clean up voice connections for finished room', ); } } }
async handleParticipantJoined(event: WebhookEvent): Promise<void> { if (event.event !== 'participant_joined') { return; } const {participant} = event; if (!participant?.metadata) { Logger.debug('Participant joined without metadata, skipping'); return; } const parsed = parseParticipantMetadataWithRaw(participant.metadata); if (!parsed) { Logger.warn({metadata: participant.metadata}, 'Failed to parse participant metadata'); return; } const {context, raw} = parsed; const tokenNonce = raw.token_nonce; Logger.debug( { type: context.type, participantIdentity: participant.identity, roomName: event.room?.name ?? null, channelId: context.channelId.toString(), guildId: context.type === 'guild' ? context.guildId.toString() : undefined, connectionId: context.connectionId, tokenNonce, }, 'Processing LiveKit participant_joined event', ); try { const guildId = context.type === 'guild' ? context.guildId : undefined; Logger.debug( { type: context.type, guildId: guildId?.toString(), channelId: context.channelId.toString(), connectionId: context.connectionId, participantIdentity: participant.identity, }, 'LiveKit participant_joined - confirming voice connection', ); const result = await this.gatewayService.confirmVoiceConnection({ guildId, channelId: context.channelId, connectionId: context.connectionId, tokenNonce, }); Logger.debug( { type: context.type, guildId: guildId?.toString(), channelId: context.channelId.toString(), connectionId: context.connectionId, success: result.success, error: result.error, }, 'LiveKit voice connection confirm result', ); if (!result.success) { if (result.error === 'connection_not_found') { Logger.warn( { type: context.type, guildId: guildId?.toString(), channelId: context.channelId.toString(), connectionId: context.connectionId, error: result.error, participantIdentity: participant.identity, }, 'LiveKit participant_joined did not match gateway state; leaving participant connected for reconciliation', ); return; } Logger.warn( { type: context.type, guildId: guildId?.toString(), channelId: context.channelId.toString(), connectionId: context.connectionId, error: result.error, participantIdentity: participant.identity, }, 'LiveKit participant_joined rejected - disconnecting participant', ); await this.markVoicePresenceHeartbeatEnded({ channelId: context.channelId, userId: context.userId, connectionId: context.connectionId, }); try { await this.liveKitService.disconnectParticipant({ guildId, channelId: context.channelId, userId: context.userId, connectionId: context.connectionId, regionId: raw.region_id ?? '', serverId: raw.server_id ?? '', }); } catch (disconnectError) { Logger.error({error: disconnectError}, 'Failed to disconnect rejected participant'); } return; } } catch (error) { Logger.error({error, type: context.type}, 'Error processing participant_joined'); } }
private async isParticipantStillInRoom(params: { participantIdentity: string; context: VoiceWebhookParticipantContext; regionId?: string; serverId?: string; }): Promise<'present' | 'absent' | 'unknown'> { const {participantIdentity, context} = params; const guildId = context.type === 'guild' ? context.guildId : undefined; let regionId = params.regionId; let serverId = params.serverId; if (!regionId || !serverId) { const pinnedServer = await this.voiceRoomStore.getPinnedRoomServer(guildId, context.channelId); if (pinnedServer) { regionId = pinnedServer.regionId; serverId = pinnedServer.serverId; } } if (!regionId || !serverId) { return 'unknown'; } const result = await this.liveKitService.listParticipants({ guildId, channelId: context.channelId, regionId, serverId, }); if (result.status === 'error') { Logger.warn( {errorCode: result.errorCode, retryable: result.retryable, participantIdentity}, 'Cannot determine participant presence due to LiveKit lookup failure', ); return 'unknown'; } return result.participants.some((p) => p.identity === participantIdentity) ? 'present' : 'absent'; }
async handleParticipantLeft(event: WebhookEvent): Promise<void> { if (event.event !== 'participant_left' && event.event !== 'participant_connection_aborted') { return; } const {participant} = event; if (!participant?.metadata) { Logger.debug('Participant left without metadata, skipping'); return; } const parsed = parseParticipantMetadataWithRaw(participant.metadata); if (!parsed) { Logger.warn({metadata: participant.metadata}, 'Failed to parse participant metadata'); return; } const {context, raw} = parsed; Logger.debug( { type: context.type, participantIdentity: participant.identity, roomName: event.room?.name ?? null, channelId: context.channelId.toString(), guildId: context.type === 'guild' ? context.guildId.toString() : undefined, userId: context.userId.toString(), connectionId: context.connectionId, }, `Processing LiveKit ${event.event} event`, ); try { if (raw.region_id && raw.server_id) { const guildId = context.type === 'guild' ? context.guildId : undefined; const pinnedServer = await this.voiceRoomStore.getPinnedRoomServer(guildId, context.channelId); if (pinnedServer && (pinnedServer.regionId !== raw.region_id || pinnedServer.serverId !== raw.server_id)) { Logger.debug( { type: context.type, participantIdentity: participant.identity, channelId: context.channelId.toString(), guildId: context.type === 'guild' ? context.guildId.toString() : undefined, connectionId: context.connectionId, eventRegionId: raw.region_id, eventServerId: raw.server_id, currentRegionId: pinnedServer.regionId, currentServerId: pinnedServer.serverId, }, 'Ignoring participant_left from stale server - room has migrated to a different server', ); return; } } const presenceStatus = await this.isParticipantStillInRoom({ participantIdentity: participant.identity, context, regionId: raw.region_id, serverId: raw.server_id, }); if (presenceStatus === 'present') { Logger.warn( { type: context.type, participantIdentity: participant.identity, channelId: context.channelId.toString(), guildId: context.type === 'guild' ? context.guildId.toString() : undefined, connectionId: context.connectionId, }, 'Ignoring stale participant_left event because participant is still present in room', ); return; } if (presenceStatus === 'unknown') { Logger.warn( { type: context.type, participantIdentity: participant.identity, channelId: context.channelId.toString(), guildId: context.type === 'guild' ? context.guildId.toString() : undefined, connectionId: context.connectionId, }, 'Skipping participant_left disconnect because participant presence is uncertain', ); return; } await this.markVoicePresenceHeartbeatEnded({ channelId: context.channelId, userId: context.userId, connectionId: context.connectionId, }); const guildId = context.type === 'guild' ? context.guildId : undefined; Logger.info( { type: context.type, guildId: guildId?.toString(), userId: context.userId.toString(), channelId: context.channelId.toString(), connectionId: context.connectionId, }, 'LiveKit participant_left - disconnecting voice user', ); const result = await this.gatewayService.disconnectVoiceUserIfInChannel({ guildId, channelId: context.channelId, userId: context.userId, connectionId: context.connectionId, }); Logger.debug( { type: context.type, guildId: guildId?.toString(), userId: context.userId.toString(), channelId: context.channelId.toString(), connectionId: context.connectionId, result, }, 'LiveKit participant_left voice disconnect result', ); } catch (error) { Logger.error({error, type: context.type}, 'Error processing participant_left'); } }
async handleTrackPublished(event: WebhookEvent, apiKey: string): Promise<void> { if (event.event !== 'track_published') { return; } const {room, participant, track} = event; if (!room || !participant || !track) { Logger.debug('Track published without required data, skipping'); return; } Logger.debug( { apiKey, roomName: room.name, participantIdentity: participant.identity, trackType: track.type, width: track.width, height: track.height, }, 'Processing LiveKit track_published event', ); if (track.type !== 1) { return; } if (track.source !== TrackSource.CAMERA && track.source !== TrackSource.SCREEN_SHARE) { return; } const trackSourceLabel = track.source === TrackSource.SCREEN_SHARE ? 'screen_share' : 'camera'; try { const identity = parseParticipantIdentity(participant.identity); if (!identity) { Logger.warn({identity: participant.identity}, 'Unexpected participant identity format'); return; } const {userId, connectionId} = identity; const user = await this.userRepository.findUnique(userId); if (!user) { Logger.warn({userId: userId.toString()}, 'User not found for track_published event'); return; } if (Config.instance.selfHosted) { return; } const ctx = createLimitMatchContext({user}); const hasHigherQuality = resolveLimitSafe( this.limitConfigService.getConfigSnapshot(), ctx, 'feature_higher_video_quality', 0, ); const canUseHigherQuality = hasHigherQuality > 0 && !user.isBot; if (canUseHigherQuality) { return; } const FREE_MAX_WIDTH = 1280; const FREE_MAX_HEIGHT = 720; const exceedsResolution = track.width > FREE_MAX_WIDTH || track.height > FREE_MAX_HEIGHT; if (!exceedsResolution) { return; } Logger.warn( { userId: userId.toString(), isBot: user.isBot, width: track.width, height: track.height, trackSource: trackSourceLabel, }, 'User without higher video quality entitlement attempted to publish video exceeding free tier limits - disconnecting', ); const roomContext = parseRoomName(room.name); if (!roomContext) { Logger.warn({roomName: room.name}, 'Unknown room name format, cannot disconnect'); return; } let regionId: string | undefined; let serverId: string | undefined; if (participant.metadata) { const parsed = parseParticipantMetadataWithRaw(participant.metadata); if (parsed) { regionId = parsed.raw.region_id; serverId = parsed.raw.server_id; } } if (!regionId || !serverId) { const serverInfo = this.serverMap.get(apiKey); if (serverInfo) { regionId = serverInfo.regionId; serverId = serverInfo.serverId; } } if (!regionId || !serverId) { const guildId = isDMRoom(roomContext) ? undefined : roomContext.guildId; const pinnedServer = await this.voiceRoomStore.getPinnedRoomServer(guildId, roomContext.channelId); if (pinnedServer) { regionId = pinnedServer.regionId; serverId = pinnedServer.serverId; } } if (!regionId || !serverId) { Logger.warn( {participantId: participant.identity, roomName: room.name, apiKey}, 'Missing region or server info, cannot disconnect', ); return; } const guildId = isDMRoom(roomContext) ? undefined : roomContext.guildId; Logger.debug( { userId: userId.toString(), type: roomContext.type, guildId: guildId?.toString(), channelId: roomContext.channelId.toString(), regionId, serverId, isBot: user.isBot, width: track.width, height: track.height, trackSource: trackSourceLabel, }, 'Disconnecting user without higher video quality entitlement for exceeding video quality limits', ); await this.liveKitService.disconnectParticipant({ userId, guildId, channelId: roomContext.channelId, connectionId, regionId, serverId, }); await this.gatewayService.disconnectVoiceUserIfInChannel({ guildId, channelId: roomContext.channelId, userId, connectionId, }); Logger.info( { userId: userId.toString(), type: roomContext.type, guildId: guildId?.toString(), channelId: roomContext.channelId.toString(), isBot: user.isBot, width: track.width, height: track.height, trackSource: trackSourceLabel, }, 'Disconnected user without higher video quality entitlement for exceeding video quality limits', ); } catch (error) { Logger.error({error}, 'Error processing track_published event'); } }
async processEvent(data: {event: WebhookEvent; apiKey: string}): Promise<void> { const {event, apiKey} = data; Logger.debug({event: event.event, apiKey}, 'Dispatching LiveKit webhook event'); switch (event.event) { case 'participant_joined': await this.handleParticipantJoined(event); break; case 'participant_left': case 'participant_connection_aborted': await this.handleParticipantLeft(event); break; case 'room_finished': await this.handleRoomFinished(event, apiKey); break; case 'track_published': await this.handleTrackPublished(event, apiKey); break; default: Logger.debug({event: event.event}, 'Ignoring LiveKit webhook event'); } Logger.debug({event: event.event, apiKey}, 'Finished LiveKit webhook event'); }}