refactor(voice): remove heartbeat and debug logging sessions (#2121)

This commit is contained in:
Hampus
2026-08-30 21:08:46 +02:00
committed by GitHub
parent 36b85512c6
commit b4d9cdc584
33 changed files with 566 additions and 3783 deletions
+5 -263
View File
@@ -8905,231 +8905,6 @@
}
}
},
"/admin/voice/diagnostics/objects": {
"get": {
"operationId": "list_voice_diagnostics_objects",
"summary": "List voice diagnostics objects",
"tags": ["Admin"],
"responses": {
"200": {
"description": "Success",
"content": {
"application/json": {"schema": {"$ref": "#/components/schemas/VoiceDiagnosticsObjectListResponse"}}
}
},
"400": {
"description": "Bad Request - The request was malformed or contained invalid data",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
},
"401": {
"description": "Unauthorized - Authentication is required or the token is invalid",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
},
"403": {
"description": "Forbidden - You do not have permission to perform this action",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
},
"429": {
"description": "Too Many Requests - You are being rate limited",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}},
"headers": {
"Retry-After": {
"description": "Number of seconds to wait before retrying (only on 429)",
"schema": {"type": "integer"}
},
"X-RateLimit-Limit": {
"description": "The number of requests that can be made in the current window",
"schema": {"type": "integer"}
},
"X-RateLimit-Remaining": {
"description": "The number of remaining requests that can be made",
"schema": {"type": "integer"}
},
"X-RateLimit-Reset": {
"description": "Unix timestamp when the rate limit resets",
"schema": {"type": "integer"}
}
}
},
"500": {
"description": "Internal Server Error - An unexpected error occurred",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
}
},
"description": "Lists raw voice diagnostics NDJSON S3 objects for a channel and time range. Requires VOICE_DIAGNOSTICS_VIEW permission.",
"security": [{"adminApiKey": []}],
"parameters": [
{
"name": "channel_id",
"in": "query",
"required": true,
"schema": {
"type": "string",
"pattern": "^(0|[1-9][0-9]*)$",
"description": "Channel id to query diagnostics for"
}
},
{
"name": "start_ms",
"in": "query",
"required": true,
"schema": {
"type": "integer",
"minimum": 0,
"maximum": 8640000000000000,
"format": "int53",
"description": "Inclusive start timestamp in milliseconds"
}
},
{
"name": "end_ms",
"in": "query",
"required": true,
"schema": {
"type": "integer",
"minimum": 0,
"maximum": 8640000000000000,
"format": "int53",
"description": "Inclusive end timestamp in milliseconds"
}
},
{
"name": "session_id",
"in": "query",
"required": false,
"schema": {
"type": "string",
"minLength": 1,
"maxLength": 128,
"pattern": "^[A-Za-z0-9_.:-]+$",
"description": "Optional voice diagnostics session id filter"
}
},
{
"name": "limit_objects",
"in": "query",
"required": false,
"schema": {
"type": "integer",
"minimum": 1,
"maximum": 5000,
"format": "int32",
"description": "Maximum matching S3 objects to include"
}
}
]
}
},
"/admin/voice/diagnostics/raw": {
"get": {
"operationId": "stream_voice_diagnostics_raw",
"summary": "Stream raw voice diagnostics",
"tags": ["Admin"],
"responses": {
"204": {"description": "No Content"},
"400": {
"description": "Bad Request - The request was malformed or contained invalid data",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
},
"401": {
"description": "Unauthorized - Authentication is required or the token is invalid",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
},
"403": {
"description": "Forbidden - You do not have permission to perform this action",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
},
"429": {
"description": "Too Many Requests - You are being rate limited",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}},
"headers": {
"Retry-After": {
"description": "Number of seconds to wait before retrying (only on 429)",
"schema": {"type": "integer"}
},
"X-RateLimit-Limit": {
"description": "The number of requests that can be made in the current window",
"schema": {"type": "integer"}
},
"X-RateLimit-Remaining": {
"description": "The number of remaining requests that can be made",
"schema": {"type": "integer"}
},
"X-RateLimit-Reset": {
"description": "Unix timestamp when the rate limit resets",
"schema": {"type": "integer"}
}
}
},
"500": {
"description": "Internal Server Error - An unexpected error occurred",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
}
},
"description": "Streams matching voice diagnostics NDJSON objects for local processing. Requires VOICE_DIAGNOSTICS_VIEW permission.",
"security": [{"adminApiKey": []}],
"parameters": [
{
"name": "channel_id",
"in": "query",
"required": true,
"schema": {
"type": "string",
"pattern": "^(0|[1-9][0-9]*)$",
"description": "Channel id to query diagnostics for"
}
},
{
"name": "start_ms",
"in": "query",
"required": true,
"schema": {
"type": "integer",
"minimum": 0,
"maximum": 8640000000000000,
"format": "int53",
"description": "Inclusive start timestamp in milliseconds"
}
},
{
"name": "end_ms",
"in": "query",
"required": true,
"schema": {
"type": "integer",
"minimum": 0,
"maximum": 8640000000000000,
"format": "int53",
"description": "Inclusive end timestamp in milliseconds"
}
},
{
"name": "session_id",
"in": "query",
"required": false,
"schema": {
"type": "string",
"minLength": 1,
"maxLength": 128,
"pattern": "^[A-Za-z0-9_.:-]+$",
"description": "Optional voice diagnostics session id filter"
}
},
{
"name": "limit_objects",
"in": "query",
"required": false,
"schema": {
"type": "integer",
"minimum": 1,
"maximum": 5000,
"format": "int32",
"description": "Maximum matching S3 objects to include"
}
}
]
}
},
"/admin/voice/regions/create": {
"post": {
"operationId": "create_voice_region",
@@ -9718,7 +9493,7 @@
"acls": {
"type": "array",
"items": {"type": "string"},
"maxItems": 115,
"maxItems": 114,
"description": "List of access control permissions for the key"
}
},
@@ -10078,7 +9853,7 @@
"acls": {
"type": "array",
"items": {"type": "string"},
"maxItems": 115,
"maxItems": 114,
"description": "List of access control permissions for the key"
}
},
@@ -10108,7 +9883,7 @@
"acls": {
"type": "array",
"items": {"type": "string"},
"maxItems": 115,
"maxItems": 114,
"description": "List of access control permissions for the key"
}
},
@@ -14918,7 +14693,7 @@
"pending_bulk_message_deletion_at": {"nullable": true, "type": "string"},
"deletion_reason_code": {"nullable": true, "allOf": [{"$ref": "#/components/schemas/Int32Type"}]},
"deletion_public_reason": {"nullable": true, "type": "string"},
"acls": {"type": "array", "items": {"type": "string"}, "maxItems": 115},
"acls": {"type": "array", "items": {"type": "string"}, "maxItems": 114},
"traits": {"type": "array", "items": {"type": "string"}, "maxItems": 100},
"has_totp": {"type": "boolean"},
"authenticator_types": {"type": "array", "items": {"$ref": "#/components/schemas/Int32Type"}, "maxItems": 10},
@@ -15486,7 +15261,7 @@
"acls": {
"type": "array",
"items": {"type": "string"},
"maxItems": 115,
"maxItems": 114,
"description": "List of access control permissions to assign"
}
},
@@ -15617,39 +15392,6 @@
"type": "object",
"properties": {"reason": {"type": "string"}, "notes": {"type": "string"}}
},
"VoiceDiagnosticsObjectListResponse": {
"type": "object",
"properties": {
"bucket": {"type": "string", "description": "Diagnostics S3 bucket"},
"objects": {
"type": "array",
"items": {
"type": "object",
"properties": {
"key": {"type": "string", "description": "S3 object key"},
"session_id": {"type": "string", "description": "Voice diagnostics session id"},
"start_ns": {
"type": "string",
"description": "First client event timestamp in the object, in nanoseconds"
},
"end_ns": {
"type": "string",
"description": "Last client event timestamp in the object, in nanoseconds"
},
"last_modified": {
"description": "S3 object last-modified timestamp",
"nullable": true,
"type": "string",
"format": "date-time"
}
},
"required": ["key", "session_id", "start_ns", "end_ns", "last_modified"]
},
"description": "Matching diagnostics objects"
}
},
"required": ["bucket", "objects"]
},
"CreateVoiceRegionResponse": {
"type": "object",
"properties": {
@@ -1,93 +0,0 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {Readable} from 'node:stream';
import {AdminACLs} from '@fluxer/constants/src/AdminACLs';
import {
VoiceDiagnosticsObjectListResponse,
VoiceDiagnosticsQueryRequest,
} from '@fluxer/schema/src/domains/admin/AdminVoiceSchemas';
import {VOICE_DIAGNOSTICS_BUCKET, VoiceDiagnosticsService} from '../../channel/services/VoiceDiagnosticsService';
import {requireAdminACL} from '../../middleware/AdminMiddleware';
import {RateLimitMiddleware} from '../../middleware/RateLimitMiddleware';
import {OpenAPI} from '../../middleware/ResponseTypeMiddleware';
import {RateLimitConfigs} from '../../RateLimitConfig';
import type {HonoApp} from '../../types/HonoEnv';
import {Validator} from '../../Validator';
export function VoiceDiagnosticsAdminController(app: HonoApp) {
app.get(
'/admin/voice/diagnostics/objects',
RateLimitMiddleware(RateLimitConfigs.ADMIN_LOOKUP),
requireAdminACL(AdminACLs.VOICE_DIAGNOSTICS_VIEW),
Validator('query', VoiceDiagnosticsQueryRequest),
OpenAPI({
operationId: 'list_voice_diagnostics_objects',
summary: 'List voice diagnostics objects',
responseSchema: VoiceDiagnosticsObjectListResponse,
statusCode: 200,
security: 'adminApiKey',
tags: 'Admin',
description:
'Lists raw voice diagnostics NDJSON S3 objects for a channel and time range. Requires VOICE_DIAGNOSTICS_VIEW permission.',
}),
async (ctx) => {
const query = ctx.req.valid('query');
const service = new VoiceDiagnosticsService(
ctx.get('cacheService'),
ctx.get('channelService'),
ctx.get('gatewayService'),
ctx.get('storageService'),
);
const objects = await service.listObjects({
channelId: query.channel_id,
startMs: query.start_ms,
endMs: query.end_ms,
sessionId: query.session_id,
limitObjects: query.limit_objects,
});
return ctx.json({bucket: VOICE_DIAGNOSTICS_BUCKET, objects});
},
);
app.get(
'/admin/voice/diagnostics/raw',
RateLimitMiddleware(RateLimitConfigs.ADMIN_LOOKUP),
requireAdminACL(AdminACLs.VOICE_DIAGNOSTICS_VIEW),
Validator('query', VoiceDiagnosticsQueryRequest),
OpenAPI({
operationId: 'stream_voice_diagnostics_raw',
summary: 'Stream raw voice diagnostics',
responseSchema: null,
statusCode: 200,
security: 'adminApiKey',
tags: 'Admin',
description:
'Streams matching voice diagnostics NDJSON objects for local processing. Requires VOICE_DIAGNOSTICS_VIEW permission.',
}),
async (ctx) => {
const query = ctx.req.valid('query');
const service = new VoiceDiagnosticsService(
ctx.get('cacheService'),
ctx.get('channelService'),
ctx.get('gatewayService'),
ctx.get('storageService'),
);
const objects = await service.listObjects({
channelId: query.channel_id,
startMs: query.start_ms,
endMs: query.end_ms,
sessionId: query.session_id,
limitObjects: query.limit_objects,
});
const stream = service.createRawObjectStream(objects);
return new Response(Readable.toWeb(stream) as ReadableStream, {
status: 200,
headers: {
'Content-Type': 'application/x-ndjson',
'Cache-Control': 'no-store, private',
'X-Fluxer-Diagnostics-Bucket': VOICE_DIAGNOSTICS_BUCKET,
'X-Fluxer-Diagnostics-Object-Count': objects.length.toString(),
},
});
},
);
}
@@ -23,7 +23,6 @@ import {SystemAdminController} from './SystemAdminController';
import {SystemDmAdminController} from './SystemDmAdminController';
import {UserAdminController} from './UserAdminController';
import {VoiceAdminController} from './VoiceAdminController';
import {VoiceDiagnosticsAdminController} from './VoiceDiagnosticsAdminController';
export function registerAdminControllers(app: HonoApp) {
AdminApiKeyAdminController(app);
@@ -42,7 +41,6 @@ export function registerAdminControllers(app: HonoApp) {
ReportAdminController(app);
BillingAdminController(app);
VoiceAdminController(app);
VoiceDiagnosticsAdminController(app);
GatewayAdminController(app);
SearchAdminController(app);
DiscoveryAdminController(app);
@@ -1,110 +0,0 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {UserFlags} from '@fluxer/constants/src/UserConstants';
import {MissingPermissionsError} from '@fluxer/errors/src/domains/core/MissingPermissionsError';
import {
VoiceDebugLoggingEventsBodySchema,
VoiceDebugLoggingToggleBodySchema,
} from '@fluxer/schema/src/domains/channel/ChannelRequestSchemas';
import {
VoiceDebugLoggingEventsResponse,
VoiceDebugLoggingStatusResponse,
} from '@fluxer/schema/src/domains/channel/ChannelSchemas';
import {ChannelIdParam} from '@fluxer/schema/src/domains/common/CommonParamSchemas';
import type {Context} from 'hono';
import {createChannelID} from '../../BrandedTypes';
import {DefaultUserOnly, LoginRequired} from '../../middleware/AuthMiddleware';
import {RateLimitMiddleware} from '../../middleware/RateLimitMiddleware';
import {OpenAPI} from '../../middleware/ResponseTypeMiddleware';
import {RateLimitConfigs} from '../../RateLimitConfig';
import type {HonoApp, HonoEnv} from '../../types/HonoEnv';
import {Validator} from '../../Validator';
import {VoiceDiagnosticsService} from '../services/VoiceDiagnosticsService';
function makeVoiceDiagnosticsService(ctx: Context<HonoEnv>): VoiceDiagnosticsService {
return new VoiceDiagnosticsService(
ctx.get('cacheService'),
ctx.get('channelService'),
ctx.get('gatewayService'),
ctx.get('storageService'),
);
}
export function VoiceDiagnosticsController(app: HonoApp) {
app.get(
'/channels/:channel_id/voice-debug-logging/session',
RateLimitMiddleware(RateLimitConfigs.CHANNEL_VOICE_DEBUG_LOGGING_STATUS),
LoginRequired,
DefaultUserOnly,
Validator('param', ChannelIdParam),
OpenAPI({
operationId: 'get_voice_debug_logging_status',
summary: 'Get voice debug logging status',
description:
'Returns whether staff-enabled voice debug logging is active for this channel. Clients poll this while connected to decide whether to upload diagnostics.',
responseSchema: VoiceDebugLoggingStatusResponse,
statusCode: 200,
security: ['bearerToken', 'sessionToken'],
tags: 'Channels',
}),
async (ctx) => {
const userId = ctx.get('user').id;
const channelId = createChannelID(ctx.req.valid('param').channel_id);
const service = makeVoiceDiagnosticsService(ctx);
return ctx.json(await service.getStatus({userId, channelId}));
},
);
app.put(
'/channels/:channel_id/voice-debug-logging/session',
RateLimitMiddleware(RateLimitConfigs.CHANNEL_VOICE_DEBUG_LOGGING_TOGGLE),
LoginRequired,
DefaultUserOnly,
Validator('param', ChannelIdParam),
Validator('json', VoiceDebugLoggingToggleBodySchema),
OpenAPI({
operationId: 'set_voice_debug_logging_status',
summary: 'Toggle voice debug logging',
description:
'Allows staff to start or stop a channel-scoped voice debug logging session. Non-staff users cannot activate or stop sessions.',
responseSchema: VoiceDebugLoggingStatusResponse,
statusCode: 200,
security: ['bearerToken', 'sessionToken'],
tags: 'Channels',
}),
async (ctx) => {
const user = ctx.get('user');
if ((user.flags & UserFlags.STAFF) === 0n) {
throw new MissingPermissionsError();
}
const channelId = createChannelID(ctx.req.valid('param').channel_id);
const {enabled, duration_ms} = ctx.req.valid('json');
const service = makeVoiceDiagnosticsService(ctx);
return ctx.json(await service.setSession({userId: user.id, channelId, enabled, durationMs: duration_ms}));
},
);
app.post(
'/channels/:channel_id/voice-debug-logging/events',
RateLimitMiddleware(RateLimitConfigs.CHANNEL_VOICE_DEBUG_LOGGING_EVENTS),
LoginRequired,
DefaultUserOnly,
Validator('param', ChannelIdParam),
Validator('json', VoiceDebugLoggingEventsBodySchema),
OpenAPI({
operationId: 'upload_voice_debug_logging_events',
summary: 'Upload voice debug logging events',
description:
'Uploads a small batch of client voice diagnostics events for an active staff-enabled debug logging session.',
responseSchema: VoiceDebugLoggingEventsResponse,
statusCode: 200,
security: ['bearerToken', 'sessionToken'],
tags: 'Channels',
}),
async (ctx) => {
const userId = ctx.get('user').id;
const channelId = createChannelID(ctx.req.valid('param').channel_id);
const body = ctx.req.valid('json');
const service = makeVoiceDiagnosticsService(ctx);
return ctx.json(await service.ingestEvents({userId, channelId, body}));
},
);
}
@@ -1,84 +0,0 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {VoicePresenceHeartbeatBodySchema} from '@fluxer/schema/src/domains/channel/ChannelRequestSchemas';
import {
VoicePresenceHeartbeatEndResponse,
VoicePresenceHeartbeatResponse,
} from '@fluxer/schema/src/domains/channel/ChannelSchemas';
import {ChannelIdParam} from '@fluxer/schema/src/domains/common/CommonParamSchemas';
import {createChannelID} from '../../BrandedTypes';
import {DefaultUserOnly, LoginRequired} from '../../middleware/AuthMiddleware';
import {RateLimitMiddleware} from '../../middleware/RateLimitMiddleware';
import {OpenAPI} from '../../middleware/ResponseTypeMiddleware';
import {RateLimitConfigs} from '../../RateLimitConfig';
import type {HonoApp} from '../../types/HonoEnv';
import {Validator} from '../../Validator';
import {VoicePresenceHeartbeatStore} from '../../voice/VoicePresenceHeartbeatStore';
export function VoicePresenceController(app: HonoApp) {
app.post(
'/channels/:channel_id/voice-presence/heartbeat',
RateLimitMiddleware(RateLimitConfigs.CHANNEL_VOICE_PRESENCE_HEARTBEAT),
LoginRequired,
DefaultUserOnly,
Validator('param', ChannelIdParam),
Validator('json', VoicePresenceHeartbeatBodySchema),
OpenAPI({
operationId: 'heartbeat_voice_presence',
summary: 'Heartbeat voice presence',
description:
'Refreshes the current user voice presence marker for v2 voice reconciliation. Clients call this while connected to voice.',
responseSchema: VoicePresenceHeartbeatResponse,
statusCode: 200,
security: ['bearerToken', 'sessionToken'],
tags: 'Channels',
}),
async (ctx) => {
const userId = ctx.get('user').id;
const channelId = createChannelID(ctx.req.valid('param').channel_id);
const body = ctx.req.valid('json');
const store = new VoicePresenceHeartbeatStore(ctx.get('apiContext').services.kv);
const heartbeat = await store.recordHeartbeat({
channelId,
userId,
connectionId: body.connection_id,
});
return ctx.json({
ok: true,
heartbeat_interval_ms: heartbeat.heartbeatIntervalMs,
heartbeat_ttl_ms: heartbeat.heartbeatTtlMs,
expires_at_ms: heartbeat.expiresAtMs,
});
},
);
app.delete(
'/channels/:channel_id/voice-presence/heartbeat',
RateLimitMiddleware(RateLimitConfigs.CHANNEL_VOICE_PRESENCE_HEARTBEAT),
LoginRequired,
DefaultUserOnly,
Validator('param', ChannelIdParam),
Validator('json', VoicePresenceHeartbeatBodySchema),
OpenAPI({
operationId: 'end_voice_presence_heartbeat',
summary: 'End voice presence heartbeat',
description:
'Clears the current user active v2 voice presence marker for a voice connection while preserving v2 enrollment for fast reconciliation.',
responseSchema: VoicePresenceHeartbeatEndResponse,
statusCode: 200,
security: ['bearerToken', 'sessionToken'],
tags: 'Channels',
}),
async (ctx) => {
const userId = ctx.get('user').id;
const channelId = createChannelID(ctx.req.valid('param').channel_id);
const body = ctx.req.valid('json');
const store = new VoicePresenceHeartbeatStore(ctx.get('apiContext').services.kv);
await store.markHeartbeatEnded({
channelId,
userId,
connectionId: body.connection_id,
});
return ctx.json({ok: true});
},
);
}
@@ -7,8 +7,6 @@ import {MessageController} from './MessageController';
import {MessageInteractionController} from './MessageInteractionController';
import {ScheduledMessageController} from './ScheduledMessageController';
import {StreamController} from './StreamController';
import {VoiceDiagnosticsController} from './VoiceDiagnosticsController';
import {VoicePresenceController} from './VoicePresenceController';
export function registerChannelControllers(app: HonoApp) {
ChannelController(app);
@@ -17,6 +15,4 @@ export function registerChannelControllers(app: HonoApp) {
ScheduledMessageController(app);
CallController(app);
StreamController(app);
VoiceDiagnosticsController(app);
VoicePresenceController(app);
}
@@ -1,336 +0,0 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {Readable} from 'node:stream';
import {Permissions} from '@fluxer/constants/src/ChannelConstants';
import {MissingPermissionsError} from '@fluxer/errors/src/domains/core/MissingPermissionsError';
import type {
VoiceDebugLoggingEventSchema,
VoiceDebugLoggingEventsBodySchema,
} from '@fluxer/schema/src/domains/channel/ChannelRequestSchemas';
import type {
VoiceDebugLoggingEventsResponse,
VoiceDebugLoggingStatusResponse,
} from '@fluxer/schema/src/domains/channel/ChannelSchemas';
import type {ICacheService} from '@pkgs/cache/src/ICacheService';
import type {ChannelID, UserID} from '../../BrandedTypes';
import type {IGatewayService} from '../../infrastructure/IGatewayService';
import type {IStorageService} from '../../infrastructure/IStorageService';
import type {ChannelService} from './ChannelService';
export const VOICE_DIAGNOSTICS_BUCKET = 'fluxer-voice-diagnostics';
const VOICE_DEBUG_LOGGING_POLL_INTERVAL_MS = 10000;
const VOICE_DEBUG_LOGGING_UPLOAD_INTERVAL_MS = 2000;
const DEFAULT_SESSION_DURATION_MS = 60 * 60 * 1000;
const MAX_SESSION_DURATION_MS = 4 * 60 * 60 * 1000;
const NDJSON_CONTENT_TYPE = 'application/x-ndjson';
const TEXT_ENCODER = new TextEncoder();
interface ActiveVoiceDebugLoggingSession {
session_id: string;
channel_id: string;
activated_by_user_id: string;
started_at_ms: number;
expires_at_ms: number;
}
interface VoiceDiagnosticsObject {
key: string;
session_id: string;
start_ns: string;
end_ns: string;
last_modified: string | null;
}
function sessionCacheKey(channelId: ChannelID): string {
return `voice_debug_logging:channel:${channelId.toString()}`;
}
function createId(prefix: string): string {
const randomUUID = globalThis.crypto?.randomUUID?.bind(globalThis.crypto);
if (randomUUID) return `${prefix}_${randomUUID()}`;
return `${prefix}_${Date.now().toString(36)}_${Math.random().toString(36).slice(2)}`;
}
function datePartFromMs(timestampMs: number): string {
return new Date(timestampMs).toISOString().slice(0, 10);
}
function sanitizeKeyPart(value: string): string {
return value.replace(/[^A-Za-z0-9_.:-]/g, '_').slice(0, 256) || 'unknown';
}
function minMaxEventTimestampNs(events: ReadonlyArray<VoiceDebugLoggingEventSchema>): {
startNs: bigint;
endNs: bigint;
} {
let startNs: bigint | null = null;
let endNs: bigint | null = null;
for (const event of events) {
const timestampNs = BigInt(event.timestamp_ns);
if (startNs === null || timestampNs < startNs) startNs = timestampNs;
if (endNs === null || timestampNs > endNs) endNs = timestampNs;
}
return {startNs: startNs ?? 0n, endNs: endNs ?? 0n};
}
function buildDiagnosticsKey(params: {
channelId: ChannelID;
sessionId: string;
userId: UserID;
connectionId?: string;
participantIdentity?: string;
startNs: bigint;
endNs: bigint;
}): string {
const timestampMs = Number(params.startNs / 1000000n);
const date = datePartFromMs(timestampMs > 0 ? timestampMs : Date.now());
return [
`channel_id=${params.channelId.toString()}`,
`date=${date}`,
`session_id=${sanitizeKeyPart(params.sessionId)}`,
`start_ns=${params.startNs.toString()}_end_ns=${params.endNs.toString()}`,
`participant=${sanitizeKeyPart(params.userId.toString())}`,
`connection=${sanitizeKeyPart(params.connectionId ?? 'unknown')}`,
`identity=${sanitizeKeyPart(params.participantIdentity ?? 'unknown')}`,
`batch=${createId('batch')}.ndjson`,
].join('/');
}
function parseDiagnosticsKey(key: string): VoiceDiagnosticsObject | null {
const match = key.match(/^channel_id=[^/]+\/date=[^/]+\/session_id=([^/]+)\/start_ns=([0-9]+)_end_ns=([0-9]+)\//);
if (!match) return null;
const [, sessionId, startNs, endNs] = match;
if (!sessionId || !startNs || !endNs) return null;
return {
key,
session_id: sessionId,
start_ns: startNs,
end_ns: endNs,
last_modified: null,
};
}
function buildDatePrefixes(params: {
channelId: string;
startMs: number;
endMs: number;
sessionId?: string;
}): Array<string> {
const prefixes: Array<string> = [];
const cursor = new Date(params.startMs);
cursor.setUTCHours(0, 0, 0, 0);
const end = new Date(params.endMs);
end.setUTCHours(0, 0, 0, 0);
while (cursor.getTime() <= end.getTime()) {
const base = `channel_id=${params.channelId}/date=${datePartFromMs(cursor.getTime())}/`;
prefixes.push(params.sessionId ? `${base}session_id=${sanitizeKeyPart(params.sessionId)}/` : base);
cursor.setUTCDate(cursor.getUTCDate() + 1);
}
return prefixes;
}
async function assertVoiceChannelAccess(params: {
channelService: ChannelService;
gatewayService: IGatewayService;
userId: UserID;
channelId: ChannelID;
}): Promise<void> {
const channel = await params.channelService.channelData.operations.getChannel({
userId: params.userId,
channelId: params.channelId,
});
if (!channel.guildId) return;
const hasConnect = await params.gatewayService.checkPermission({
guildId: channel.guildId,
channelId: params.channelId,
userId: params.userId,
permission: Permissions.CONNECT,
});
if (!hasConnect) {
throw new MissingPermissionsError();
}
}
export class VoiceDiagnosticsService {
constructor(
private readonly cacheService: ICacheService,
private readonly channelService: ChannelService,
private readonly gatewayService: IGatewayService,
private readonly storageService: IStorageService,
) {}
private async getActiveSession(channelId: ChannelID): Promise<ActiveVoiceDebugLoggingSession | null> {
const session = await this.cacheService.get<ActiveVoiceDebugLoggingSession>(sessionCacheKey(channelId));
if (!session) return null;
if (session.expires_at_ms <= Date.now()) {
await this.cacheService.delete(sessionCacheKey(channelId));
return null;
}
return session;
}
private statusFromSession(session: ActiveVoiceDebugLoggingSession | null): VoiceDebugLoggingStatusResponse {
return {
active: session !== null,
session_id: session?.session_id ?? null,
activated_by_user_id: session?.activated_by_user_id ?? null,
started_at_ms: session?.started_at_ms ?? null,
expires_at_ms: session?.expires_at_ms ?? null,
poll_interval_ms: VOICE_DEBUG_LOGGING_POLL_INTERVAL_MS,
upload_interval_ms: VOICE_DEBUG_LOGGING_UPLOAD_INTERVAL_MS,
};
}
async getStatus(params: {userId: UserID; channelId: ChannelID}): Promise<VoiceDebugLoggingStatusResponse> {
await assertVoiceChannelAccess({
channelService: this.channelService,
gatewayService: this.gatewayService,
userId: params.userId,
channelId: params.channelId,
});
return this.statusFromSession(await this.getActiveSession(params.channelId));
}
async setSession(params: {
userId: UserID;
channelId: ChannelID;
enabled: boolean;
durationMs?: number;
}): Promise<VoiceDebugLoggingStatusResponse> {
await assertVoiceChannelAccess({
channelService: this.channelService,
gatewayService: this.gatewayService,
userId: params.userId,
channelId: params.channelId,
});
const key = sessionCacheKey(params.channelId);
if (!params.enabled) {
await this.cacheService.delete(key);
return this.statusFromSession(null);
}
const now = Date.now();
const durationMs = Math.min(params.durationMs ?? DEFAULT_SESSION_DURATION_MS, MAX_SESSION_DURATION_MS);
const session: ActiveVoiceDebugLoggingSession = {
session_id: createId('voice_debug'),
channel_id: params.channelId.toString(),
activated_by_user_id: params.userId.toString(),
started_at_ms: now,
expires_at_ms: now + durationMs,
};
await this.cacheService.set(key, session, Math.ceil(durationMs / 1000));
return this.statusFromSession(session);
}
async ingestEvents(params: {
userId: UserID;
channelId: ChannelID;
body: VoiceDebugLoggingEventsBodySchema;
}): Promise<VoiceDebugLoggingEventsResponse> {
await assertVoiceChannelAccess({
channelService: this.channelService,
gatewayService: this.gatewayService,
userId: params.userId,
channelId: params.channelId,
});
const session = await this.getActiveSession(params.channelId);
if (!session || session.session_id !== params.body.session_id) {
return {accepted: false, active: session !== null, stored_event_count: 0};
}
const {startNs, endNs} = minMaxEventTimestampNs(params.body.events);
const key = buildDiagnosticsKey({
channelId: params.channelId,
sessionId: session.session_id,
userId: params.userId,
connectionId: params.body.connection_id,
participantIdentity: params.body.participant_identity,
startNs,
endNs,
});
const serverReceivedAtMs = Date.now();
const serverReceivedMonotonicNs = process.hrtime.bigint().toString();
const lines = params.body.events.map((event) =>
JSON.stringify({
schema_version: 1,
server_received_at_ms: serverReceivedAtMs,
server_received_monotonic_ns: serverReceivedMonotonicNs,
channel_id: params.channelId.toString(),
session_id: session.session_id,
activated_by_user_id: session.activated_by_user_id,
participant_user_id: params.userId.toString(),
connection_id: params.body.connection_id ?? null,
participant_identity: params.body.participant_identity ?? null,
event,
}),
);
await this.storageService.uploadObject({
bucket: VOICE_DIAGNOSTICS_BUCKET,
key,
body: TEXT_ENCODER.encode(`${lines.join('\n')}\n`),
contentType: NDJSON_CONTENT_TYPE,
});
return {accepted: true, active: true, stored_event_count: params.body.events.length};
}
async listObjects(params: {
channelId: string;
startMs: number;
endMs: number;
sessionId?: string;
limitObjects: number;
}): Promise<Array<VoiceDiagnosticsObject>> {
const startNs = BigInt(Math.floor(params.startMs)) * 1000000n;
const endNs = BigInt(Math.floor(params.endMs)) * 1000000n;
const prefixes = buildDatePrefixes({
channelId: params.channelId,
startMs: params.startMs,
endMs: params.endMs,
sessionId: params.sessionId,
});
const listed = await Promise.all(
prefixes.map((prefix) => this.storageService.listObjects({bucket: VOICE_DIAGNOSTICS_BUCKET, prefix})),
);
return listed
.flat()
.flatMap((object): Array<VoiceDiagnosticsObject> => {
const parsed = parseDiagnosticsKey(object.key);
if (!parsed) return [];
return [
{
...parsed,
last_modified: object.lastModified?.toISOString() ?? null,
},
];
})
.filter((object) => {
const objectStartNs = BigInt(object.start_ns);
const objectEndNs = BigInt(object.end_ns);
return objectEndNs >= startNs && objectStartNs <= endNs;
})
.sort((a, b) => {
const delta = BigInt(a.start_ns) - BigInt(b.start_ns);
if (delta < 0n) return -1;
if (delta > 0n) return 1;
return a.key.localeCompare(b.key);
})
.slice(0, params.limitObjects);
}
createRawObjectStream(objects: ReadonlyArray<VoiceDiagnosticsObject>): Readable {
const storageService = this.storageService;
async function* streamObjects(): AsyncGenerator<Uint8Array> {
for (const object of objects) {
const stream = await storageService.streamObject({
bucket: VOICE_DIAGNOSTICS_BUCKET,
key: object.key,
});
if (!stream) continue;
for await (const chunk of stream.body) {
yield typeof chunk === 'string' ? TEXT_ENCODER.encode(chunk) : new Uint8Array(chunk);
}
yield TEXT_ENCODER.encode('\n');
}
}
return Readable.from(streamObjects());
}
}
@@ -2,14 +2,13 @@
import type {WebhookEvent} from 'livekit-server-sdk';
import {TrackSource, WebhookReceiver} from 'livekit-server-sdk';
import type {ChannelID, GuildID, UserID} from '../BrandedTypes';
import type {ChannelID, GuildID} 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';
@@ -39,7 +38,6 @@ export class LiveKitWebhookService {
private liveKitService: ILiveKitService,
private voiceTopology: VoiceTopology,
private limitConfigService: LimitConfigService,
private voicePresenceHeartbeatStore: VoicePresenceHeartbeatStore | null = null,
) {
this.receivers = new Map();
this.serverMap = new Map();
@@ -140,43 +138,6 @@ export class LiveKitWebhookService {
);
}
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;
@@ -210,7 +171,6 @@ export class LiveKitWebhookService {
);
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 {
@@ -227,7 +187,6 @@ export class LiveKitWebhookService {
);
return;
}
await this.markChannelVoicePresenceHeartbeatsEnded(context.channelId);
await this.voiceRoomStore.deleteRoomServer(context.guildId, context.channelId);
Logger.debug(
{guildId: context.guildId.toString(), channelId: context.channelId.toString()},
@@ -338,11 +297,6 @@ export class LiveKitWebhookService {
},
'LiveKit participant_joined rejected - disconnecting participant',
);
await this.markVoicePresenceHeartbeatEnded({
channelId: context.channelId,
userId: context.userId,
connectionId: context.connectionId,
});
try {
await this.liveKitService.disconnectParticipant({
guildId,
@@ -479,11 +433,6 @@ export class LiveKitWebhookService {
);
return;
}
await this.markVoicePresenceHeartbeatEnded({
channelId: context.channelId,
userId: context.userId,
connectionId: context.connectionId,
});
const guildId = context.type === 'guild' ? context.guildId : undefined;
Logger.info(
{
@@ -81,7 +81,6 @@ import {UserContentRequestService} from '../user/services/UserContentRequestServ
import {UserRelationshipRequestService} from '../user/services/UserRelationshipRequestService';
import {UserService} from '../user/services/UserService';
import {resolveRequestClientIp} from '../utils/IpUtils';
import {VoicePresenceHeartbeatStore} from '../voice/VoicePresenceHeartbeatStore';
import {VoiceService} from '../voice/VoiceService';
import {WebhookRequestService} from '../webhook/WebhookRequestService';
import {WebhookService} from '../webhook/WebhookService';
@@ -334,7 +333,6 @@ function getLiveKitWebhookService(): LiveKitWebhookService | null {
liveKitService,
voiceTopology,
getLimitConfigService(),
new VoicePresenceHeartbeatStore(getKVClient()),
);
}
}
-646
View File
@@ -5468,410 +5468,6 @@
]
}
},
"/channels/{channel_id}/voice-debug-logging/events": {
"post": {
"operationId": "upload_voice_debug_logging_events",
"summary": "Upload voice debug logging events",
"tags": ["Channels"],
"responses": {
"200": {
"description": "Success",
"content": {
"application/json": {"schema": {"$ref": "#/components/schemas/VoiceDebugLoggingEventsResponse"}}
}
},
"400": {
"description": "Bad Request - The request was malformed or contained invalid data",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
},
"401": {
"description": "Unauthorized - Authentication is required or the token is invalid",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
},
"403": {
"description": "Forbidden - You do not have permission to perform this action",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
},
"429": {
"description": "Too Many Requests - You are being rate limited",
"content": {
"application/json": {
"schema": {
"type": "object",
"properties": {
"code": {"type": "string", "enum": ["RATE_LIMITED"]},
"message": {"type": "string"},
"retry_after": {"type": "number", "description": "Seconds to wait before retrying"},
"global": {"type": "boolean", "description": "Whether this is a global rate limit"}
},
"required": ["code", "message", "retry_after"]
}
}
},
"headers": {
"Retry-After": {
"description": "Number of seconds to wait before retrying (only on 429)",
"schema": {"type": "integer"}
},
"X-RateLimit-Limit": {
"description": "The number of requests that can be made in the current window",
"schema": {"type": "integer"}
},
"X-RateLimit-Remaining": {
"description": "The number of remaining requests that can be made",
"schema": {"type": "integer"}
},
"X-RateLimit-Reset": {
"description": "Unix timestamp when the rate limit resets",
"schema": {"type": "integer"}
}
}
},
"500": {
"description": "Internal Server Error - An unexpected error occurred",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
}
},
"x-mint": {"metadata": {"title": "Upload voice debug logging events"}},
"description": "Uploads a small batch of client voice diagnostics events for an active staff-enabled debug logging session.",
"security": [{"sessionToken": []}],
"parameters": [
{
"name": "channel_id",
"in": "path",
"required": true,
"schema": {"$ref": "#/components/schemas/SnowflakeType"},
"description": "The ID of the channel"
}
],
"requestBody": {
"required": true,
"content": {
"application/json": {"schema": {"$ref": "#/components/schemas/VoiceDebugLoggingEventsBodySchema"}}
}
}
}
},
"/channels/{channel_id}/voice-debug-logging/session": {
"get": {
"operationId": "get_voice_debug_logging_status",
"summary": "Get voice debug logging status",
"tags": ["Channels"],
"responses": {
"200": {
"description": "Success",
"content": {
"application/json": {"schema": {"$ref": "#/components/schemas/VoiceDebugLoggingStatusResponse"}}
}
},
"400": {
"description": "Bad Request - The request was malformed or contained invalid data",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
},
"401": {
"description": "Unauthorized - Authentication is required or the token is invalid",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
},
"403": {
"description": "Forbidden - You do not have permission to perform this action",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
},
"429": {
"description": "Too Many Requests - You are being rate limited",
"content": {
"application/json": {
"schema": {
"type": "object",
"properties": {
"code": {"type": "string", "enum": ["RATE_LIMITED"]},
"message": {"type": "string"},
"retry_after": {"type": "number", "description": "Seconds to wait before retrying"},
"global": {"type": "boolean", "description": "Whether this is a global rate limit"}
},
"required": ["code", "message", "retry_after"]
}
}
},
"headers": {
"Retry-After": {
"description": "Number of seconds to wait before retrying (only on 429)",
"schema": {"type": "integer"}
},
"X-RateLimit-Limit": {
"description": "The number of requests that can be made in the current window",
"schema": {"type": "integer"}
},
"X-RateLimit-Remaining": {
"description": "The number of remaining requests that can be made",
"schema": {"type": "integer"}
},
"X-RateLimit-Reset": {
"description": "Unix timestamp when the rate limit resets",
"schema": {"type": "integer"}
}
}
},
"500": {
"description": "Internal Server Error - An unexpected error occurred",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
}
},
"x-mint": {"metadata": {"title": "Get voice debug logging status"}},
"description": "Returns whether staff-enabled voice debug logging is active for this channel. Clients poll this while connected to decide whether to upload diagnostics.",
"security": [{"sessionToken": []}],
"parameters": [
{
"name": "channel_id",
"in": "path",
"required": true,
"schema": {"$ref": "#/components/schemas/SnowflakeType"},
"description": "The ID of the channel"
}
]
},
"put": {
"operationId": "set_voice_debug_logging_status",
"summary": "Toggle voice debug logging",
"tags": ["Channels"],
"responses": {
"200": {
"description": "Success",
"content": {
"application/json": {"schema": {"$ref": "#/components/schemas/VoiceDebugLoggingStatusResponse"}}
}
},
"400": {
"description": "Bad Request - The request was malformed or contained invalid data",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
},
"401": {
"description": "Unauthorized - Authentication is required or the token is invalid",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
},
"403": {
"description": "Forbidden - You do not have permission to perform this action",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
},
"429": {
"description": "Too Many Requests - You are being rate limited",
"content": {
"application/json": {
"schema": {
"type": "object",
"properties": {
"code": {"type": "string", "enum": ["RATE_LIMITED"]},
"message": {"type": "string"},
"retry_after": {"type": "number", "description": "Seconds to wait before retrying"},
"global": {"type": "boolean", "description": "Whether this is a global rate limit"}
},
"required": ["code", "message", "retry_after"]
}
}
},
"headers": {
"Retry-After": {
"description": "Number of seconds to wait before retrying (only on 429)",
"schema": {"type": "integer"}
},
"X-RateLimit-Limit": {
"description": "The number of requests that can be made in the current window",
"schema": {"type": "integer"}
},
"X-RateLimit-Remaining": {
"description": "The number of remaining requests that can be made",
"schema": {"type": "integer"}
},
"X-RateLimit-Reset": {
"description": "Unix timestamp when the rate limit resets",
"schema": {"type": "integer"}
}
}
},
"500": {
"description": "Internal Server Error - An unexpected error occurred",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
}
},
"x-mint": {"metadata": {"title": "Toggle voice debug logging"}},
"description": "Allows staff to start or stop a channel-scoped voice debug logging session. Non-staff users cannot activate or stop sessions.",
"security": [{"sessionToken": []}],
"parameters": [
{
"name": "channel_id",
"in": "path",
"required": true,
"schema": {"$ref": "#/components/schemas/SnowflakeType"},
"description": "The ID of the channel"
}
],
"requestBody": {
"required": true,
"content": {
"application/json": {"schema": {"$ref": "#/components/schemas/VoiceDebugLoggingToggleBodySchema"}}
}
}
}
},
"/channels/{channel_id}/voice-presence/heartbeat": {
"post": {
"operationId": "heartbeat_voice_presence",
"summary": "Heartbeat voice presence",
"tags": ["Channels"],
"responses": {
"200": {
"description": "Success",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/VoicePresenceHeartbeatResponse"}}}
},
"400": {
"description": "Bad Request - The request was malformed or contained invalid data",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
},
"401": {
"description": "Unauthorized - Authentication is required or the token is invalid",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
},
"403": {
"description": "Forbidden - You do not have permission to perform this action",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
},
"429": {
"description": "Too Many Requests - You are being rate limited",
"content": {
"application/json": {
"schema": {
"type": "object",
"properties": {
"code": {"type": "string", "enum": ["RATE_LIMITED"]},
"message": {"type": "string"},
"retry_after": {"type": "number", "description": "Seconds to wait before retrying"},
"global": {"type": "boolean", "description": "Whether this is a global rate limit"}
},
"required": ["code", "message", "retry_after"]
}
}
},
"headers": {
"Retry-After": {
"description": "Number of seconds to wait before retrying (only on 429)",
"schema": {"type": "integer"}
},
"X-RateLimit-Limit": {
"description": "The number of requests that can be made in the current window",
"schema": {"type": "integer"}
},
"X-RateLimit-Remaining": {
"description": "The number of remaining requests that can be made",
"schema": {"type": "integer"}
},
"X-RateLimit-Reset": {
"description": "Unix timestamp when the rate limit resets",
"schema": {"type": "integer"}
}
}
},
"500": {
"description": "Internal Server Error - An unexpected error occurred",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
}
},
"x-mint": {"metadata": {"title": "Heartbeat voice presence"}},
"description": "Refreshes the current user voice presence marker for v2 voice reconciliation. Clients call this while connected to voice.",
"security": [{"sessionToken": []}],
"parameters": [
{
"name": "channel_id",
"in": "path",
"required": true,
"schema": {"$ref": "#/components/schemas/SnowflakeType"},
"description": "The ID of the channel"
}
],
"requestBody": {
"required": true,
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/VoicePresenceHeartbeatBodySchema"}}}
}
},
"delete": {
"operationId": "end_voice_presence_heartbeat",
"summary": "End voice presence heartbeat",
"tags": ["Channels"],
"responses": {
"200": {
"description": "Success",
"content": {
"application/json": {"schema": {"$ref": "#/components/schemas/VoicePresenceHeartbeatEndResponse"}}
}
},
"400": {
"description": "Bad Request - The request was malformed or contained invalid data",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
},
"401": {
"description": "Unauthorized - Authentication is required or the token is invalid",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
},
"403": {
"description": "Forbidden - You do not have permission to perform this action",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
},
"429": {
"description": "Too Many Requests - You are being rate limited",
"content": {
"application/json": {
"schema": {
"type": "object",
"properties": {
"code": {"type": "string", "enum": ["RATE_LIMITED"]},
"message": {"type": "string"},
"retry_after": {"type": "number", "description": "Seconds to wait before retrying"},
"global": {"type": "boolean", "description": "Whether this is a global rate limit"}
},
"required": ["code", "message", "retry_after"]
}
}
},
"headers": {
"Retry-After": {
"description": "Number of seconds to wait before retrying (only on 429)",
"schema": {"type": "integer"}
},
"X-RateLimit-Limit": {
"description": "The number of requests that can be made in the current window",
"schema": {"type": "integer"}
},
"X-RateLimit-Remaining": {
"description": "The number of remaining requests that can be made",
"schema": {"type": "integer"}
},
"X-RateLimit-Reset": {
"description": "Unix timestamp when the rate limit resets",
"schema": {"type": "integer"}
}
}
},
"500": {
"description": "Internal Server Error - An unexpected error occurred",
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}
}
},
"x-mint": {"metadata": {"title": "End voice presence heartbeat"}},
"description": "Clears the current user active v2 voice presence marker for a voice connection while preserving v2 enrollment for fast reconciliation.",
"security": [{"sessionToken": []}],
"parameters": [
{
"name": "channel_id",
"in": "path",
"required": true,
"schema": {"$ref": "#/components/schemas/SnowflakeType"},
"description": "The ID of the channel"
}
],
"requestBody": {
"required": true,
"content": {"application/json": {"schema": {"$ref": "#/components/schemas/VoicePresenceHeartbeatBodySchema"}}}
}
}
},
"/channels/{channel_id}/webhooks": {
"get": {
"operationId": "list_channel_webhooks",
@@ -29516,248 +29112,6 @@
},
"required": ["rate_limit_per_user", "retry_after_ms", "next_send_allowed_at", "can_bypass"]
},
"VoiceDebugLoggingEventsResponse": {
"type": "object",
"properties": {
"accepted": {"type": "boolean", "description": "Whether the telemetry batch was accepted for storage"},
"active": {
"type": "boolean",
"description": "Whether the server still considers this logging session active"
},
"stored_event_count": {
"type": "integer",
"minimum": 0,
"maximum": 2147483647,
"format": "int32",
"description": "Number of events written to diagnostics storage"
}
},
"required": ["accepted", "active", "stored_event_count"]
},
"VoiceDebugLoggingEventsBodySchema": {
"type": "object",
"properties": {
"session_id": {"type": "string", "description": "Active voice debug logging session id"},
"connection_id": {"type": "string", "description": "Client voice connection id"},
"participant_identity": {"type": "string", "description": "LiveKit participant identity"},
"events": {
"type": "array",
"items": {"$ref": "#/components/schemas/VoiceDebugLoggingEventSchema"},
"minItems": 1,
"maxItems": 200,
"description": "NDJSON batch events to store"
}
},
"required": ["session_id", "events"]
},
"VoiceDebugLoggingEventSchema": {
"type": "object",
"properties": {
"type": {"type": "string", "description": "Client-side diagnostic event type"},
"timestamp_ns": {
"type": "string",
"pattern": "^[0-9]{1,32}$",
"description": "Client wall-clock Unix timestamp in nanoseconds"
},
"monotonic_ns": {
"type": "string",
"pattern": "^[0-9]{1,32}$",
"description": "Client monotonic timestamp in nanoseconds"
},
"data": {
"type": "object",
"additionalProperties": {
"anyOf": [
{"type": "string"},
{"type": "number"},
{"type": "boolean"},
{
"type": "object",
"additionalProperties": {
"anyOf": [
{"type": "string"},
{"type": "number"},
{"type": "boolean"},
{
"type": "object",
"additionalProperties": {
"anyOf": [
{"type": "string"},
{"type": "number"},
{"type": "boolean"},
{"type": "object", "additionalProperties": true},
{"type": "null"}
]
}
},
{
"type": "array",
"items": {
"anyOf": [
{"type": "string"},
{"type": "number"},
{"type": "boolean"},
{"type": "object", "additionalProperties": true},
{"type": "null"}
]
}
},
{"type": "null"}
]
}
},
{
"type": "array",
"items": {
"anyOf": [
{"type": "string"},
{"type": "number"},
{"type": "boolean"},
{
"type": "object",
"additionalProperties": {
"anyOf": [
{"type": "string"},
{"type": "number"},
{"type": "boolean"},
{"type": "object", "additionalProperties": true},
{"type": "null"}
]
}
},
{
"type": "array",
"items": {
"anyOf": [
{"type": "string"},
{"type": "number"},
{"type": "boolean"},
{"type": "object", "additionalProperties": true},
{"type": "null"}
]
}
},
{"type": "null"}
]
}
},
{"type": "null"}
]
},
"description": "Event-specific diagnostic payload"
}
},
"required": ["type", "timestamp_ns"]
},
"VoiceDebugLoggingStatusResponse": {
"type": "object",
"properties": {
"active": {
"type": "boolean",
"description": "Whether clients in this channel should currently send voice diagnostics"
},
"session_id": {
"anyOf": [{"type": "string"}, {"type": "null"}],
"description": "Current debug logging session id, if active"
},
"activated_by_user_id": {
"anyOf": [{"$ref": "#/components/schemas/SnowflakeType"}, {"type": "null"}],
"description": "Staff user that activated the session, if active"
},
"started_at_ms": {
"anyOf": [
{"type": "integer", "minimum": 0, "maximum": 9007199254740991, "format": "int53"},
{"type": "null"}
],
"description": "Session start Unix timestamp in milliseconds"
},
"expires_at_ms": {
"anyOf": [
{"type": "integer", "minimum": 0, "maximum": 9007199254740991, "format": "int53"},
{"type": "null"}
],
"description": "Session expiration Unix timestamp in milliseconds"
},
"poll_interval_ms": {
"type": "integer",
"minimum": 0,
"maximum": 2147483647,
"format": "int32",
"description": "Recommended client polling interval in milliseconds"
},
"upload_interval_ms": {
"type": "integer",
"minimum": 0,
"maximum": 2147483647,
"format": "int32",
"description": "Recommended client telemetry batch upload interval in milliseconds"
}
},
"required": [
"active",
"session_id",
"activated_by_user_id",
"started_at_ms",
"expires_at_ms",
"poll_interval_ms",
"upload_interval_ms"
]
},
"VoiceDebugLoggingToggleBodySchema": {
"type": "object",
"properties": {
"enabled": {
"type": "boolean",
"description": "Whether voice debug logging should be active for this channel"
},
"duration_ms": {
"type": "integer",
"minimum": 60000,
"maximum": 14400000,
"format": "int32",
"description": "Optional activation duration in milliseconds. Defaults to one hour and is capped at four hours."
}
},
"required": ["enabled"]
},
"VoicePresenceHeartbeatResponse": {
"type": "object",
"properties": {
"ok": {"type": "boolean", "description": "Whether the heartbeat was accepted"},
"heartbeat_interval_ms": {
"type": "integer",
"minimum": 0,
"maximum": 2147483647,
"format": "int32",
"description": "Recommended client heartbeat interval in milliseconds"
},
"heartbeat_ttl_ms": {
"type": "integer",
"minimum": 0,
"maximum": 2147483647,
"format": "int32",
"description": "Server-side heartbeat expiration window in milliseconds"
},
"expires_at_ms": {
"type": "integer",
"minimum": 0,
"maximum": 9007199254740991,
"format": "int53",
"description": "Unix timestamp in milliseconds when this heartbeat expires"
}
},
"required": ["ok", "heartbeat_interval_ms", "heartbeat_ttl_ms", "expires_at_ms"]
},
"VoicePresenceHeartbeatBodySchema": {
"type": "object",
"properties": {"connection_id": {"type": "string", "description": "Client voice connection id"}},
"required": ["connection_id"]
},
"VoicePresenceHeartbeatEndResponse": {
"type": "object",
"properties": {"ok": {"type": "boolean", "description": "Whether the heartbeat was ended"}},
"required": ["ok"]
},
"WebhookResponse": {
"type": "object",
"properties": {
@@ -96,22 +96,6 @@ export const ChannelRateLimitConfigs = {
bucket: 'channel:call:stop_ringing::channel_id',
config: {limit: 20, windowMs: ms('10 seconds')},
} as RouteRateLimitConfig,
CHANNEL_VOICE_DEBUG_LOGGING_STATUS: {
bucket: 'channel:voice_debug_logging:status::channel_id',
config: {limit: 60, windowMs: ms('10 seconds')},
} as RouteRateLimitConfig,
CHANNEL_VOICE_DEBUG_LOGGING_TOGGLE: {
bucket: 'channel:voice_debug_logging:toggle::channel_id',
config: {limit: 10, windowMs: ms('1 minute')},
} as RouteRateLimitConfig,
CHANNEL_VOICE_DEBUG_LOGGING_EVENTS: {
bucket: 'channel:voice_debug_logging:events::channel_id::user_id',
config: {limit: 60, windowMs: ms('10 seconds')},
} as RouteRateLimitConfig,
CHANNEL_VOICE_PRESENCE_HEARTBEAT: {
bucket: 'channel:voice_presence:heartbeat::channel_id::user_id',
config: {limit: 20, windowMs: ms('10 seconds')},
} as RouteRateLimitConfig,
CHANNEL_STREAM_UPDATE: {
bucket: 'channel:stream:update::stream_key',
config: {limit: 20, windowMs: ms('10 seconds')},
@@ -1,109 +0,0 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {createHash} from 'node:crypto';
import type {IKVProvider} from '@pkgs/kv_client/src/IKVProvider';
import type {ChannelID, UserID} from '../BrandedTypes';
export type VoicePresenceHeartbeatState = 'active' | 'expired' | 'legacy';
interface VoicePresenceHeartbeatParams {
channelId: ChannelID;
userId: UserID;
connectionId: string;
}
interface VoicePresenceHeartbeatResult {
heartbeatIntervalMs: number;
heartbeatTtlMs: number;
expiresAtMs: number;
}
const VOICE_PRESENCE_HEARTBEAT_INTERVAL_MS = 15000;
const VOICE_PRESENCE_HEARTBEAT_TTL_SECONDS = 45;
const VOICE_PRESENCE_ENROLLMENT_TTL_SECONDS = 3600;
const VOICE_PRESENCE_HEARTBEAT_ACTIVE_PREFIX = 'voice:presence:v2:active:';
const VOICE_PRESENCE_HEARTBEAT_ENROLLED_PREFIX = 'voice:presence:v2:enrolled:';
const VOICE_PRESENCE_HEARTBEAT_CHANNEL_PREFIX = 'voice:presence:v2:channel:';
export class VoicePresenceHeartbeatStore {
constructor(private readonly kvClient: IKVProvider) {}
async recordHeartbeat(params: VoicePresenceHeartbeatParams): Promise<VoicePresenceHeartbeatResult> {
const now = Date.now();
const payload = JSON.stringify({
version: 2,
channelId: params.channelId.toString(),
userId: params.userId.toString(),
connectionId: params.connectionId,
heartbeatAtMs: now,
});
await Promise.all([
this.kvClient.setex(voicePresenceHeartbeatActiveKey(params), VOICE_PRESENCE_HEARTBEAT_TTL_SECONDS, payload),
this.kvClient.setex(voicePresenceHeartbeatEnrollmentKey(params), VOICE_PRESENCE_ENROLLMENT_TTL_SECONDS, payload),
this.kvClient.sadd(voicePresenceHeartbeatChannelKey(params.channelId), voicePresenceHeartbeatKeySuffix(params)),
this.kvClient.expire(voicePresenceHeartbeatChannelKey(params.channelId), VOICE_PRESENCE_ENROLLMENT_TTL_SECONDS),
]);
return {
heartbeatIntervalMs: VOICE_PRESENCE_HEARTBEAT_INTERVAL_MS,
heartbeatTtlMs: VOICE_PRESENCE_HEARTBEAT_TTL_SECONDS * 1000,
expiresAtMs: now + VOICE_PRESENCE_HEARTBEAT_TTL_SECONDS * 1000,
};
}
async getHeartbeatState(params: VoicePresenceHeartbeatParams): Promise<VoicePresenceHeartbeatState> {
const active = await this.kvClient.get(voicePresenceHeartbeatActiveKey(params));
if (active !== null) {
return 'active';
}
const enrolled = await this.kvClient.get(voicePresenceHeartbeatEnrollmentKey(params));
return enrolled !== null ? 'expired' : 'legacy';
}
async markHeartbeatEnded(params: VoicePresenceHeartbeatParams): Promise<void> {
const enrollmentKey = voicePresenceHeartbeatEnrollmentKey(params);
const enrolled = await this.kvClient.get(enrollmentKey);
const payload = JSON.stringify({
version: 2,
channelId: params.channelId.toString(),
userId: params.userId.toString(),
connectionId: params.connectionId,
endedAtMs: Date.now(),
});
await Promise.all([
this.kvClient.del(voicePresenceHeartbeatActiveKey(params)),
this.kvClient.srem(voicePresenceHeartbeatChannelKey(params.channelId), voicePresenceHeartbeatKeySuffix(params)),
enrolled === null
? Promise.resolve()
: this.kvClient.setex(enrollmentKey, VOICE_PRESENCE_ENROLLMENT_TTL_SECONDS, payload),
]);
}
async markChannelHeartbeatsEnded(channelId: ChannelID): Promise<void> {
const channelKey = voicePresenceHeartbeatChannelKey(channelId);
const suffixes = await this.kvClient.smembers(channelKey);
const activeKeys = suffixes.map((suffix) => `${VOICE_PRESENCE_HEARTBEAT_ACTIVE_PREFIX}${suffix}`);
if (activeKeys.length > 0) {
await this.kvClient.del(...activeKeys, channelKey);
return;
}
await this.kvClient.del(channelKey);
}
}
function voicePresenceHeartbeatActiveKey(params: VoicePresenceHeartbeatParams): string {
return `${VOICE_PRESENCE_HEARTBEAT_ACTIVE_PREFIX}${voicePresenceHeartbeatKeySuffix(params)}`;
}
function voicePresenceHeartbeatEnrollmentKey(params: VoicePresenceHeartbeatParams): string {
return `${VOICE_PRESENCE_HEARTBEAT_ENROLLED_PREFIX}${voicePresenceHeartbeatKeySuffix(params)}`;
}
function voicePresenceHeartbeatChannelKey(channelId: ChannelID): string {
return `${VOICE_PRESENCE_HEARTBEAT_CHANNEL_PREFIX}${channelId.toString()}`;
}
function voicePresenceHeartbeatKeySuffix(params: VoicePresenceHeartbeatParams): string {
const connectionHash = createHash('sha256').update(params.connectionId).digest('base64url');
return `channel:${params.channelId.toString()}:user:${params.userId.toString()}:connection:${connectionHash}`;
}
@@ -12,7 +12,6 @@ import type {GatewayVoiceStateEntry, IGatewayService} from '../infrastructure/IG
import type {ILiveKitService, LiveKitRoomLocation} from '../infrastructure/ILiveKitService';
import type {IVoiceRoomStore} from '../infrastructure/IVoiceRoomStore';
import {parseParticipantIdentity, parseRoomName} from '../infrastructure/VoiceRoomContext';
import {type VoicePresenceHeartbeatState, VoicePresenceHeartbeatStore} from './VoicePresenceHeartbeatStore';
interface GatewayPendingJoinEntry {
readonly connectionId: string;
@@ -27,7 +26,6 @@ interface VoiceReconciliationWorkerOptions {
voiceRoomStore: IVoiceRoomStore;
kvClient: IKVProvider;
logger: ILogger;
voicePresenceHeartbeatStore?: VoicePresenceHeartbeatStore;
intervalMs?: number;
staggerDelayMs?: number;
lockTtlSeconds?: number;
@@ -122,7 +120,6 @@ export class VoiceReconciliationWorker {
private readonly liveKitService: ILiveKitService;
private readonly voiceRoomStore: IVoiceRoomStore;
private readonly kvClient: IKVProvider;
private readonly voicePresenceHeartbeatStore: VoicePresenceHeartbeatStore;
private readonly logger: ILogger;
private readonly intervalMs: number;
private readonly staggerDelayMs: number;
@@ -140,8 +137,6 @@ export class VoiceReconciliationWorker {
this.liveKitService = options.liveKitService;
this.voiceRoomStore = options.voiceRoomStore;
this.kvClient = options.kvClient;
this.voicePresenceHeartbeatStore =
options.voicePresenceHeartbeatStore ?? new VoicePresenceHeartbeatStore(this.kvClient);
this.logger = options.logger.child({worker: 'VoiceReconciliationWorker'});
this.intervalMs = options.intervalMs ?? DEFAULT_INTERVAL_MS;
this.staggerDelayMs = options.staggerDelayMs ?? DEFAULT_STAGGER_DELAY_MS;
@@ -575,34 +570,6 @@ export class VoiceReconciliationWorker {
livekitOnlyDeferred++;
continue;
}
const heartbeatState = await this.getVoicePresenceHeartbeatState(room.channelId, participant);
if (heartbeatState === 'active') {
await this.clearLiveKitOnlyCandidate(room.guildId, room.channelId, participant);
this.logger.debug(
{
roomName: room.roomName,
userId: participant.userId.toString(),
connectionId: participant.connectionId,
},
'Deferring LiveKit-only participant because v2 voice presence heartbeat is active',
);
livekitOnlyDeferred++;
continue;
}
if (heartbeatState === 'expired') {
await this.clearLiveKitOnlyCandidate(room.guildId, room.channelId, participant);
this.logger.warn(
{
roomName: room.roomName,
userId: participant.userId.toString(),
connectionId: participant.connectionId,
},
'Disconnecting LiveKit-only participant because v2 voice presence heartbeat expired',
);
await this.disconnectLiveKitOnlyParticipant(room, participant);
livekitOnlyDisconnected++;
continue;
}
if (!liveKitSnapshot.completed) {
this.logger.warn(
{
@@ -649,7 +616,7 @@ export class VoiceReconciliationWorker {
gatewayOnlySkipped++;
continue;
}
const shouldRemove = await this.shouldRemoveGatewayOnlyState(room, voiceState);
const shouldRemove = await this.confirmGatewayOnlyCandidate(room.guildId, room.channelId, voiceState);
if (!shouldRemove) {
gatewayOnlyDeferred++;
continue;
@@ -954,67 +921,6 @@ export class VoiceReconciliationWorker {
}
}
private async shouldRemoveGatewayOnlyState(
room: DiscoveredRoom,
voiceState: GatewayVoiceStateEntry,
): Promise<boolean> {
let userId: UserID;
try {
userId = createUserID(BigInt(voiceState.userId));
} catch (error) {
this.logger.warn(
{error, roomName: room.roomName, userId: voiceState.userId, connectionId: voiceState.connectionId},
'Falling back to legacy gateway-only reconciliation because voice state user id was invalid',
);
return this.confirmGatewayOnlyCandidate(room.guildId, room.channelId, voiceState);
}
const heartbeatState = await this.getVoicePresenceHeartbeatState(room.channelId, {
userId,
connectionId: voiceState.connectionId,
});
if (heartbeatState === 'active') {
await this.clearGatewayOnlyCandidate(room.guildId, room.channelId, voiceState.connectionId);
this.logger.debug(
{roomName: room.roomName, userId: voiceState.userId, connectionId: voiceState.connectionId},
'Deferring gateway-only state because v2 voice presence heartbeat is active',
);
return false;
}
if (heartbeatState === 'expired') {
await this.clearGatewayOnlyCandidate(room.guildId, room.channelId, voiceState.connectionId);
this.logger.warn(
{roomName: room.roomName, userId: voiceState.userId, connectionId: voiceState.connectionId},
'Removing gateway-only state because v2 voice presence heartbeat expired',
);
return true;
}
return this.confirmGatewayOnlyCandidate(room.guildId, room.channelId, voiceState);
}
private async getVoicePresenceHeartbeatState(
channelId: ChannelID,
connection: {userId: UserID; connectionId: string},
): Promise<VoicePresenceHeartbeatState> {
try {
return await this.voicePresenceHeartbeatStore.getHeartbeatState({
channelId,
userId: connection.userId,
connectionId: connection.connectionId,
});
} catch (error) {
this.logger.warn(
{
error,
channelId: channelId.toString(),
userId: connection.userId.toString(),
connectionId: connection.connectionId,
},
'Falling back to legacy reconciliation because v2 voice presence heartbeat lookup failed',
);
return 'legacy';
}
}
private async confirmGatewayOnlyCandidate(
guildId: GuildID | undefined,
channelId: ChannelID,
@@ -83,9 +83,6 @@ export const Endpoints = {
CHANNEL_CALL: (channelId: string) => `/channels/${channelId}/call`,
CHANNEL_CALL_RING: (channelId: string) => `/channels/${channelId}/call/ring`,
CHANNEL_CALL_STOP_RINGING: (channelId: string) => `/channels/${channelId}/call/stop-ringing`,
CHANNEL_VOICE_DEBUG_LOGGING_SESSION: (channelId: string) => `/channels/${channelId}/voice-debug-logging/session`,
CHANNEL_VOICE_DEBUG_LOGGING_EVENTS: (channelId: string) => `/channels/${channelId}/voice-debug-logging/events`,
CHANNEL_VOICE_PRESENCE_HEARTBEAT: (channelId: string) => `/channels/${channelId}/voice-presence/heartbeat`,
GUILDS: '/guilds',
GUILD: (guildId: string) => `/guilds/${guildId}`,
GUILD_CHANNELS: (guildId: string) => `/guilds/${guildId}/channels`,
@@ -9,7 +9,7 @@ import {
renderVoiceDebugStatsHtml,
renderVoiceDebugStatsUnavailableHtml,
} from '@app/features/voice/diagnostics/VoiceDebugStatsHtml';
import voiceEngineV2AppDebugLoggingHostAdapter from '@app/features/voice/engine/v2/VoiceEngineV2AppDebugLoggingHostAdapter';
import voiceEngineV2AppDebugEventSinkHostAdapter from '@app/features/voice/engine/v2/VoiceEngineV2AppDebugEventSinkHostAdapter';
import {collectStatsForNerdsSnapshot} from '@app/features/voice/utils/StatsForNerdsCopy';
export function canOpenVoiceDebugEventSinkPopout(): boolean {
@@ -41,5 +41,5 @@ function publishBrowserOpeningStatsSnapshot(): void {
export async function openVoiceDebugEventSinkPopout(): Promise<void> {
publishBrowserOpeningStatsSnapshot();
await voiceEngineV2AppDebugLoggingHostAdapter.openEventSinkPopout();
await voiceEngineV2AppDebugEventSinkHostAdapter.openEventSinkPopout();
}
@@ -1,72 +0,0 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {Endpoints} from '@app/features/app/constants/Endpoints';
import {http} from '@app/features/platform/transport/RestTransport';
import type {
VoiceDebugLoggingEventsResponse,
VoiceDebugLoggingStatusResponse,
} from '@fluxer/schema/src/domains/channel/ChannelSchemas';
export interface VoiceDebugLoggingEvent {
type: string;
timestamp_ns: string;
monotonic_ns?: string;
data?: Record<string, unknown>;
}
export async function fetchStatus(channelId: string): Promise<VoiceDebugLoggingStatusResponse> {
const response = await http.get<VoiceDebugLoggingStatusResponse>(
Endpoints.CHANNEL_VOICE_DEBUG_LOGGING_SESSION(channelId),
{retries: 1},
);
return (
response.body ?? {
active: false,
session_id: null,
activated_by_user_id: null,
started_at_ms: null,
expires_at_ms: null,
poll_interval_ms: 10000,
upload_interval_ms: 2000,
}
);
}
export async function setEnabled(channelId: string, enabled: boolean): Promise<VoiceDebugLoggingStatusResponse> {
const response = await http.put<VoiceDebugLoggingStatusResponse>(
Endpoints.CHANNEL_VOICE_DEBUG_LOGGING_SESSION(channelId),
{body: {enabled}, retries: 1},
);
const fallback = response.body ?? {
active: false,
session_id: null,
activated_by_user_id: null,
started_at_ms: null,
expires_at_ms: null,
poll_interval_ms: 10000,
upload_interval_ms: 2000,
};
return await fetchStatus(channelId).catch(() => fallback);
}
export async function uploadEvents(params: {
channelId: string;
sessionId: string;
connectionId: string | null;
participantIdentity: string | null;
events: Array<VoiceDebugLoggingEvent>;
}): Promise<VoiceDebugLoggingEventsResponse> {
const response = await http.post<VoiceDebugLoggingEventsResponse>(
Endpoints.CHANNEL_VOICE_DEBUG_LOGGING_EVENTS(params.channelId),
{
body: {
session_id: params.sessionId,
connection_id: params.connectionId ?? undefined,
participant_identity: params.participantIdentity ?? undefined,
events: params.events,
},
retries: 1,
},
);
return response.body ?? {accepted: false, active: false, stored_event_count: 0};
}
@@ -1,51 +0,0 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {Endpoints} from '@app/features/app/constants/Endpoints';
import {http} from '@app/features/platform/transport/RestTransport';
import type {
VoicePresenceHeartbeatEndResponse,
VoicePresenceHeartbeatResponse,
} from '@fluxer/schema/src/domains/channel/ChannelSchemas';
export async function heartbeat(params: {
channelId: string;
connectionId: string;
}): Promise<VoicePresenceHeartbeatResponse> {
const response = await http.post<VoicePresenceHeartbeatResponse>(
Endpoints.CHANNEL_VOICE_PRESENCE_HEARTBEAT(params.channelId),
{
body: {connection_id: params.connectionId},
mode: 'silent',
retries: 1,
timeoutMs: 5000,
},
);
if (response.ok && response.body && typeof response.body.ok === 'boolean') {
return response.body;
}
return {
ok: false,
heartbeat_interval_ms: 15000,
heartbeat_ttl_ms: 45000,
expires_at_ms: 0,
};
}
export async function end(params: {
channelId: string;
connectionId: string;
}): Promise<VoicePresenceHeartbeatEndResponse> {
const response = await http.delete<VoicePresenceHeartbeatEndResponse>(
Endpoints.CHANNEL_VOICE_PRESENCE_HEARTBEAT(params.channelId),
{
body: {connection_id: params.connectionId},
mode: 'silent',
retries: 1,
timeoutMs: 5000,
},
);
if (response.ok && response.body && typeof response.body.ok === 'boolean') {
return response.body;
}
return {ok: false};
}
@@ -647,9 +647,10 @@ export const VoiceMoreOptionsMenu: React.FC<VoiceMoreOptionsMenuProps> = observe
const layoutMode = VoiceCallLayout.layoutMode;
const isGrid = layoutMode === 'grid';
const connectedChannelId = MediaEngine.channelId;
const canControlDebugLogging = connectedChannelId != null && (Users.currentUser?.isStaff() ?? false);
const canOpenDebugEventSink =
canControlDebugLogging && VoiceDebugEventSinkCommands.canOpenVoiceDebugEventSinkPopout();
connectedChannelId != null &&
(Users.currentUser?.isStaff() ?? false) &&
VoiceDebugEventSinkCommands.canOpenVoiceDebugEventSinkPopout();
const isDmVoiceCall = connectedChannelId != null && (MediaEngine.guildId ?? null) === null;
const currentRegion =
isDmVoiceCall && connectedChannelId
@@ -841,25 +842,6 @@ export const VoiceMoreOptionsMenu: React.FC<VoiceMoreOptionsMenuProps> = observe
>
{i18n._(PRIORITIZE_SPEAKERS_DESCRIPTOR)}
</CheckboxItem>
{canControlDebugLogging && (
<CheckboxItem
icon={
<ChartBarIcon
weight="fill"
className={styles.icon}
data-flx="voice.voice-settings-menus.voice-more-options-menu.icon.debug-logging"
/>
}
checked={MediaEngine.voiceDebugLoggingActive}
disabled={MediaEngine.voiceDebugLoggingToggleInFlight}
onCheckedChange={(checked) => {
void MediaEngine.setVoiceDebugLoggingEnabled(checked);
}}
data-flx="voice.voice-settings-menus.voice-more-options-menu.checkbox-item.debug-logging"
>
<Trans>Start debug logging session</Trans>
</CheckboxItem>
)}
{canOpenDebugEventSink && (
<MenuItem
icon={
@@ -59,11 +59,7 @@ import {
saveCurrentVoiceSessionRestoreSnapshot,
type VoiceSessionRestoreSyncHandle,
} from '@app/features/voice/engine/media_engine_facade/VoiceSessionRestoreSync';
import {getScreenShareCaptureDiagnosticSnapshot} from '@app/features/voice/engine/ScreenShareCaptureDiagnostics';
import ScreenShareCodecNegotiation, {
buildLocalCodecAdvertisements,
getScreenShareCodecPreferenceOrder,
} from '@app/features/voice/engine/ScreenShareCodecNegotiation';
import ScreenShareCodecNegotiation from '@app/features/voice/engine/ScreenShareCodecNegotiation';
import ScreenSharePublicationMigration from '@app/features/voice/engine/ScreenSharePublicationMigration';
import {Store, useStoreVersion} from '@app/features/voice/engine/Store';
import {shouldMoveToAfkOnTick} from '@app/features/voice/engine/VoiceAfkTracking';
@@ -118,7 +114,7 @@ import {
createVoiceEngineV2AppControllerHost,
type VoiceEngineV2AppControllerHost,
} from '@app/features/voice/engine/v2/VoiceEngineV2AppControllerHost';
import voiceEngineV2AppDebugLoggingHostAdapter from '@app/features/voice/engine/v2/VoiceEngineV2AppDebugLoggingHostAdapter';
import voiceEngineV2AppDebugEventSinkHostAdapter from '@app/features/voice/engine/v2/VoiceEngineV2AppDebugEventSinkHostAdapter';
import {createVoiceEngineV2AppEventLogSpillLoggerSink} from '@app/features/voice/engine/v2/VoiceEngineV2AppEventLogSpillLoggerSink';
import {createVoiceEngineV2AppIngestionPort} from '@app/features/voice/engine/v2/VoiceEngineV2AppHostPorts';
import type {VoiceEngineV2AppLifecycleDisposable} from '@app/features/voice/engine/v2/VoiceEngineV2AppLifecycleAdapter';
@@ -153,18 +149,9 @@ import VoiceCallLayout from '@app/features/voice/state/VoiceCallLayout';
import VoiceRegionTeleport from '@app/features/voice/state/VoiceRegionTeleport';
import VoiceSessionRestore, {type VoiceSessionRestoreSnapshot} from '@app/features/voice/state/VoiceSessionRestore';
import VoiceSettings from '@app/features/voice/state/VoiceSettings';
import {type CodecPreference, getCodecCapabilityReport} from '@app/features/voice/utils/CodecCapabilityDetector';
import {getGpuEncoderReportSync, loadGpuEncoderReport} from '@app/features/voice/utils/GpuEncoderCapabilities';
import {
getNativeAudioCaptureDiagnosticState,
setNativeAudioCaptureBridgeLifecycleBridge,
} from '@app/features/voice/utils/NativeAudioCaptureBridge';
import {getOpenH264StatusSync, loadOpenH264Status} from '@app/features/voice/utils/OpenH264Status';
import type {CodecPreference} from '@app/features/voice/utils/CodecCapabilityDetector';
import {setNativeAudioCaptureBridgeLifecycleBridge} from '@app/features/voice/utils/NativeAudioCaptureBridge';
import {areOrderedStringArraysEqual} from '@app/features/voice/utils/StringArrayUtils';
import {
getVideoDecoderExclusionsSync,
loadVideoDecoderExclusions,
} from '@app/features/voice/utils/VideoDecoderCapabilities';
import {buildVoiceParticipantIdentity} from '@app/features/voice/utils/VoiceParticipantIdentity';
import {
isVoiceServerMuteActive,
@@ -212,8 +199,6 @@ import {
import {makeObservable, observable} from 'mobx';
const logger = new Logger('MediaEngineFacade');
const MAX_DIAGNOSTIC_LINUX_AUDIO_TARGETS = 200;
const MAX_DIAGNOSTIC_PIPEWIRE_GRAPH_RECORDS = 500;
function isCameraPermissionDeniedCommandFailure(error: unknown): boolean {
if (!(error instanceof Error)) return false;
@@ -318,137 +303,6 @@ function buildRoomEventDependencies(facade: {
};
}
function voiceDiagnosticError(error: unknown): Record<string, unknown> {
if (error instanceof Error) {
return {
name: error.name,
message: error.message,
stack: error.stack,
};
}
return {message: String(error)};
}
function limitLinuxAudioTargets(result: unknown): unknown {
if (!result || typeof result !== 'object' || !('targets' in result)) return result;
const record = result as {targets?: unknown};
if (!Array.isArray(record.targets)) return result;
return {
...result,
targets: record.targets.slice(0, MAX_DIAGNOSTIC_LINUX_AUDIO_TARGETS),
targetCount: record.targets.length,
truncated: record.targets.length > MAX_DIAGNOSTIC_LINUX_AUDIO_TARGETS,
};
}
function limitGraphArray(record: Record<string, unknown>, key: string): Record<string, unknown> {
const value = record[key];
if (!Array.isArray(value)) return record;
return {
...record,
[key]: value.slice(0, MAX_DIAGNOSTIC_PIPEWIRE_GRAPH_RECORDS),
[`${key}Count`]: value.length,
[`${key}Truncated`]: value.length > MAX_DIAGNOSTIC_PIPEWIRE_GRAPH_RECORDS,
};
}
function limitPipewireRoutingGraph(graph: unknown): unknown {
if (!graph || typeof graph !== 'object') return graph;
let limited = graph as Record<string, unknown>;
limited = limitGraphArray(limited, 'nodes');
limited = limitGraphArray(limited, 'ports');
limited = limitGraphArray(limited, 'ownedLinks');
return limited;
}
function limitPipewireRoutingGraphResult(result: unknown): unknown {
if (!result || typeof result !== 'object') return result;
const record = result as Record<string, unknown>;
if ('graph' in record) {
return {...record, graph: limitPipewireRoutingGraph(record.graph)};
}
if (Array.isArray(record.graphs)) {
return {
...record,
graphs: record.graphs.map((entry) =>
entry && typeof entry === 'object'
? {...entry, graph: limitPipewireRoutingGraph((entry as Record<string, unknown>).graph)}
: entry,
),
};
}
return result;
}
function hasPipewireRoutingGraph(result: unknown): boolean {
if (!result || typeof result !== 'object') return false;
const record = result as Record<string, unknown>;
if (record.ok !== true) return false;
if (record.graph && typeof record.graph === 'object') return true;
return (
Array.isArray(record.graphs) &&
record.graphs.some((entry) => {
if (!entry || typeof entry !== 'object') return false;
const graph = (entry as Record<string, unknown>).graph;
return Boolean(graph && typeof graph === 'object');
})
);
}
async function createVoiceAudioDiagnosticsSnapshot(): Promise<Record<string, unknown>> {
const electron = getElectronAPI();
const [
nativeAudioAvailability,
nativeAudioApplications,
nativeAudioRoutingGraph,
virtmicAvailability,
virtmicTargets,
virtmicRoutingGraph,
] = await Promise.all([
electron?.nativeAudio
? electron.nativeAudio.getAvailability().catch((error) => ({error: voiceDiagnosticError(error)}))
: Promise.resolve(null),
electron?.nativeAudio
? electron.nativeAudio.listAudibleApplications().catch((error) => ({error: voiceDiagnosticError(error)}))
: Promise.resolve(null),
electron?.nativeAudio
? electron.nativeAudio.getRoutingGraph().catch((error) => ({error: voiceDiagnosticError(error)}))
: Promise.resolve(null),
electron?.virtmic
? electron.virtmic.getAvailability().catch((error) => ({error: voiceDiagnosticError(error)}))
: Promise.resolve(null),
electron?.virtmic
? electron.virtmic.listTargets({granular: true}).catch((error) => ({error: voiceDiagnosticError(error)}))
: Promise.resolve(null),
electron?.virtmic
? electron.virtmic.getRoutingGraph().catch((error) => ({error: voiceDiagnosticError(error)}))
: Promise.resolve(null),
]);
return {
platform: electron?.platform ?? null,
nativeAudio: {
availability: nativeAudioAvailability,
audibleApplications: nativeAudioApplications,
capture: getNativeAudioCaptureDiagnosticState(),
pipewireRoutingGraph: limitPipewireRoutingGraphResult(nativeAudioRoutingGraph),
},
linuxRouting: {
virtmicAvailability,
pipewireNodeInventory: limitLinuxAudioTargets(virtmicTargets),
pipewireRoutingGraph: limitPipewireRoutingGraphResult(virtmicRoutingGraph),
activeLinkGraphExported:
hasPipewireRoutingGraph(virtmicRoutingGraph) || hasPipewireRoutingGraph(nativeAudioRoutingGraph),
routingSettings: {
screenShareAudioSourceMode: VoiceSettings.getScreenShareAudioSourceMode(),
screenShareAudioIncludeSources: VoiceSettings.getScreenShareAudioIncludeSources(),
screenShareAudioExcludeSources: VoiceSettings.getScreenShareAudioExcludeSources(),
shareAppAudio: VoiceSettings.getShareAppAudio(),
shareDesktopAudio: VoiceSettings.getShareDesktopAudio(),
},
},
};
}
interface ConnectToVoiceChannelOptions {
skipChannelGate?: boolean;
deferNavigationUntilConnected?: boolean;
@@ -656,7 +510,6 @@ class MediaEngineFacade extends Store {
voiceEngineV2AppScreenShareExecutionAdapter.subscribe(forwardChange);
VoiceEngineV2AppSubscriptionAdapter.subscribe(forwardChange);
VoiceEngineV2AppPermissionAdapter.subscribe(forwardChange);
voiceEngineV2AppDebugLoggingHostAdapter.subscribe(forwardChange);
ScreenSharePublicationMigration.subscribe(forwardChange);
}
@@ -817,18 +670,6 @@ class MediaEngineFacade extends Store {
return voiceEngineV2AppConnectionHostAdapter.voiceServerEndpoint;
}
get voiceDebugLoggingActive(): boolean {
return voiceEngineV2AppDebugLoggingHostAdapter.active;
}
get voiceDebugLoggingToggleInFlight(): boolean {
return voiceEngineV2AppDebugLoggingHostAdapter.toggleInFlight;
}
async setVoiceDebugLoggingEnabled(enabled: boolean): Promise<void> {
await voiceEngineV2AppDebugLoggingHostAdapter.setEnabled(enabled);
}
get connectionContext(): VoiceEngineConnectionContext {
return {
guildId: this.guildId,
@@ -2522,19 +2363,18 @@ class MediaEngineFacade extends Store {
this.startAfkTracking();
const channelId = voiceEngineV2AppConnectionHostAdapter.channelId;
if (channelId) {
void voiceEngineV2AppDebugLoggingHostAdapter.start({
voiceEngineV2AppDebugEventSinkHostAdapter.start({
guildId: voiceEngineV2AppConnectionHostAdapter.guildId,
channelId,
connectionId: voiceEngineV2AppConnectionHostAdapter.connectionId,
room,
collectSnapshot: () => this.createVoiceDiagnosticsSnapshot(),
});
}
logger.info('All tracking started');
}
private stopTracking(): void {
void voiceEngineV2AppDebugLoggingHostAdapter.stop('tracking-stopped');
voiceEngineV2AppDebugEventSinkHostAdapter.stop('tracking-stopped');
this.statsHostAdapter.stopLatencyTracking();
this.statsHostAdapter.stopStatsTracking();
VoiceEngineV2AppRemoteSpeakingAdapter.clear();
@@ -2555,83 +2395,6 @@ class MediaEngineFacade extends Store {
logger.info('All tracking stopped');
}
private async createVoiceDiagnosticsSnapshot(): Promise<Record<string, unknown>> {
await Promise.allSettled([loadGpuEncoderReport(), loadOpenH264Status(), loadVideoDecoderExclusions()]);
return {
connection: {
guildId: voiceEngineV2AppConnectionHostAdapter.guildId,
channelId: voiceEngineV2AppConnectionHostAdapter.channelId,
connectionId: voiceEngineV2AppConnectionHostAdapter.connectionId,
connected: voiceEngineV2AppConnectionHostAdapter.connected,
connecting: voiceEngineV2AppConnectionHostAdapter.connecting,
reconnecting: voiceEngineV2AppConnectionHostAdapter.reconnecting,
disconnecting: voiceEngineV2AppConnectionHostAdapter.disconnecting,
voiceServerEndpoint: voiceEngineV2AppConnectionHostAdapter.voiceServerEndpoint,
regionHotSwapInProgress: voiceEngineV2AppConnectionHostAdapter.regionHotSwapInProgress,
reconnectAttempts: voiceEngineV2AppConnectionHostAdapter.reconnectAttempts,
},
localVoiceState: {
selfMute: LocalVoiceState.getSelfMute(),
selfDeaf: LocalVoiceState.getSelfDeaf(),
selfVideo: LocalVoiceState.getSelfVideo(),
selfStream: LocalVoiceState.getSelfStream(),
selfStreamAudio: LocalVoiceState.getSelfStreamAudio(),
selfStreamAudioMute: LocalVoiceState.getSelfStreamAudioMute(),
viewerStreamKeys: LocalVoiceState.getViewerStreamKeys(),
hasUserSetMute: LocalVoiceState.getHasUserSetMute(),
hasUserSetDeaf: LocalVoiceState.getHasUserSetDeaf(),
},
effectiveAudioState: getEffectiveAudioState(),
connectionVoiceStates: voiceEngineV2AppVoiceStateAdapter.getConnectionVoiceStates(),
participants: Object.values(this.voiceEngineV2Participants.participants),
stats: {
currentLatency: this.currentLatency,
averageLatency: this.averageLatency,
displayLatency: this.displayLatency,
estimatedLatency: this.estimatedLatency,
reconnectionCount: this.reconnectionCount,
voiceStats: this.voiceStats,
perTrackStats: this.perTrackStats,
statsTimeSeriesTail: this.statsTimeSeries.slice(-30),
latencyHistoryTail: this.latencyHistory.slice(-30),
publisherTransport: this.publisherTransport,
subscriberTransport: this.subscriberTransport,
},
screenShare: {
selectedCodec: ScreenShareCodecNegotiation.getSelectedCodec(),
preferenceOrder: getScreenShareCodecPreferenceOrder(),
capture: await getScreenShareCaptureDiagnosticSnapshot(),
audioCapture: await createVoiceAudioDiagnosticsSnapshot(),
codecCapabilities: {
localAdvertisements: buildLocalCodecAdvertisements(),
report: getCodecCapabilityReport(),
gpuEncoderReport: getGpuEncoderReportSync(),
openH264Status: getOpenH264StatusSync(),
decoderExclusions: getVideoDecoderExclusionsSync(),
},
adaptiveQuality: AdaptiveScreenShareEngine.qualitySnapshot,
settings: {
preferredScreenShareCodec: VoiceSettings.getPreferredScreenShareCodec(),
screenShareEncoderMode: VoiceSettings.getScreenShareEncoderMode(),
screenShareSoftwareQuality: VoiceSettings.getScreenShareSoftwareQuality(),
screenShareScalabilityMode: VoiceSettings.getScreenShareScalabilityMode(),
screenShareBackupCodecMode: VoiceSettings.getScreenShareBackupCodecMode(),
screenShareMaxBitrateMbps: VoiceSettings.getScreenShareMaxBitrateMbps(),
adaptiveScreenShareQuality: VoiceSettings.getAdaptiveScreenShareQuality(),
screenshareResolution: VoiceSettings.getScreenshareResolution(),
videoFrameRate: VoiceSettings.getVideoFrameRate(),
streamingMode: VoiceSettings.getStreamingMode(),
shareAppAudio: VoiceSettings.getShareAppAudio(),
shareDesktopAudio: VoiceSettings.getShareDesktopAudio(),
shareDeviceAudio: VoiceSettings.getShareDeviceAudio(),
screenShareAudioSourceMode: VoiceSettings.getScreenShareAudioSourceMode(),
screenShareAudioIncludeSources: VoiceSettings.getScreenShareAudioIncludeSources(),
screenShareAudioExcludeSources: VoiceSettings.getScreenShareAudioExcludeSources(),
},
},
};
}
private abortVoiceConnection(): void {
this.transitionFacadeState({type: 'screenShareReconnect.clear'});
voiceEngineV2AppConnectionHostAdapter.abortConnection();
@@ -1,226 +0,0 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {getElectronAPI} from '@app/features/ui/utils/NativeUtils';
import type {
NativeScreenCaptureAvailability,
NativeScreenCaptureDiagnostics,
NativeScreenCaptureSourceKind,
} from '@app/types/electron.d';
export type ScreenShareCaptureMethod =
| 'display-media'
| 'native-screen-capture'
| 'native-voice-engine-screen-capture'
| 'device-media';
interface ScreenShareCaptureState {
active: boolean;
generation: number;
method: ScreenShareCaptureMethod;
startedAtMs: number;
updatedAtMs: number;
endedAtMs?: number;
endReason?: string;
platform?: string | null;
sourceId?: string | null;
desktopSourceId?: string | null;
sourceKind?: NativeScreenCaptureSourceKind | null;
captureId?: string | null;
displayShareEnvironment?: string | null;
displayMediaSettings?: Record<string, unknown> | null;
device?: {
videoDeviceId?: string;
audioDeviceId?: string;
};
nativeAvailability?: NativeScreenCaptureAvailability | null;
nativeDiagnostics?: NativeScreenCaptureDiagnostics | null;
}
interface CaptureStateUpdate {
method: ScreenShareCaptureMethod;
sourceId?: string | null;
desktopSourceId?: string | null;
sourceKind?: NativeScreenCaptureSourceKind | null;
captureId?: string | null;
displayShareEnvironment?: string | null;
displayMediaSettings?: Record<string, unknown> | null;
device?: {
videoDeviceId?: string;
audioDeviceId?: string;
};
nativeAvailability?: NativeScreenCaptureAvailability | null;
nativeDiagnostics?: NativeScreenCaptureDiagnostics | null;
}
interface DisplayMediaTrackSettingsUpdate {
sourceId?: string | null;
desktopSourceId?: string | null;
displayShareEnvironment?: string | null;
}
interface MediaTrackSettingsReader {
getSettings(): MediaTrackSettings;
}
const MAX_CAPTURE_HISTORY = 20;
let activeCapture: ScreenShareCaptureState | null = null;
let lastCapture: ScreenShareCaptureState | null = null;
let captureGeneration = 0;
const captureHistory: Array<ScreenShareCaptureState | null> = new Array(MAX_CAPTURE_HISTORY).fill(null);
let captureHistoryHead = 0;
let captureHistoryLength = 0;
function nowMs(): number {
return Date.now();
}
function readPlatform(): string | null {
return getElectronAPI()?.platform ?? null;
}
function cloneCaptureState(state: ScreenShareCaptureState | null): ScreenShareCaptureState | null {
if (!state) return null;
const cloned = {...state};
if (state.device) {
cloned.device = {...state.device};
}
return cloned;
}
function cloneCaptureHistoryTail(limit: number): Array<ScreenShareCaptureState | null> {
const tailCount = Math.min(captureHistoryLength, limit);
const tail = new Array<ScreenShareCaptureState | null>(tailCount);
const startIndex = captureHistoryLength - tailCount;
for (let i = 0; i < tailCount; i += 1) {
const index = (captureHistoryHead + startIndex + i) % MAX_CAPTURE_HISTORY;
tail[i] = cloneCaptureState(captureHistory[index] ?? null);
}
return tail;
}
function pushHistory(state: ScreenShareCaptureState): void {
const cloned = cloneCaptureState(state);
if (!cloned) return;
if (captureHistoryLength < MAX_CAPTURE_HISTORY) {
const index = (captureHistoryHead + captureHistoryLength) % MAX_CAPTURE_HISTORY;
captureHistory[index] = cloned;
captureHistoryLength += 1;
return;
}
captureHistory[captureHistoryHead] = cloned;
captureHistoryHead = (captureHistoryHead + 1) % MAX_CAPTURE_HISTORY;
}
export function markScreenShareCaptureActive(update: CaptureStateUpdate): void {
const currentTime = nowMs();
const canUpdateCurrent =
activeCapture != null &&
activeCapture.method === update.method &&
(update.captureId == null || activeCapture.captureId === update.captureId);
const base = canUpdateCurrent ? activeCapture : null;
const device = update.device ?? base?.device;
const nextCapture: ScreenShareCaptureState = {
active: true,
generation: base?.generation ?? ++captureGeneration,
method: update.method,
startedAtMs: base?.startedAtMs ?? currentTime,
updatedAtMs: currentTime,
platform: readPlatform(),
sourceId: update.sourceId ?? base?.sourceId ?? null,
desktopSourceId: update.desktopSourceId ?? base?.desktopSourceId ?? null,
sourceKind: update.sourceKind ?? base?.sourceKind ?? null,
captureId: update.captureId ?? base?.captureId ?? null,
displayShareEnvironment: update.displayShareEnvironment ?? base?.displayShareEnvironment ?? null,
displayMediaSettings: update.displayMediaSettings ?? base?.displayMediaSettings ?? null,
nativeAvailability: update.nativeAvailability ?? base?.nativeAvailability ?? null,
nativeDiagnostics: update.nativeDiagnostics ?? base?.nativeDiagnostics ?? null,
};
if (device) {
nextCapture.device = {...device};
}
activeCapture = nextCapture;
}
export function updateScreenShareDisplayMediaSettings(
track: MediaTrackSettingsReader,
extra?: DisplayMediaTrackSettingsUpdate,
): void {
const settings = track.getSettings();
markScreenShareCaptureActive({
method: 'display-media',
...extra,
displayMediaSettings: {
width: settings.width,
height: settings.height,
frameRate: settings.frameRate,
aspectRatio: settings.aspectRatio,
deviceId: settings.deviceId,
displaySurface: (settings as MediaTrackSettings & {displaySurface?: string}).displaySurface,
cursor: (settings as MediaTrackSettings & {cursor?: string}).cursor,
logicalSurface: (settings as MediaTrackSettings & {logicalSurface?: boolean}).logicalSurface,
},
});
}
export function markScreenShareCaptureEnded(reason: string): void {
if (!activeCapture) return;
const ended = {
...activeCapture,
active: false,
updatedAtMs: nowMs(),
endedAtMs: nowMs(),
endReason: reason,
};
lastCapture = ended;
pushHistory(ended);
activeCapture = null;
}
function inferWindowsCaptureMethod(
capture: ScreenShareCaptureState | null,
diagnostics: NativeScreenCaptureDiagnostics | null,
): string | null {
if (!capture || capture.platform !== 'win32') return null;
if (capture.method === 'display-media') return 'chromium-getDisplayMedia-desktopCapturer';
if (capture.method === 'device-media') return 'getUserMedia-device';
if (capture.method !== 'native-screen-capture' && capture.method !== 'native-voice-engine-screen-capture') {
return null;
}
if (diagnostics?.activeStrategy) return diagnostics.activeStrategy;
if (capture.sourceKind === 'game') return 'wgc';
if (capture.sourceKind === 'screen') return 'wgc';
if (capture.sourceKind === 'window') return 'dxgi-duplication';
return 'native-screen-capture';
}
async function readNativeAvailability(): Promise<NativeScreenCaptureAvailability | null> {
const api = getElectronAPI()?.nativeScreenCapture;
if (!api) return null;
return api.getAvailability().catch(() => null);
}
async function readNativeDiagnostics(
captureId: string | null | undefined,
): Promise<NativeScreenCaptureDiagnostics | null> {
const api = getElectronAPI()?.nativeScreenCapture;
if (!api || !captureId) return null;
return api.getDiagnostics(captureId).catch(() => null);
}
export async function getScreenShareCaptureDiagnosticSnapshot(): Promise<Record<string, unknown>> {
const active = cloneCaptureState(activeCapture);
const nativeAvailability = await readNativeAvailability();
const nativeDiagnostics = await readNativeDiagnostics(active?.captureId);
if (active) {
active.nativeAvailability = nativeAvailability ?? active.nativeAvailability ?? null;
active.nativeDiagnostics = nativeDiagnostics ?? active.nativeDiagnostics ?? null;
}
return {
active,
last: cloneCaptureState(lastCapture),
historyTail: cloneCaptureHistoryTail(5),
nativeAvailability,
windowsCaptureMethod: inferWindowsCaptureMethod(active, nativeDiagnostics ?? active?.nativeDiagnostics ?? null),
};
}
@@ -14,7 +14,6 @@ import {
hasAnyTerminalTransport,
isMutedOrDeafened,
isPermissionDeniedError,
isPresenceConnectionReady,
isReadyToRepublishTrack,
} from './VoiceEngineV2AppAdapterAssertions';
@@ -188,26 +187,6 @@ describe('VoiceEngineV2AppAdapterAssertions', () => {
});
});
describe('isPresenceConnectionReady', () => {
it('returns true when connected and ids are present', () => {
expect(isPresenceConnectionReady(true, 'channel', 'connection')).toBe(true);
});
it('returns false when disconnected', () => {
expect(isPresenceConnectionReady(false, 'channel', 'connection')).toBe(false);
});
it('returns false on missing channelId', () => {
expect(isPresenceConnectionReady(true, null, 'connection')).toBe(false);
expect(isPresenceConnectionReady(true, '', 'connection')).toBe(false);
});
it('returns false on missing connectionId', () => {
expect(isPresenceConnectionReady(true, 'channel', null)).toBe(false);
expect(isPresenceConnectionReady(true, 'channel', '')).toBe(false);
});
});
describe('hasAnyTerminalTransport', () => {
it('returns true when any slot is set', () => {
const sentinel = {};
@@ -144,17 +144,6 @@ export function assertDisconnectReason(
assert.fail(`${fieldName} must be one of 'user' | 'error' | 'server'`);
}
export function isPresenceConnectionReady(
connected: boolean,
channelId: string | null,
connectionId: string | null,
): boolean {
if (!connected) return false;
if (channelId === null || channelId.length === 0) return false;
if (connectionId === null || connectionId.length === 0) return false;
return true;
}
export interface TerminalTransportSnapshot {
readonly current: {readonly room?: unknown};
readonly hotSwap: {readonly pendingRoom?: unknown; readonly previousRoom?: unknown};
@@ -3,7 +3,6 @@
import assert from 'node:assert/strict';
import {isElectronPlatform} from '@app/features/platform/types/Platform';
import {Logger} from '@app/features/platform/utils/AppLogger';
import * as VoicePresenceHeartbeatCommands from '@app/features/voice/commands/VoicePresenceHeartbeatCommands';
import {Store} from '@app/features/voice/engine/Store';
import {sendVoiceStateDisconnect} from '@app/features/voice/engine/VoiceChannelConnector';
import {
@@ -35,7 +34,6 @@ import {
assertOptionalNonEmptyString,
assertVoiceServerUpdateShape,
hasAnyTerminalTransport,
isPresenceConnectionReady,
isReadyToRepublishTrack,
} from '@app/features/voice/engine/v2/VoiceEngineV2AppAdapterAssertions';
import {VoiceEngineV2AppReconnectPolicy} from '@app/features/voice/engine/v2/VoiceEngineV2AppReconnectPolicy';
@@ -60,7 +58,6 @@ import {timer} from 'rxjs';
const logger = new Logger('VoiceEngineV2AppConnectionHostAdapter');
const VOICE_SERVER_TIMEOUT_MS = 5000;
const VIDEO_DECODER_EXCLUSION_TIMEOUT_MS = 500;
const VOICE_PRESENCE_HEARTBEAT_INTERVAL_MS = 15000;
export interface VoiceServerUpdateData {
token: string;
@@ -133,19 +130,6 @@ async function getRoomVideoDecoderExclusions(): Promise<RoomOptions['subscriberV
}
}
function isSamePresenceConnection(
activeSub: Subscription | null,
activeConnection: {channelId: string; connectionId: string} | null,
channelId: string,
connectionId: string,
): boolean {
if (!activeSub) return false;
if (!activeConnection) return false;
if (activeConnection.channelId !== channelId) return false;
if (activeConnection.connectionId !== connectionId) return false;
return true;
}
function createWebAudioMixOption(): RoomOptions['webAudioMix'] {
const audioContext = getSharedVoiceAudioContext();
if (audioContext) {
@@ -210,8 +194,6 @@ export class VoiceEngineV2AppConnectionHostAdapter extends Store {
private reconnect = new VoiceEngineV2AppReconnectPolicy();
private voiceServerTimeoutSub: Subscription | null = null;
private hotSwapTimeoutSub: Subscription | null = null;
private voicePresenceHeartbeatSub: Subscription | null = null;
private voicePresenceHeartbeatConnection: {channelId: string; connectionId: string} | null = null;
private isLocalDisconnecting = false;
private hotSwapOperationQueue: Array<HotSwapQueuedOperation> = [];
@@ -338,71 +320,6 @@ export class VoiceEngineV2AppConnectionHostAdapter extends Store {
});
}
private startVoicePresenceHeartbeatForCurrentConnection(): void {
const {channelId, connectionId, connected} = this.connectionState;
if (!isPresenceConnectionReady(connected, channelId, connectionId)) {
this.stopVoicePresenceHeartbeat();
return;
}
const presenceChannelId = channelId as string;
const presenceConnectionId = connectionId as string;
if (
isSamePresenceConnection(
this.voicePresenceHeartbeatSub,
this.voicePresenceHeartbeatConnection,
presenceChannelId,
presenceConnectionId,
)
) {
return;
}
this.stopVoicePresenceHeartbeat();
this.voicePresenceHeartbeatConnection = {channelId: presenceChannelId, connectionId: presenceConnectionId};
this.voicePresenceHeartbeatSub = timer(0, VOICE_PRESENCE_HEARTBEAT_INTERVAL_MS).subscribe(() => {
void this.sendVoicePresenceHeartbeat(presenceChannelId, presenceConnectionId, {requireConnected: true});
});
}
private stopVoicePresenceHeartbeat(options: {markEnded?: boolean} = {}): void {
const connection = this.voicePresenceHeartbeatConnection;
if (this.voicePresenceHeartbeatSub) {
this.voicePresenceHeartbeatSub.unsubscribe();
this.voicePresenceHeartbeatSub = null;
}
this.voicePresenceHeartbeatConnection = null;
if (options.markEnded && connection) {
void this.markVoicePresenceHeartbeatEnded(connection);
}
}
private async sendVoicePresenceHeartbeat(
channelId: string,
connectionId: string,
options: {requireConnected: boolean},
): Promise<void> {
const current = this.connectionState;
if (
current.channelId !== channelId ||
current.connectionId !== connectionId ||
(options.requireConnected && !current.connected)
) {
return;
}
try {
await VoicePresenceHeartbeatCommands.heartbeat({channelId, connectionId});
} catch (error) {
logger.warn('Voice presence heartbeat failed', {channelId, connectionId, error});
}
}
private async markVoicePresenceHeartbeatEnded(connection: {channelId: string; connectionId: string}): Promise<void> {
try {
await VoicePresenceHeartbeatCommands.end(connection);
} catch (error) {
logger.warn('Voice presence heartbeat end failed', {...connection, error});
}
}
get lastConnectedChannel(): {
guildId: string;
channelId: string;
@@ -637,7 +554,6 @@ export class VoiceEngineV2AppConnectionHostAdapter extends Store {
this.update(() => {
this.transitionConnection({type: 'connection.failed', reason: 'error'});
});
this.stopVoicePresenceHeartbeat({markEnded: true});
this.throttle.setInFlightConnect(false);
this.reconnect.setReconnectState('error');
void onConnectFailed?.(guildId, resolvedChannelId, connectionId, attemptId, error ?? new Error(message));
@@ -675,7 +591,6 @@ export class VoiceEngineV2AppConnectionHostAdapter extends Store {
} catch (error) {
logger.warn('Failed to disconnect stale room', error);
}
this.stopVoicePresenceHeartbeat({markEnded: true});
return;
}
logger.info('Initializing voice connection');
@@ -691,7 +606,6 @@ export class VoiceEngineV2AppConnectionHostAdapter extends Store {
this.update(() => {
this.transitionConnection({type: 'connection.failed', reason: 'error'});
});
this.stopVoicePresenceHeartbeat({markEnded: true});
this.throttle.setInFlightConnect(false);
this.reconnect.setReconnectState('error');
void onConnectFailed?.(guildId, resolvedChannelId, connectionId, attemptId, error);
@@ -834,7 +748,6 @@ export class VoiceEngineV2AppConnectionHostAdapter extends Store {
this.transitionConnection({type: 'hotSwap.complete', room: newRoom, endpoint, connectionId});
});
this.clearHotSwapTimeout();
this.startVoicePresenceHeartbeatForCurrentConnection();
onHotSwapComplete?.(newRoom, attemptId, guildId, channelId);
await this.drainHotSwapQueue();
try {
@@ -960,7 +873,6 @@ export class VoiceEngineV2AppConnectionHostAdapter extends Store {
this.reconnect.setLastConnectedChannel(guildId, channelId);
this.throttle.setInFlightConnect(false);
this.reconnect.resetOnConnection();
this.startVoicePresenceHeartbeatForCurrentConnection();
assert.ok(this.connectionState.connected, 'markConnected post-condition: connection state reflects connected');
logger.info('Connection established');
}
@@ -973,7 +885,6 @@ export class VoiceEngineV2AppConnectionHostAdapter extends Store {
this.invalidateThrottleAttempt();
this.throttle.setInFlightConnect(false);
this.reconnect.setReconnectState(reason);
this.stopVoicePresenceHeartbeat({markEnded: true});
logger.info('Connection terminated', {reason});
}
@@ -991,7 +902,6 @@ export class VoiceEngineV2AppConnectionHostAdapter extends Store {
this.transitionConnection({type: 'connection.reconnected'});
});
this.reconnect.resetOnConnection();
this.startVoicePresenceHeartbeatForCurrentConnection();
logger.info('Connection reconnected');
}
@@ -1013,7 +923,6 @@ export class VoiceEngineV2AppConnectionHostAdapter extends Store {
});
this.invalidateThrottleAttempt();
this.reconnect.setReconnectState(reason);
this.stopVoicePresenceHeartbeat({markEnded: true});
this.update(() => {
this.isLocalDisconnecting = false;
});
@@ -1034,7 +943,6 @@ export class VoiceEngineV2AppConnectionHostAdapter extends Store {
this.transitionConnection({type: 'connection.disconnectForChannelMove'});
});
this.invalidateThrottleAttempt();
this.stopVoicePresenceHeartbeat({markEnded: true});
logger.info('Disconnected for channel move (preserving connectionId)');
}
@@ -1070,7 +978,6 @@ export class VoiceEngineV2AppConnectionHostAdapter extends Store {
this.disconnectRoomForTerminalUnload(previousRoom, 'previous-hot-swap');
}
this.disconnectRoomForTerminalUnload(room, 'current');
this.stopVoicePresenceHeartbeat({markEnded: true});
this.throttle.setInFlightConnect(false);
this.update(() => {
this.isLocalDisconnecting = false;
@@ -1188,7 +1095,6 @@ export class VoiceEngineV2AppConnectionHostAdapter extends Store {
this.transitionConnection({type: 'connection.reset'});
});
this.invalidateThrottleAttempt();
this.stopVoicePresenceHeartbeat({markEnded: true});
this.throttle.setInFlightConnect(false);
}
@@ -1206,7 +1112,6 @@ export class VoiceEngineV2AppConnectionHostAdapter extends Store {
this.transitionConnection({type: 'connection.abort'});
});
this.invalidateThrottleAttempt();
this.stopVoicePresenceHeartbeat({markEnded: true});
this.throttle.setInFlightConnect(false);
logger.info('Connection aborted due to gateway error');
}
@@ -0,0 +1,545 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import assert from 'node:assert/strict';
import {Logger} from '@app/features/platform/utils/AppLogger';
import {getElectronAPI} from '@app/features/ui/utils/NativeUtils';
import {
appendBrowserVoiceDebugEventSinkEntries,
openBrowserVoiceDebugEventSinkPopout,
} from '@app/features/voice/diagnostics/VoiceDebugBrowserEventSinkPopout';
import {asVoiceTrackSource, VoiceTrackSource} from '@app/features/voice/engine/VoiceTrackSource';
import type {DesktopVoiceDebugEventSinkEntry} from '@app/types/electron.d';
import type {
LocalTrackPublication,
Participant,
RemoteTrack,
RemoteTrackPublication,
Room,
TrackPublication,
} from 'livekit-client';
import {RoomEvent} from 'livekit-client';
import {assertNonNullObject, assertString} from './VoiceEngineV2AppAdapterAssertions';
const logger = new Logger('VoiceEngineV2AppDebugEventSinkHostAdapter');
export const VOICE_ENGINE_V2_APP_DEBUG_EVENT_SINK_MAX_ENTRIES = 1000;
export const VOICE_ENGINE_V2_APP_DEBUG_EVENT_SINK_MAX_LINE_CHARS = 262_144;
const SCREEN_SHARE_CODEC_NEGOTIATION_TOPIC = 'fluxer.rtc.codec-negotiation.v1';
const TEXT_DECODER = new TextDecoder();
interface VoiceDebugEventSinkEvent {
type: string;
timestamp_ns: string;
monotonic_ns?: string;
data?: Record<string, unknown>;
}
class BoundedRing<T> {
private readonly items: Array<T | null>;
private readonly capacity: number;
private readonly label: string;
private head = 0;
private tail = 0;
private count = 0;
constructor(capacity: number, label: string) {
assert.ok(capacity > 0, `${label} capacity must be positive`);
assertString(label, 'label');
this.capacity = capacity;
this.label = label;
this.items = new Array(capacity).fill(null);
}
get length(): number {
return this.count;
}
pushDropOldest(item: T): T | null {
let dropped: T | null = null;
if (this.count >= this.capacity) {
dropped = this.popFront();
}
this.pushBack(item);
assert.ok(this.count <= this.capacity, `${this.label} must stay bounded`);
return dropped;
}
toArray(): Array<T> {
const out: Array<T> = [];
for (let index = 0; index < this.count; index += 1) {
const slot = (this.head + index) % this.capacity;
const item = this.items[slot];
if (item === null) {
assert.fail(`${this.label} slot must exist`);
}
out.push(item);
}
return out;
}
private pushBack(item: T): void {
assert.ok(this.count < this.capacity, `${this.label} must have free capacity before push`);
this.items[this.tail] = item;
this.tail = (this.tail + 1) % this.capacity;
this.count += 1;
}
private popFront(): T | null {
if (this.count === 0) return null;
const item = this.items[this.head];
assert.notEqual(item, null, `${this.label} front slot must exist`);
this.items[this.head] = null;
this.head = (this.head + 1) % this.capacity;
this.count -= 1;
return item;
}
}
type VoiceDebugEventSinkRoomEventHandler = (...args: Array<never>) => void;
type VoiceDebugEventSinkRoomEventBinding = [RoomEvent, VoiceDebugEventSinkRoomEventHandler];
interface VoiceDebugEventSinkStartOptions {
guildId: string | null;
channelId: string;
connectionId: string | null;
room: Room;
}
interface RoomParticipantSummary {
identity: string;
sid: string;
isLocal: boolean;
name?: string;
metadata?: string;
attributes?: Record<string, string>;
connectionQuality?: string;
isSpeaking?: boolean;
permissions?: unknown;
trackPublications: Array<TrackPublicationSummary>;
}
interface TrackPublicationSummary {
trackSid: string;
trackName?: string;
source?: string;
kind?: string;
mimeType?: string;
isMuted?: boolean;
isSubscribed?: boolean;
isEnabled?: boolean;
dimensions?: {
width: number;
height: number;
};
}
function millisecondsToNanosecondsString(milliseconds: number): string {
if (!Number.isFinite(milliseconds) || milliseconds < 0) return '0';
const wholeMs = Math.trunc(milliseconds);
const fractionalNs = Math.round((milliseconds - wholeMs) * 1000000);
return (BigInt(wholeMs) * 1000000n + BigInt(fractionalNs)).toString();
}
function createDiagnosticEvent(type: string, data?: Record<string, unknown>): VoiceDebugEventSinkEvent {
const monotonicNow = typeof performance !== 'undefined' ? performance.now() : Date.now();
const timeOrigin =
typeof performance !== 'undefined' && Number.isFinite(performance.timeOrigin)
? performance.timeOrigin
: Date.now() - monotonicNow;
return {
type,
timestamp_ns: millisecondsToNanosecondsString(timeOrigin + monotonicNow),
monotonic_ns: millisecondsToNanosecondsString(monotonicNow),
...(data ? {data} : {}),
};
}
function errorToData(error: unknown): Record<string, unknown> {
if (error instanceof Error) {
return {
name: error.name,
message: error.message,
stack: error.stack,
};
}
return {message: String(error)};
}
function truncateEventSinkLine(line: string): string {
assertString(line, 'event sink line');
if (line.length <= VOICE_ENGINE_V2_APP_DEBUG_EVENT_SINK_MAX_LINE_CHARS) return line;
const omittedChars = line.length - VOICE_ENGINE_V2_APP_DEBUG_EVENT_SINK_MAX_LINE_CHARS;
return `${line.slice(0, VOICE_ENGINE_V2_APP_DEBUG_EVENT_SINK_MAX_LINE_CHARS)}... [truncated ${omittedChars} chars]`;
}
function stringifyEventSinkEntry(sequence: number, event: VoiceDebugEventSinkEvent): string {
assert.ok(Number.isSafeInteger(sequence), 'event sink sequence must be a safe integer');
assert.ok(sequence >= 1, 'event sink sequence must be >= 1');
try {
return truncateEventSinkLine(JSON.stringify({sequence, ...event}));
} catch (error) {
return truncateEventSinkLine(
JSON.stringify({
sequence,
type: event.type,
timestamp_ns: event.timestamp_ns,
monotonic_ns: event.monotonic_ns,
stringifyError: errorToData(error),
}),
);
}
}
function createEventSinkEntry(sequence: number, event: VoiceDebugEventSinkEvent): DesktopVoiceDebugEventSinkEntry {
assert.ok(event !== null && typeof event === 'object', 'event sink event must be an object');
assertString(event.type, 'event sink event type');
assertString(event.timestamp_ns, 'event sink event timestamp');
return {
sequence,
line: stringifyEventSinkEntry(sequence, event),
};
}
function getTrackDimensions(publication: TrackPublication): TrackPublicationSummary['dimensions'] | undefined {
const track = publication.track;
if (!track || !('dimensions' in track)) return undefined;
const dimensions = (track as {dimensions?: {width?: number; height?: number}}).dimensions;
if (typeof dimensions?.width !== 'number' || typeof dimensions.height !== 'number') return undefined;
return {
width: dimensions.width,
height: dimensions.height,
};
}
function summarizePublication(publication: TrackPublication): TrackPublicationSummary {
return {
trackSid: publication.trackSid,
trackName: publication.trackName,
source: publication.source,
kind: publication.kind,
mimeType: publication.mimeType,
isMuted: publication.isMuted,
isSubscribed: 'isSubscribed' in publication ? Boolean(publication.isSubscribed) : undefined,
isEnabled: 'isEnabled' in publication ? Boolean(publication.isEnabled) : undefined,
dimensions: getTrackDimensions(publication),
};
}
function summarizeParticipant(participant: Participant | undefined): RoomParticipantSummary | null {
if (!participant) return null;
const trackPublications: Array<TrackPublicationSummary> = [];
participant.trackPublications.forEach((publication) => {
trackPublications.push(summarizePublication(publication));
});
return {
identity: participant.identity,
sid: participant.sid,
isLocal: participant.isLocal,
name: participant.name,
metadata: participant.metadata,
attributes: participant.attributes,
connectionQuality: participant.connectionQuality,
isSpeaking: participant.isSpeaking,
permissions: participant.permissions,
trackPublications,
};
}
function summarizeRoom(room: Room | null): Record<string, unknown> | null {
if (!room) return null;
return {
name: room.name,
state: room.state,
numParticipants: room.numParticipants,
localParticipant: summarizeParticipant(room.localParticipant),
remoteParticipants: Array.from(room.remoteParticipants.values()).map((participant) =>
summarizeParticipant(participant),
),
};
}
function summarizeTrackEvent(
publication: TrackPublication | RemoteTrackPublication | LocalTrackPublication,
participant: Participant | undefined,
): Record<string, unknown> {
return {
publication: summarizePublication(publication as TrackPublication),
participant: summarizeParticipant(participant),
isScreenShare: asVoiceTrackSource(publication.source) === VoiceTrackSource.ScreenShare,
isScreenShareAudio: asVoiceTrackSource(publication.source) === VoiceTrackSource.ScreenShareAudio,
};
}
function summarizeRemoteTrack(track: RemoteTrack): Record<string, unknown> {
return {
sid: track.sid,
kind: track.kind,
source: track.source,
mediaStreamTrackId: track.mediaStreamTrack?.id,
readyState: track.mediaStreamTrack?.readyState,
muted: track.mediaStreamTrack?.muted,
};
}
function parseAllowedDataMessage(payload: Uint8Array, topic: string | undefined): Record<string, unknown> {
if (topic !== SCREEN_SHARE_CODEC_NEGOTIATION_TOPIC) {
return {
topic: topic ?? null,
payloadBytes: payload.byteLength,
decoded: null,
};
}
try {
return {
topic,
payloadBytes: payload.byteLength,
decoded: JSON.parse(TEXT_DECODER.decode(payload)) as unknown,
};
} catch (error) {
return {
topic,
payloadBytes: payload.byteLength,
decodeError: errorToData(error),
};
}
}
export class VoiceEngineV2AppDebugEventSinkHostAdapter {
private channelId: string | null = null;
private connectionId: string | null = null;
private room: Room | null = null;
private roomDisposer: (() => void) | null = null;
private readonly eventSinkEntries = new BoundedRing<DesktopVoiceDebugEventSinkEntry>(
VOICE_ENGINE_V2_APP_DEBUG_EVENT_SINK_MAX_ENTRIES,
'voice debug event sink history',
);
private eventSinkSequence = 0;
private eventSinkForwardFailureCount = 0;
getEventSinkEntries(): Array<DesktopVoiceDebugEventSinkEntry> {
assert.ok(
this.eventSinkEntries.length <= VOICE_ENGINE_V2_APP_DEBUG_EVENT_SINK_MAX_ENTRIES,
'event sink history must stay bounded',
);
return this.eventSinkEntries.toArray();
}
async openEventSinkPopout(): Promise<void> {
const electron = getElectronAPI();
const entries = this.getEventSinkEntries();
if (electron?.openVoiceDebugEventSinkPopout) {
try {
await electron.openVoiceDebugEventSinkPopout(entries);
return;
} catch (error) {
logger.warn('Failed to open voice debug event sink desktop popout', {error});
}
}
try {
const opened = await openBrowserVoiceDebugEventSinkPopout(entries);
if (!opened) {
logger.warn('Failed to open voice debug event sink browser popout');
}
} catch (error) {
logger.warn('Failed to open voice debug event sink browser popout', {error});
}
}
private isStartIdempotent(options: VoiceDebugEventSinkStartOptions): boolean {
if (this.channelId !== options.channelId) return false;
if (this.connectionId !== options.connectionId) return false;
return this.room === options.room;
}
start(options: VoiceDebugEventSinkStartOptions): void {
assertNonNullObject(options, 'options');
assertString(options.channelId, 'options.channelId');
assert.ok(options.channelId.length > 0, 'options.channelId must not be empty');
assertNonNullObject(options.room, 'options.room');
if (this.isStartIdempotent(options)) return;
this.stop('replaced');
this.channelId = options.channelId;
this.connectionId = options.connectionId;
this.room = options.room;
this.bindRoom(options.room);
this.record('voice.debug_event_sink.tracking_started', {
guildId: options.guildId,
channelId: options.channelId,
connectionId: options.connectionId,
room: summarizeRoom(options.room),
});
}
stop(reason = 'stopped'): void {
assertString(reason, 'reason');
if (this.channelId !== null) {
this.record('voice.debug_event_sink.tracking_stopped', {
reason,
channelId: this.channelId,
connectionId: this.connectionId,
room: summarizeRoom(this.room),
});
}
this.roomDisposer?.();
this.roomDisposer = null;
this.channelId = null;
this.connectionId = null;
this.room = null;
}
private buildRoomLifecycleEventBindings(room: Room): Array<VoiceDebugEventSinkRoomEventBinding> {
return [
[RoomEvent.Connected, () => this.record('livekit.room.connected', {room: summarizeRoom(room)})],
[
RoomEvent.Disconnected,
(reason?: unknown) =>
this.record('livekit.room.disconnected', {reason: String(reason ?? 'unknown'), room: summarizeRoom(room)}),
],
[RoomEvent.Reconnecting, () => this.record('livekit.room.reconnecting', {room: summarizeRoom(room)})],
[RoomEvent.Reconnected, () => this.record('livekit.room.reconnected', {room: summarizeRoom(room)})],
];
}
private buildParticipantEventBindings(): Array<VoiceDebugEventSinkRoomEventBinding> {
return [
[
RoomEvent.ParticipantConnected,
(participant: Participant) =>
this.record('livekit.participant.connected', {participant: summarizeParticipant(participant)}),
],
[
RoomEvent.ParticipantDisconnected,
(participant: Participant) =>
this.record('livekit.participant.disconnected', {participant: summarizeParticipant(participant)}),
],
[
RoomEvent.ActiveSpeakersChanged,
(speakers: Array<Participant>) =>
this.record('livekit.active_speakers.changed', {
speakers: speakers.map((participant) => summarizeParticipant(participant)),
}),
],
];
}
private buildTrackEventBindings(): Array<VoiceDebugEventSinkRoomEventBinding> {
return [
[
RoomEvent.TrackPublished,
(publication: RemoteTrackPublication, participant: Participant) =>
this.record('livekit.track.published', summarizeTrackEvent(publication, participant)),
],
[
RoomEvent.TrackUnpublished,
(publication: RemoteTrackPublication, participant: Participant) =>
this.record('livekit.track.unpublished', summarizeTrackEvent(publication, participant)),
],
[
RoomEvent.TrackSubscribed,
(track: RemoteTrack, publication: RemoteTrackPublication, participant: Participant) =>
this.record('livekit.track.subscribed', {
...summarizeTrackEvent(publication, participant),
track: summarizeRemoteTrack(track),
}),
],
[
RoomEvent.TrackUnsubscribed,
(track: RemoteTrack, publication: RemoteTrackPublication, participant: Participant) =>
this.record('livekit.track.unsubscribed', {
...summarizeTrackEvent(publication, participant),
track: summarizeRemoteTrack(track),
}),
],
[
RoomEvent.TrackMuted,
(publication: TrackPublication, participant: Participant) =>
this.record('livekit.track.muted', summarizeTrackEvent(publication, participant)),
],
[
RoomEvent.TrackUnmuted,
(publication: TrackPublication, participant: Participant) =>
this.record('livekit.track.unmuted', summarizeTrackEvent(publication, participant)),
],
[
RoomEvent.LocalTrackPublished,
(publication: LocalTrackPublication, participant: Participant) =>
this.record('livekit.local_track.published', summarizeTrackEvent(publication, participant)),
],
[
RoomEvent.LocalTrackUnpublished,
(publication: LocalTrackPublication, participant: Participant) =>
this.record('livekit.local_track.unpublished', summarizeTrackEvent(publication, participant)),
],
];
}
private buildDataEventBindings(): Array<VoiceDebugEventSinkRoomEventBinding> {
return [
[
RoomEvent.DataReceived,
(payload: Uint8Array, participant: Participant | undefined, kind: unknown, topic?: string) =>
this.record('livekit.data.received', {
participant: summarizeParticipant(participant),
kind: String(kind),
...parseAllowedDataMessage(payload, topic),
}),
],
];
}
private buildRoomEventBindings(room: Room): Array<VoiceDebugEventSinkRoomEventBinding> {
return [
...this.buildRoomLifecycleEventBindings(room),
...this.buildParticipantEventBindings(),
...this.buildTrackEventBindings(),
...this.buildDataEventBindings(),
];
}
private bindRoom(room: Room): void {
assertNonNullObject(room, 'room');
this.roomDisposer?.();
const bindings = this.buildRoomEventBindings(room);
assert.ok(bindings.length > 0, 'expected at least one room event binding');
for (const [event, handler] of bindings) {
room.on(event, handler);
}
this.roomDisposer = () => {
for (const [event, handler] of bindings) {
room.off(event, handler);
}
};
}
private record(type: string, data?: Record<string, unknown>): void {
const event = createDiagnosticEvent(type, data);
this.appendEventSinkEntry(event);
}
private appendEventSinkEntry(event: VoiceDebugEventSinkEvent): void {
this.eventSinkSequence += 1;
assert.ok(Number.isSafeInteger(this.eventSinkSequence), 'event sink sequence must stay safe');
const entry = createEventSinkEntry(this.eventSinkSequence, event);
this.eventSinkEntries.pushDropOldest(entry);
this.forwardEventSinkEntries([entry]);
}
private forwardEventSinkEntries(entries: Array<DesktopVoiceDebugEventSinkEntry>): void {
assert.ok(entries.length >= 1, 'event sink forward requires at least one entry');
const electron = getElectronAPI();
if (electron?.appendVoiceDebugEventSinkEntries) {
try {
electron.appendVoiceDebugEventSinkEntries(entries);
this.eventSinkForwardFailureCount = 0;
} catch (error) {
this.eventSinkForwardFailureCount += 1;
if (this.eventSinkForwardFailureCount === 1) {
logger.warn('Failed to forward voice debug event sink entries to desktop popout', {error});
}
}
}
appendBrowserVoiceDebugEventSinkEntries(entries);
}
}
export default new VoiceEngineV2AppDebugEventSinkHostAdapter();
@@ -3,7 +3,6 @@
import assert from 'node:assert/strict';
import {Platform} from '@app/features/platform/types/Platform';
import AdaptiveScreenShareEngine from '@app/features/voice/engine/AdaptiveScreenShareEngine';
import {markScreenShareCaptureEnded} from '@app/features/voice/engine/ScreenShareCaptureDiagnostics';
import type {NegotiationReason} from '@app/features/voice/engine/ScreenShareCodecNegotiation';
import ScreenSharePublicationMigration from '@app/features/voice/engine/ScreenSharePublicationMigration';
import {updateLocalParticipantFromRoom} from '@app/features/voice/engine/VoiceMediaEngineBridge';
@@ -408,7 +407,6 @@ export class VoiceEngineV2AppScreenShareCodecMigration {
this.adapter.ensureScreenShareKeepAliveSinkInternal(ctx.participant);
} else {
await this.adapter.cleanupLingeringScreenShareTracks(ctx.participant);
markScreenShareCaptureEnded('screen-share-codec-migration-failed');
}
this.adapter.transitionScreenShareLifecycleInternal({
type: 'share.reject',
@@ -3,10 +3,6 @@
import assert from 'node:assert/strict';
import {isDesktop, isNativeMacOS} from '@app/features/ui/utils/NativeUtils';
import AdaptiveScreenShareEngine from '@app/features/voice/engine/AdaptiveScreenShareEngine';
import {
markScreenShareCaptureActive,
markScreenShareCaptureEnded,
} from '@app/features/voice/engine/ScreenShareCaptureDiagnostics';
import {updateLocalParticipantFromRoom} from '@app/features/voice/engine/VoiceMediaEngineBridge';
import {
enforceLocalMediaPublicationCap,
@@ -150,7 +146,6 @@ export class VoiceEngineV2AppScreenShareLiveKitFlows {
? null
: async () => {
await this.adapter.cleanupLingeringScreenShareTracks(participant, stopCleanupSnapshot ?? undefined);
markScreenShareCaptureEnded('screen-share-disabled');
},
updateLocalParticipant: true,
audioSync: {kind: 'participant-after-watch'},
@@ -179,7 +174,6 @@ export class VoiceEngineV2AppScreenShareLiveKitFlows {
): Promise<void> {
assert.ok(participant);
const cancelled = isUserCancelledScreenShareError(error);
const endedReason = cancelled ? 'screen-share-cancelled' : 'screen-share-failed';
if (cancelled) {
logger.debug('User cancelled or permission denied', {name: (error as Error).name});
} else {
@@ -192,7 +186,6 @@ export class VoiceEngineV2AppScreenShareLiveKitFlows {
}
if (!actual) {
await this.adapter.cleanupLingeringScreenShareTracks(participant, stopCleanupSnapshot ?? undefined);
markScreenShareCaptureEnded(endedReason);
}
settleScreenShareFailure({
adapter: this.adapter,
@@ -356,8 +349,7 @@ export class VoiceEngineV2AppScreenShareLiveKitFlows {
playSound: boolean,
): Promise<void> {
assert.ok(participant);
const {videoDeviceId, audioDeviceId} = options || {};
markScreenShareCaptureActive({method: 'device-media', device: {videoDeviceId, audioDeviceId}});
const {videoDeviceId} = options || {};
await runScreenShareActivationRitual({
adapter: this.adapter,
room,
@@ -409,7 +401,7 @@ export class VoiceEngineV2AppScreenShareLiveKitFlows {
participant,
actual: participant.isScreenShareEnabled,
applyState,
onInactiveAfterSync: () => markScreenShareCaptureEnded('device-screen-share-failed'),
onInactiveAfterSync: null,
monitorEndOnActive: false,
playSound: false,
buildTransition: (actualNow) =>
@@ -1,6 +1,5 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {markScreenShareCaptureActive} from '@app/features/voice/engine/ScreenShareCaptureDiagnostics';
import {
type CapturedScreenShareTracks,
type DeviceScreenShareCaptureOptions,
@@ -134,13 +133,6 @@ export async function createDeviceReplacementTracks(
throw new Error('No video track found in device screen share capture');
}
const audioTrack = stream.getAudioTracks()[0];
markScreenShareCaptureActive({
method: 'device-media',
device: {
videoDeviceId: options?.videoDeviceId,
audioDeviceId: options?.audioDeviceId,
},
});
stopUnselectedStreamTracks(stream, [videoTrack, audioTrack]);
return {
videoTrack,
@@ -1,12 +1,10 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {updateScreenShareDisplayMediaSettings} from '@app/features/voice/engine/ScreenShareCaptureDiagnostics';
import {
type CapturedScreenShareTracks,
stopMediaTrack,
stopUnselectedStreamTracks,
} from '@app/features/voice/engine/voice_screen_share_manager/shared';
import ActiveScreenShareSource from '@app/features/voice/state/ActiveScreenShareSource';
import type {ScreenShareCaptureOptions} from 'livekit-client';
type DisplayMediaVideoConstraints = MediaTrackConstraints & {
@@ -108,9 +106,6 @@ export async function createDisplayScreenShareTracks(
videoTrack.contentHint = options.contentHint;
}
await videoTrack.applyConstraints({colorSpace: 'rec709'} as MediaTrackConstraints).catch(() => undefined);
updateScreenShareDisplayMediaSettings(videoTrack, {
sourceId: ActiveScreenShareSource.getSourceId(),
});
const cursor = resolveCapturedDisplayMediaCursorCapture(videoTrack, options);
if ((videoTrack.getSettings() as DisplayMediaTrackSettings).cursor !== cursor) {
await videoTrack.applyConstraints({cursor} as MediaTrackConstraints).catch(() => undefined);
-1
View File
@@ -111,7 +111,6 @@ export const AdminACLs = {
VOICE_REGION_DELETE: 'voice:region:delete',
VOICE_REGION_LIST: 'voice:region:list',
VOICE_REGION_UPDATE: 'voice:region:update',
VOICE_DIAGNOSTICS_VIEW: 'voice:diagnostics:view',
VOICE_SERVER_CREATE: 'voice:server:create',
VOICE_SERVER_DELETE: 'voice:server:delete',
VOICE_SERVER_LIST: 'voice:server:list',
@@ -1,11 +1,6 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {
coerceNumberFromString,
createStringType,
SnowflakeStringType,
SnowflakeType,
} from '@fluxer/schema/src/primitives/SchemaPrimitives';
import {createStringType, SnowflakeStringType, SnowflakeType} from '@fluxer/schema/src/primitives/SchemaPrimitives';
import {z} from 'zod';
function areServerCoordinatesPaired(
@@ -266,64 +261,6 @@ export const GetVoiceServerResponse = z.object({
export type GetVoiceServerResponse = z.infer<typeof GetVoiceServerResponse>;
const DiagnosticsTimestampMs = coerceNumberFromString(z.number().int().min(0).max(8640000000000000)).describe(
'Unix timestamp in milliseconds',
);
const DiagnosticsSessionId = z
.string()
.min(1)
.max(128)
.regex(/^[A-Za-z0-9_.:-]+$/)
.optional()
.describe('Optional voice diagnostics session id filter');
export const VoiceDiagnosticsQueryRequest = z
.object({
channel_id: SnowflakeStringType.describe('Channel id to query diagnostics for'),
start_ms: DiagnosticsTimestampMs.describe('Inclusive start timestamp in milliseconds'),
end_ms: DiagnosticsTimestampMs.describe('Inclusive end timestamp in milliseconds'),
session_id: DiagnosticsSessionId,
limit_objects: coerceNumberFromString(z.number().int().min(1).max(5000))
.optional()
.default(1000)
.describe('Maximum matching S3 objects to include'),
})
.superRefine((value, ctx) => {
if (value.end_ms < value.start_ms) {
ctx.addIssue({
code: 'custom',
path: ['end_ms'],
message: 'INVALID_FORMAT',
});
}
const maxRangeMs = 31 * 24 * 60 * 60 * 1000;
if (value.end_ms - value.start_ms > maxRangeMs) {
ctx.addIssue({
code: 'custom',
path: ['end_ms'],
message: 'INVALID_FORMAT',
});
}
});
export type VoiceDiagnosticsQueryRequest = z.infer<typeof VoiceDiagnosticsQueryRequest>;
const VoiceDiagnosticsObjectResponse = z.object({
key: z.string().describe('S3 object key'),
session_id: z.string().describe('Voice diagnostics session id'),
start_ns: z.string().describe('First client event timestamp in the object, in nanoseconds'),
end_ns: z.string().describe('Last client event timestamp in the object, in nanoseconds'),
last_modified: z.iso.datetime().nullable().describe('S3 object last-modified timestamp'),
});
export const VoiceDiagnosticsObjectListResponse = z.object({
bucket: z.string().describe('Diagnostics S3 bucket'),
objects: z.array(VoiceDiagnosticsObjectResponse).describe('Matching diagnostics objects'),
});
export type VoiceDiagnosticsObjectListResponse = z.infer<typeof VoiceDiagnosticsObjectListResponse>;
export const CreateVoiceServerResponse = z.object({
server: VoiceServerAdminResponse.describe('Created voice server'),
});
@@ -24,7 +24,6 @@ import {createBase64StringType} from '@fluxer/schema/src/primitives/FileValidato
import {ContentWarningLevelSchema} from '@fluxer/schema/src/primitives/GuildValidators';
import {QueryBooleanType} from '@fluxer/schema/src/primitives/QueryValidators';
import {
coerceNumberFromString,
createNamedLiteral,
createNamedLiteralUnion,
createStringType,
@@ -292,46 +291,6 @@ export const CallRingBodySchema = z.object({
export type CallRingBodySchema = z.infer<typeof CallRingBodySchema>;
export const VoiceDebugLoggingToggleBodySchema = z.object({
enabled: z.boolean().describe('Whether voice debug logging should be active for this channel'),
duration_ms: coerceNumberFromString(z.number().int().min(60000).max(14400000))
.optional()
.describe('Optional activation duration in milliseconds. Defaults to one hour and is capped at four hours.'),
});
export type VoiceDebugLoggingToggleBodySchema = z.infer<typeof VoiceDebugLoggingToggleBodySchema>;
const VoiceDebugLoggingTimestampNs = z
.string()
.regex(/^[0-9]{1,32}$/)
.describe('Nanosecond timestamp encoded as an unsigned decimal string');
export const VoiceDebugLoggingEventSchema = z
.object({
type: createStringType(1, 128).describe('Client-side diagnostic event type'),
timestamp_ns: VoiceDebugLoggingTimestampNs.describe('Client wall-clock Unix timestamp in nanoseconds'),
monotonic_ns: VoiceDebugLoggingTimestampNs.optional().describe('Client monotonic timestamp in nanoseconds'),
data: z.record(z.string(), z.unknown()).optional().describe('Event-specific diagnostic payload'),
})
.passthrough();
export type VoiceDebugLoggingEventSchema = z.infer<typeof VoiceDebugLoggingEventSchema>;
export const VoiceDebugLoggingEventsBodySchema = z.object({
session_id: createStringType(1, 128).describe('Active voice debug logging session id'),
connection_id: createStringType(1, 128).optional().describe('Client voice connection id'),
participant_identity: createStringType(1, 256).optional().describe('LiveKit participant identity'),
events: z.array(VoiceDebugLoggingEventSchema).min(1).max(200).describe('NDJSON batch events to store'),
});
export type VoiceDebugLoggingEventsBodySchema = z.infer<typeof VoiceDebugLoggingEventsBodySchema>;
export const VoicePresenceHeartbeatBodySchema = z.object({
connection_id: createStringType(1, 128).describe('Client voice connection id'),
});
export type VoicePresenceHeartbeatBodySchema = z.infer<typeof VoicePresenceHeartbeatBodySchema>;
export const StreamUpdateBodySchema = z.object({
region: createStringType(RTC_REGION_ID_MIN_LENGTH, RTC_REGION_ID_MAX_LENGTH)
.optional()
@@ -47,46 +47,6 @@ export const CallEligibilityResponse = z.object({
export type CallEligibilityResponse = z.infer<typeof CallEligibilityResponse>;
export const VoiceDebugLoggingStatusResponse = z.object({
active: z.boolean().describe('Whether clients in this channel should currently send voice diagnostics'),
session_id: z.string().nullable().describe('Current debug logging session id, if active'),
activated_by_user_id: SnowflakeStringType.nullable().describe('Staff user that activated the session, if active'),
started_at_ms: z.number().int().nonnegative().nullable().describe('Session start Unix timestamp in milliseconds'),
expires_at_ms: z
.number()
.int()
.nonnegative()
.nullable()
.describe('Session expiration Unix timestamp in milliseconds'),
poll_interval_ms: Int32Type.describe('Recommended client polling interval in milliseconds'),
upload_interval_ms: Int32Type.describe('Recommended client telemetry batch upload interval in milliseconds'),
});
export type VoiceDebugLoggingStatusResponse = z.infer<typeof VoiceDebugLoggingStatusResponse>;
export const VoiceDebugLoggingEventsResponse = z.object({
accepted: z.boolean().describe('Whether the telemetry batch was accepted for storage'),
active: z.boolean().describe('Whether the server still considers this logging session active'),
stored_event_count: Int32Type.describe('Number of events written to diagnostics storage'),
});
export type VoiceDebugLoggingEventsResponse = z.infer<typeof VoiceDebugLoggingEventsResponse>;
export const VoicePresenceHeartbeatResponse = z.object({
ok: z.boolean().describe('Whether the heartbeat was accepted'),
heartbeat_interval_ms: Int32Type.describe('Recommended client heartbeat interval in milliseconds'),
heartbeat_ttl_ms: Int32Type.describe('Server-side heartbeat expiration window in milliseconds'),
expires_at_ms: z.number().int().nonnegative().describe('Unix timestamp in milliseconds when this heartbeat expires'),
});
export type VoicePresenceHeartbeatResponse = z.infer<typeof VoicePresenceHeartbeatResponse>;
export const VoicePresenceHeartbeatEndResponse = z.object({
ok: z.boolean().describe('Whether the heartbeat was ended'),
});
export type VoicePresenceHeartbeatEndResponse = z.infer<typeof VoicePresenceHeartbeatEndResponse>;
export const ChannelResponse = z.object({
id: SnowflakeStringType.describe('The unique identifier (snowflake) for this channel'),
guild_id: SnowflakeStringType.optional().describe('The ID of the guild this channel belongs to'),