diff --git a/server/lib/Room.js b/server/lib/Room.js index 0130569c..b22b368e 100644 --- a/server/lib/Room.js +++ b/server/lib/Room.js @@ -64,7 +64,7 @@ const ROUTER_SCALE_SIZE = config.routerScaleSize || 40; class Room extends EventEmitter { - static getLeastLoadedRouter(mediasoupWorkers, peers, mediasoupRouters) + static getLeastLoadedRouter(mediasoupWorkers, peers, mediasoupRouters, pipedRoutersIds = []) { const routerLoads = new Map(); @@ -107,25 +107,58 @@ class Room extends EventEmitter } } - let load = Infinity; - let pid; + const sortedWorkerLoads = new Map([...workerLoads.entries()].sort((a, b) => a[1] - b[1])); - for (const workerId of workerLoads.keys()) + // we don't care about if router is piped, just choose the least loaded worker + if (pipedRoutersIds.length === 0 || pipedRoutersIds.length === mediasoupRouters.size) { - const workerLoad = workerLoads.get(workerId); + const workerId = sortedWorkerLoads.keys().next().value; - if (workerLoad < load) + for (const worker of mediasoupWorkers) { - pid = workerId; - load = workerLoad; + if (worker._pid === workerId) + { + for (const router of worker._routers) + { + const routerId = router._internal.routerId; + + if (mediasoupRouters.has(routerId)) + { + return routerId; + } + } + } + } + } + else + { + //find if there is a piped router that is on a worker that is below limit + for (const [workerId, workerLoad] of sortedWorkerLoads.entries()) + { + for (const worker of mediasoupWorkers) + { + if (worker._pid === workerId) + { + for (const router of worker._routers) + { + const routerId = router._internal.routerId; + + //on purpose we check if the worker loat is below the limit, as in reality the worker load is imortant, not the router load + if (mediasoupRouters.has(routerId) && pipedRoutersIds.includes(routerId) && workerLoad < ROUTER_SCALE_SIZE) + { + return routerId; + } + } + } } } - let id; + //no piped router found, we need to return router from least loaded worker + const workerId = sortedWorkerLoads.keys().next().value; for (const worker of mediasoupWorkers) { - if (worker._pid === pid) + if (worker._pid === workerId) { for (const router of worker._routers) @@ -134,15 +167,13 @@ class Room extends EventEmitter if (mediasoupRouters.has(routerId)) { - id = routerId; - break; + return routerId; + } } } } } - - return id; } /** @@ -171,7 +202,7 @@ class Room extends EventEmitter } const firstRouter = mediasoupRouters.get(Room.getLeastLoadedRouter( - mediasoupWorkers, peers, mediasoupRouters, {})); + mediasoupWorkers, peers, mediasoupRouters)); // Create a mediasoup AudioLevelObserver on first router const audioLevelObserver = await firstRouter.createAudioLevelObserver( @@ -247,8 +278,6 @@ class Room extends EventEmitter // Array of mediasoup Router instances. this._mediasoupRouters = mediasoupRouters; - this._currentRouter = firstRouter; - // mediasoup AudioLevelObserver. this._audioLevelObserver = audioLevelObserver; @@ -302,8 +331,6 @@ class Room extends EventEmitter router.close(); } - this._currentRouter = null; - this._allPeers = null; this._pipedRoutersIds = null; @@ -2029,11 +2056,13 @@ class Room extends EventEmitter } } - async _pipeProducersToNewRouter() + async _pipeProducersToNewRouter(routerId) { + const router = this._mediasoupRouters.get(routerId); + const peersToPipe = Object.values(this._peers) - .filter((peer) => peer.routerId !== this._currentRouter.id); + .filter((peer) => peer.routerId !== routerId); for (const peer of peersToPipe) { @@ -2043,7 +2072,7 @@ class Room extends EventEmitter { await srcRouter.pipeToRouter({ producerId, - router : this._currentRouter + router : router }); } } @@ -2051,23 +2080,17 @@ class Room extends EventEmitter async _getRouterId() { - const routerLoad = - Object.values(this._peers) - .filter((peer) => peer.routerId === this._currentRouter.id).length; + const routerId = Room.getLeastLoadedRouter( + this._mediasoupWorkers,this._allPeers, + this._mediasoupRouters, this._pipedRoutersIds); - if (routerLoad >= ROUTER_SCALE_SIZE) + if (!this._pipedRoutersIds.includes(routerId)) { - this._currentRouter = this._mediasoupRouters.get(Room.getLeastLoadedRouter( - this._mediasoupWorkers, this._allPeers, this._mediasoupRouters)); - - if (!this._pipedRoutersIds.includes(this._currentRouter.id)) - { - await this._pipeProducersToNewRouter(); - this._pipedRoutersIds.push(this._currentRouter.id); - } + await this._pipeProducersToNewRouter(routerId); + this._pipedRoutersIds.push(routerId); } - return this._currentRouter.id; + return routerId; } // Returns an array of router ids we need to pipe to