From 2559331c6531933c91d799e1c532cf861035c8fe Mon Sep 17 00:00:00 2001 From: Evgeny Melnikov Date: Tue, 10 Mar 2026 15:50:51 +0300 Subject: [PATCH] fix: memory leaks, bugs and logic fixes --- src/NetworkScoresCalculator.ts | 6 ++++-- src/detectors/BaseIssueDetector.ts | 6 ++++-- src/detectors/InboundNetworkIssueDetector.ts | 2 +- src/detectors/NetworkMediaSyncIssueDetector.ts | 7 ++++++- src/detectors/OutboundNetworkIssueDetector.ts | 4 ++-- src/helpers/calc.ts | 5 +++++ src/helpers/streams.ts | 3 ++- src/parser/RTCStatsParser.ts | 6 ++++-- src/utils/tasks.ts | 2 -- 9 files changed, 28 insertions(+), 13 deletions(-) diff --git a/src/NetworkScoresCalculator.ts b/src/NetworkScoresCalculator.ts index c570dd6..5f2a20f 100644 --- a/src/NetworkScoresCalculator.ts +++ b/src/NetworkScoresCalculator.ts @@ -6,7 +6,7 @@ import { WebRTCStatsParsed, NetworkQualityStatsSample, } from './types'; -import { scheduleTask } from './utils/tasks'; +import { createTaskScheduler } from './utils/tasks'; import { CLEANUP_PREV_STATS_TTL_MS } from './utils/constants'; type MosCalculatorResult = { @@ -17,13 +17,15 @@ type MosCalculatorResult = { class NetworkScoresCalculator implements INetworkScoresCalculator { #lastProcessedStats: { [connectionId: string]: WebRTCStatsParsed } = {}; + readonly #scheduleTask = createTaskScheduler(); + calculate(data: WebRTCStatsParsed): NetworkScores { const { connection: { id: connectionId } } = data; const { mos: outbound, stats: outboundStatsSample } = this.calculateOutboundScore(data) || {}; const { mos: inbound, stats: inboundStatsSample } = this.calculateInboundScore(data) || {}; this.#lastProcessedStats[connectionId] = data; - scheduleTask({ + this.#scheduleTask({ taskId: connectionId, delayMs: CLEANUP_PREV_STATS_TTL_MS, callback: () => (delete this.#lastProcessedStats[connectionId]), diff --git a/src/detectors/BaseIssueDetector.ts b/src/detectors/BaseIssueDetector.ts index cfe05a5..437d771 100644 --- a/src/detectors/BaseIssueDetector.ts +++ b/src/detectors/BaseIssueDetector.ts @@ -5,7 +5,7 @@ import { WebRTCStatsParsed, WebRTCStatsParsedWithNetworkScores, } from '../types'; -import { scheduleTask } from '../utils/tasks'; +import { createTaskScheduler } from '../utils/tasks'; import { CLEANUP_PREV_STATS_TTL_MS, MAX_PARSED_STATS_STORAGE_SIZE } from '../utils/constants'; export interface PrevStatsCleanupPayload { @@ -25,6 +25,8 @@ abstract class BaseIssueDetector implements IssueDetector { readonly #maxParsedStatsStorageSize: number; + readonly #scheduleTask = createTaskScheduler(); + constructor(params: BaseIssueDetectorParams = {}) { this.#statsCleanupDelayMs = params.statsCleanupTtlMs ?? CLEANUP_PREV_STATS_TTL_MS; this.#maxParsedStatsStorageSize = params.maxParsedStatsStorageSize ?? MAX_PARSED_STATS_STORAGE_SIZE; @@ -57,7 +59,7 @@ abstract class BaseIssueDetector implements IssueDetector { return; } - scheduleTask({ + this.#scheduleTask({ taskId: connectionId, delayMs: this.#statsCleanupDelayMs, callback: () => { diff --git a/src/detectors/InboundNetworkIssueDetector.ts b/src/detectors/InboundNetworkIssueDetector.ts index 14234fc..68372af 100644 --- a/src/detectors/InboundNetworkIssueDetector.ts +++ b/src/detectors/InboundNetworkIssueDetector.ts @@ -23,7 +23,7 @@ class InboundNetworkIssueDetector extends BaseIssueDetector { readonly #highRttThresholdMs: number; constructor(params: InboundNetworkIssueDetectorParams = {}) { - super(); + super(params); this.#highPacketLossThresholdPct = params.highPacketLossThresholdPct ?? 5; this.#highJitterThreshold = params.highJitterThreshold ?? 200; this.#highJitterBufferDelayThresholdMs = params.highJitterBufferDelayThresholdMs ?? 500; diff --git a/src/detectors/NetworkMediaSyncIssueDetector.ts b/src/detectors/NetworkMediaSyncIssueDetector.ts index 7fa36c9..07c7133 100644 --- a/src/detectors/NetworkMediaSyncIssueDetector.ts +++ b/src/detectors/NetworkMediaSyncIssueDetector.ts @@ -14,7 +14,7 @@ class NetworkMediaSyncIssueDetector extends BaseIssueDetector { readonly #correctedSamplesThresholdPct: number; constructor(params: NetworkMediaSyncIssueDetectorParams = {}) { - super(); + super(params); this.#correctedSamplesThresholdPct = params.correctedSamplesThresholdPct ?? 5; } @@ -47,6 +47,11 @@ class NetworkMediaSyncIssueDetector extends BaseIssueDetector { } const deltaSamplesReceived = stats.track.totalSamplesReceived - previousStreamStats.track.totalSamplesReceived; + + if (deltaSamplesReceived === 0) { + return; + } + const deltaCorrectedSamples = nowCorrectedSamples - lastCorrectedSamples; const correctedSamplesPct = Math.round((deltaCorrectedSamples * 100) / deltaSamplesReceived); const statsSample = { diff --git a/src/detectors/OutboundNetworkIssueDetector.ts b/src/detectors/OutboundNetworkIssueDetector.ts index 133a623..ca95a93 100644 --- a/src/detectors/OutboundNetworkIssueDetector.ts +++ b/src/detectors/OutboundNetworkIssueDetector.ts @@ -17,7 +17,7 @@ class OutboundNetworkIssueDetector extends BaseIssueDetector { readonly #highJitterThreshold: number; constructor(params: OutboundNetworkIssueDetectorParams = {}) { - super(); + super(params); this.#highPacketLossThresholdPct = params.highPacketLossThresholdPct ?? 5; this.#highJitterThreshold = params.highJitterThreshold ?? 200; } @@ -79,7 +79,7 @@ class OutboundNetworkIssueDetector extends BaseIssueDetector { const isHighPacketsLoss = packetLossPct > this.#highPacketLossThresholdPct; const isHighJitter = avgJitter >= this.#highJitterThreshold; const isNetworkMediaLatencyIssue = isHighPacketsLoss && isHighJitter; - const isNetworkIssue = (!isHighPacketsLoss && isHighJitter) || isHighJitter || isHighPacketsLoss; + const isNetworkIssue = isHighJitter || isHighPacketsLoss; const statsSample = { rtt, diff --git a/src/helpers/calc.ts b/src/helpers/calc.ts index 7f05eb5..d390561 100644 --- a/src/helpers/calc.ts +++ b/src/helpers/calc.ts @@ -15,6 +15,11 @@ export const calculateVolatility = (values: number[]) => { } const mean = calculateMean(values); + + if (mean === 0) { + return 0; + } + const meanAbsoluteDeviationFps = values.reduce((acc, val) => acc + Math.abs(val - mean), 0) / values.length; return (meanAbsoluteDeviationFps * 100) / mean; }; diff --git a/src/helpers/streams.ts b/src/helpers/streams.ts index 522a5d0..e30b16d 100644 --- a/src/helpers/streams.ts +++ b/src/helpers/streams.ts @@ -7,7 +7,8 @@ export const isDtxLikeBehavior = ( stdDevThreshold = 30, ): boolean => { const frameIntervals: number[] = []; - for (let i = 1; i < allProcessedStats.length - 1; i += 1) { + + for (let i = 1; i < allProcessedStats.length; i += 1) { const videoStreamStats = allProcessedStats[i]?.video?.inbound.find( (stream) => stream.ssrc === ssrc, ); diff --git a/src/parser/RTCStatsParser.ts b/src/parser/RTCStatsParser.ts index 54d96b3..933ebe8 100644 --- a/src/parser/RTCStatsParser.ts +++ b/src/parser/RTCStatsParser.ts @@ -15,7 +15,7 @@ import { Logger, } from '../types'; import { checkIsConnectionClosed, calcBitrate } from './utils'; -import { scheduleTask } from '../utils/tasks'; +import { createTaskScheduler } from '../utils/tasks'; import { CLEANUP_PREV_STATS_TTL_MS } from '../utils/constants'; interface PrevStatsItem { @@ -31,6 +31,8 @@ interface WebRTCStatsParserParams { class RTCStatsParser implements StatsParser { private readonly prevStats = new Map(); + private readonly scheduleTask = createTaskScheduler(); + private readonly allowedReportTypes: Set = new Set([ 'candidate-pair', 'inbound-rtp', @@ -133,7 +135,7 @@ class RTCStatsParser implements StatsParser { ts: Date.now(), }); - scheduleTask({ + this.scheduleTask({ taskId: connectionId, delayMs: CLEANUP_PREV_STATS_TTL_MS, callback: () => (this.prevStats.delete(connectionId)), diff --git a/src/utils/tasks.ts b/src/utils/tasks.ts index 11ac41c..f637d70 100644 --- a/src/utils/tasks.ts +++ b/src/utils/tasks.ts @@ -26,5 +26,3 @@ export const createTaskScheduler = () => { scheduledTasks.set(taskId, newTimer); }; }; - -export const scheduleTask = createTaskScheduler();