mirror of
https://github.com/fluxerapp/fluxer
synced 2026-10-07 19:22:14 +09:00
feat(voice): add the recon service and remove the old worker (#2719)
This commit is contained in:
@@ -531,15 +531,6 @@ export function buildAPIConfigFromMaster(master: MasterConfig): APIConfig {
|
||||
laneName: apiWorkerConfig?.lane,
|
||||
taskName: apiWorkerConfig?.task as WorkerTaskName | undefined,
|
||||
enableCronScheduler: apiWorkerConfig?.enable_cron_scheduler,
|
||||
enableVoiceReconciliation: apiWorkerConfig?.enable_voice_reconciliation ?? true,
|
||||
voiceReconciliation: {
|
||||
intervalMs: apiWorkerConfig?.voice_reconciliation?.interval_ms,
|
||||
staggerDelayMs: apiWorkerConfig?.voice_reconciliation?.stagger_delay_ms,
|
||||
lockTtlSeconds: apiWorkerConfig?.voice_reconciliation?.lock_ttl_seconds,
|
||||
cadenceTtlSeconds: apiWorkerConfig?.voice_reconciliation?.cadence_ttl_seconds,
|
||||
gatewayOnlyGraceMs: apiWorkerConfig?.voice_reconciliation?.gateway_only_grace_ms,
|
||||
liveKitOnlyGraceMs: apiWorkerConfig?.voice_reconciliation?.livekit_only_grace_ms,
|
||||
},
|
||||
laneConcurrencyOverrides: {
|
||||
realtime: apiWorkerConfig?.lane_concurrency_overrides?.realtime,
|
||||
unfurl: apiWorkerConfig?.lane_concurrency_overrides?.unfurl,
|
||||
|
||||
@@ -377,15 +377,6 @@ export interface APIConfig {
|
||||
laneName?: APIWorkerLaneName;
|
||||
taskName?: WorkerTaskName;
|
||||
enableCronScheduler?: boolean;
|
||||
enableVoiceReconciliation: boolean;
|
||||
voiceReconciliation: {
|
||||
intervalMs: number | undefined;
|
||||
staggerDelayMs: number | undefined;
|
||||
lockTtlSeconds: number | undefined;
|
||||
cadenceTtlSeconds: number | undefined;
|
||||
gatewayOnlyGraceMs: number | undefined;
|
||||
liveKitOnlyGraceMs: number | undefined;
|
||||
};
|
||||
laneConcurrencyOverrides: {
|
||||
realtime?: number;
|
||||
unfurl?: number;
|
||||
|
||||
@@ -134,28 +134,3 @@ export function parseParticipantMetadataWithRaw(metadata: string): {
|
||||
export function isDMRoom(context: VoiceRoomContext): context is DMRoomContext {
|
||||
return context.type === 'dm';
|
||||
}
|
||||
|
||||
const PARTICIPANT_IDENTITY_PREFIX = 'user_';
|
||||
|
||||
interface ParticipantIdentity {
|
||||
readonly userId: UserID;
|
||||
readonly connectionId: string;
|
||||
}
|
||||
|
||||
export function parseParticipantIdentity(identity: string): ParticipantIdentity | null {
|
||||
if (!identity.startsWith(PARTICIPANT_IDENTITY_PREFIX)) {
|
||||
return null;
|
||||
}
|
||||
const parts = identity.split('_');
|
||||
if (parts.length !== 3 || parts[0] !== 'user') {
|
||||
return null;
|
||||
}
|
||||
try {
|
||||
return {
|
||||
userId: createUserID(BigInt(parts[1])),
|
||||
connectionId: parts[2],
|
||||
};
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,40 +0,0 @@
|
||||
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
|
||||
import {describe, expect, it} from 'vitest';
|
||||
import {candidateTtlSecondsFor} from './VoiceReconciliationWorker';
|
||||
|
||||
const INTERVAL_MS = 15000;
|
||||
const GATEWAY_ONLY_GRACE_MS = 10000;
|
||||
|
||||
function ttlFor(observedSweepSpacingMs: number): number {
|
||||
return candidateTtlSecondsFor({
|
||||
intervalMs: INTERVAL_MS,
|
||||
observedSweepSpacingMs,
|
||||
graceMs: GATEWAY_ONLY_GRACE_MS,
|
||||
});
|
||||
}
|
||||
|
||||
describe('candidateTtlSecondsFor', () => {
|
||||
it('outlives the gap between two consecutive observations of the same key', () => {
|
||||
for (const observedSweepSpacingMs of [0, 45_000, 136_000, 300_000, 596_000, 900_000]) {
|
||||
expect(ttlFor(observedSweepSpacingMs) * 1000).toBeGreaterThan(observedSweepSpacingMs);
|
||||
}
|
||||
});
|
||||
|
||||
it('outlives a sweep gap far longer than the tick interval', () => {
|
||||
expect(ttlFor(596_000) * 1000).toBeGreaterThan(596_000);
|
||||
});
|
||||
|
||||
it('grows with the observed sweep spacing rather than the tick interval', () => {
|
||||
expect(ttlFor(596_000)).toBeGreaterThan(ttlFor(136_000));
|
||||
expect(ttlFor(136_000)).toBeGreaterThan(ttlFor(0));
|
||||
});
|
||||
|
||||
it('keeps a floor that survives a single long sweep before any spacing is observed', () => {
|
||||
expect(ttlFor(0)).toBeGreaterThanOrEqual(300);
|
||||
});
|
||||
|
||||
it('stays bounded so a stale candidate cannot outlive its connection indefinitely', () => {
|
||||
expect(ttlFor(Number.MAX_SAFE_INTEGER)).toBeLessThanOrEqual(3600);
|
||||
});
|
||||
});
|
||||
File diff suppressed because it is too large
Load Diff
@@ -1,64 +0,0 @@
|
||||
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
|
||||
import type {IKVProvider} from '@pkgs/kv_client/src/IKVProvider';
|
||||
import {describe, expect, it, vi} from 'vitest';
|
||||
import type {ILogger} from '../../ILogger';
|
||||
import type {IGatewayService} from '../../infrastructure/IGatewayService';
|
||||
import type {ILiveKitService} from '../../infrastructure/ILiveKitService';
|
||||
import type {IVoiceRoomStore} from '../../infrastructure/IVoiceRoomStore';
|
||||
import {VoiceReconciliationWorker} from '../VoiceReconciliationWorker';
|
||||
|
||||
function createLogger(): ILogger {
|
||||
const logger = {
|
||||
trace: vi.fn(),
|
||||
debug: vi.fn(),
|
||||
info: vi.fn(),
|
||||
warn: vi.fn(),
|
||||
error: vi.fn(),
|
||||
child: () => logger,
|
||||
};
|
||||
return logger as unknown as ILogger;
|
||||
}
|
||||
|
||||
function createHarness() {
|
||||
const releaseLock = vi.fn().mockResolvedValue(true);
|
||||
const kvClient = {
|
||||
acquireLock: vi.fn().mockResolvedValue(true),
|
||||
extendLock: vi.fn().mockResolvedValue(true),
|
||||
releaseLock,
|
||||
setnx: vi.fn().mockResolvedValue(true),
|
||||
get: vi.fn().mockResolvedValue(null),
|
||||
setex: vi.fn().mockResolvedValue(undefined),
|
||||
} as unknown as IKVProvider;
|
||||
let finishDiscovery: () => void = () => {};
|
||||
const discovery = new Promise<{rooms: []}>((resolve) => {
|
||||
finishDiscovery = () => resolve({rooms: []});
|
||||
});
|
||||
const getActiveVoiceRooms = vi.fn().mockReturnValue(discovery);
|
||||
const worker = new VoiceReconciliationWorker({
|
||||
gatewayService: {getActiveVoiceRooms} as unknown as IGatewayService,
|
||||
liveKitService: {
|
||||
listActiveRooms: async () => ({rooms: [], errors: [], completed: true, searchedServers: 0}),
|
||||
} as unknown as ILiveKitService,
|
||||
voiceRoomStore: {listPinnedRooms: async () => []} as unknown as IVoiceRoomStore,
|
||||
kvClient,
|
||||
logger: createLogger(),
|
||||
intervalMs: 60000,
|
||||
staggerDelayMs: 0,
|
||||
});
|
||||
return {worker, releaseLock, getActiveVoiceRooms, finishDiscovery: () => finishDiscovery()};
|
||||
}
|
||||
|
||||
describe('VoiceReconciliationWorker stop', () => {
|
||||
it('releases the reconciliation lock before stop resolves', async () => {
|
||||
const {worker, releaseLock, getActiveVoiceRooms, finishDiscovery} = createHarness();
|
||||
|
||||
worker.start();
|
||||
await vi.waitFor(() => expect(getActiveVoiceRooms).toHaveBeenCalled());
|
||||
|
||||
setTimeout(finishDiscovery, 0);
|
||||
await worker.stop();
|
||||
|
||||
expect(releaseLock).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
});
|
||||
@@ -114,7 +114,6 @@ import type {UserContactChangeLogService} from '../user/services/UserContactChan
|
||||
import {UserDeletionEligibilityService} from '../user/services/UserDeletionEligibilityService';
|
||||
import {UserHarvestRepository} from '../user/UserHarvestRepository';
|
||||
import type {UserPermissionUtils} from '../utils/UserPermissionUtils';
|
||||
import {VoiceReconciliationWorker} from '../voice/VoiceReconciliationWorker';
|
||||
import type {VoiceRepository} from '../voice/VoiceRepository';
|
||||
import type {VoiceTopology} from '../voice/VoiceTopology';
|
||||
import type {WorkerTaskName} from './WorkerLaneConfig';
|
||||
@@ -165,7 +164,6 @@ export interface WorkerDependencies {
|
||||
voiceRoomStore: IVoiceRoomStore;
|
||||
liveKitService: ILiveKitService;
|
||||
voiceTopology: VoiceTopology | null;
|
||||
voiceReconciliationWorker: VoiceReconciliationWorker | null;
|
||||
channelService: ChannelService;
|
||||
guildAuditLogService: GuildAuditLogService;
|
||||
contactChangeLogService: UserContactChangeLogService;
|
||||
@@ -231,25 +229,8 @@ export async function initializeWorkerDependencies(snowflakeService: ISnowflakeS
|
||||
const voiceRoomStore = getVoiceRoomStoreInstance() ?? new InMemoryVoiceRoomStore();
|
||||
const liveKitService = getLiveKitServiceInstance() ?? new DisabledLiveKitService();
|
||||
const voiceAvailabilityService = getVoiceAvailabilityService();
|
||||
const voiceReconciliationEnabled = Config.worker.enableVoiceReconciliation;
|
||||
const voiceReconciliationWorker =
|
||||
Config.voice.enabled && voiceTopology !== null && voiceReconciliationEnabled
|
||||
? new VoiceReconciliationWorker({
|
||||
gatewayService,
|
||||
liveKitService,
|
||||
voiceRoomStore,
|
||||
kvClient,
|
||||
logger: Logger,
|
||||
intervalMs: Config.worker.voiceReconciliation.intervalMs,
|
||||
staggerDelayMs: Config.worker.voiceReconciliation.staggerDelayMs,
|
||||
lockTtlSeconds: Config.worker.voiceReconciliation.lockTtlSeconds,
|
||||
cadenceTtlSeconds: Config.worker.voiceReconciliation.cadenceTtlSeconds,
|
||||
gatewayOnlyGraceMs: Config.worker.voiceReconciliation.gatewayOnlyGraceMs,
|
||||
liveKitOnlyGraceMs: Config.worker.voiceReconciliation.liveKitOnlyGraceMs,
|
||||
})
|
||||
: null;
|
||||
if (Config.voice.enabled && voiceTopology !== null) {
|
||||
Logger.info({reconciliationEnabled: voiceReconciliationEnabled}, 'Voice services initialized');
|
||||
Logger.info('Voice services initialized');
|
||||
}
|
||||
const inviteRepository = getInviteRepository();
|
||||
const webhookRepository = getWebhookRepository();
|
||||
@@ -337,7 +318,6 @@ export async function initializeWorkerDependencies(snowflakeService: ISnowflakeS
|
||||
voiceRoomStore,
|
||||
liveKitService,
|
||||
voiceTopology,
|
||||
voiceReconciliationWorker,
|
||||
channelService,
|
||||
guildService,
|
||||
donationRepository,
|
||||
@@ -348,11 +328,3 @@ export async function initializeWorkerDependencies(snowflakeService: ISnowflakeS
|
||||
stripe,
|
||||
};
|
||||
}
|
||||
|
||||
export async function shutdownWorkerDependencies(deps: WorkerDependencies): Promise<void> {
|
||||
Logger.info('Shutting down worker dependencies...');
|
||||
if (deps.voiceReconciliationWorker !== null) {
|
||||
await deps.voiceReconciliationWorker.stop();
|
||||
}
|
||||
Logger.info('Worker dependencies shut down successfully');
|
||||
}
|
||||
|
||||
@@ -22,7 +22,7 @@ import {initializeSearch, shutdownSearch} from '../SearchFactory';
|
||||
import {CronScheduler} from './CronScheduler';
|
||||
import {JetStreamWorkerQueue} from './JetStreamWorkerQueue';
|
||||
import {clearWorkerDependencies, setWorkerDependencies} from './WorkerContext';
|
||||
import {initializeWorkerDependencies, shutdownWorkerDependencies, type WorkerDependencies} from './WorkerDependencies';
|
||||
import {initializeWorkerDependencies, type WorkerDependencies} from './WorkerDependencies';
|
||||
import {WorkerHeartbeat} from './WorkerHeartbeat';
|
||||
import {
|
||||
resolveCronSchedulerEnabled,
|
||||
@@ -118,11 +118,8 @@ export async function startWorkerMain(): Promise<void> {
|
||||
await jsConnectionManager?.drain();
|
||||
jsConnectionManager = null;
|
||||
});
|
||||
await cleanupStep('worker dependencies', async () => {
|
||||
if (dependencies) {
|
||||
await shutdownWorkerDependencies(dependencies);
|
||||
dependencies = null;
|
||||
}
|
||||
await cleanupStep('worker dependencies', () => {
|
||||
dependencies = null;
|
||||
clearWorkerDependencies();
|
||||
setInjectedWorkerService(undefined);
|
||||
});
|
||||
@@ -268,10 +265,6 @@ export async function startWorkerMain(): Promise<void> {
|
||||
} else {
|
||||
Logger.info('Search initialisation skipped for worker lanes without search tasks');
|
||||
}
|
||||
if (dependencies.voiceReconciliationWorker !== null) {
|
||||
dependencies.voiceReconciliationWorker.start();
|
||||
Logger.info('VoiceReconciliationWorker started');
|
||||
}
|
||||
if (cronSchedulerEnabled) {
|
||||
cron.start();
|
||||
Logger.info('Cron scheduler started');
|
||||
|
||||
Reference in New Issue
Block a user