This commit is contained in:
Iñaki Baz Castillo
2017-11-02 16:38:52 +01:00
parent 625f20f547
commit 52acf81eff
100 changed files with 26907 additions and 10234 deletions
+395 -421
View File
@@ -2,123 +2,53 @@
const EventEmitter = require('events').EventEmitter;
const protooServer = require('protoo-server');
const webrtc = require('mediasoup').webrtc;
const logger = require('./logger')('Room');
const Logger = require('./Logger');
const config = require('../config');
const MAX_BITRATE = config.mediasoup.maxBitrate || 3000000;
const MIN_BITRATE = Math.min(50000 || MAX_BITRATE);
const MAX_BITRATE = config.mediasoup.maxBitrate || 1000000;
const MIN_BITRATE = Math.min(50000, MAX_BITRATE);
const BITRATE_FACTOR = 0.75;
const MIN_AUDIO_LEVEL = -50;
const logger = new Logger('Room');
class Room extends EventEmitter
{
constructor(roomId, mediaServer)
{
logger.log('constructor() [roomId:"%s"]', roomId);
logger.info('constructor() [roomId:"%s"]', roomId);
super();
this.setMaxListeners(Infinity);
// Room ID.
this._roomId = roomId;
// Protoo Room instance.
this._protooRoom = new protooServer.Room();
// mediasoup Room instance.
this._mediaRoom = null;
// Pending peers (this is because at the time we get the first peer, the
// mediasoup room does not yet exist).
this._pendingProtooPeers = [];
// Closed flag.
this._closed = false;
try
{
// Protoo Room instance.
this._protooRoom = new protooServer.Room();
// mediasoup Room instance.
this._mediaRoom = mediaServer.Room(config.mediasoup.mediaCodecs);
}
catch (error)
{
this.close();
throw error;
}
// Current max bitrate for all the participants.
this._maxBitrate = MAX_BITRATE;
// Current active speaker mediasoup Peer.
this._activeSpeaker = null;
// Create a mediasoup room.
mediaServer.createRoom(
{
mediaCodecs : config.mediasoup.roomCodecs
})
.then((room) =>
{
logger.debug('mediasoup room created');
// Current active speaker.
// @type {mediasoup.Peer}
this._currentActiveSpeaker = null;
this._mediaRoom = room;
process.nextTick(() =>
{
this._mediaRoom.on('newpeer', (peer) =>
{
this._updateMaxBitrate();
peer.on('close', () =>
{
this._updateMaxBitrate();
});
});
// TODO: FIX?
this._mediaRoom.on('audiolevels', (entries) =>
{
logger.debug('room "audiolevels" event');
for (let entry of entries)
{
logger.debug('- [peer name:%s, rtpReceiver.id:%s, audio level:%s]',
entry.peer.name, entry.rtpReceiver.id, entry.audioLevel);
}
let activeSpeaker;
let activeLevel;
if (entries.length > 0)
{
activeSpeaker = entries[0].peer;
activeLevel = entries[0].audioLevel;
if (activeLevel < MIN_AUDIO_LEVEL)
{
activeSpeaker = null;
activeLevel = undefined;
}
}
else
{
activeSpeaker = null;
}
if (this._activeSpeaker !== activeSpeaker)
{
let data = {};
if (activeSpeaker)
{
logger.debug('active speaker [peer:"%s", volume:%s]',
activeSpeaker.name, activeLevel);
data.peer = { id: activeSpeaker.name };
data.level = activeLevel;
}
else
{
logger.debug('no current speaker');
data.peer = null;
}
this._protooRoom.spread('activespeaker', data);
}
this._activeSpeaker = activeSpeaker;
});
});
// Run all the pending join requests.
for (let protooPeer of this._pendingProtooPeers)
{
this._handleProtooPeer(protooPeer);
}
});
this._handleMediaRoom();
}
get id()
@@ -130,8 +60,11 @@ class Room extends EventEmitter
{
logger.debug('close()');
this._closed = true;
// Close the protoo Room.
this._protooRoom.close();
if (this._protooRoom)
this._protooRoom.close();
// Close the mediasoup Room.
if (this._mediaRoom)
@@ -146,367 +79,409 @@ class Room extends EventEmitter
if (!this._mediaRoom)
return;
logger.log(
logger.info(
'logStatus() [room id:"%s", protoo peers:%s, mediasoup peers:%s]',
this._roomId,
this._protooRoom.peers.length,
this._mediaRoom.peers.length);
}
createProtooPeer(peerId, transport)
handleConnection(peerName, transport)
{
logger.log('createProtooPeer() [peerId:"%s"]', peerId);
logger.info('handleConnection() [peerName:"%s"]', peerName);
if (this._protooRoom.hasPeer(peerId))
if (this._protooRoom.hasPeer(peerName))
{
logger.warn('createProtooPeer() | there is already a peer with same peerId, closing the previous one [peerId:"%s"]', peerId);
logger.warn(
'handleConnection() | there is already a peer with same peerName, ' +
'closing the previous one [peerName:"%s"]',
peerName);
let protooPeer = this._protooRoom.getPeer(peerId);
const protooPeer = this._protooRoom.getPeer(peerName);
protooPeer.close();
}
return this._protooRoom.createPeer(peerId, transport)
.then((protooPeer) =>
const protooPeer = this._protooRoom.createPeer(peerName, transport);
this._handleProtooPeer(protooPeer);
}
_handleMediaRoom()
{
logger.debug('_handleMediaRoom()');
const activeSpeakerDetector = this._mediaRoom.createActiveSpeakerDetector();
activeSpeakerDetector.on('activespeakerchange', (activePeer) =>
{
if (activePeer)
{
if (this._mediaRoom)
this._handleProtooPeer(protooPeer);
else
this._pendingProtooPeers.push(protooPeer);
});
logger.info('new active speaker [peerName:"%s"]', activePeer.name);
this._currentActiveSpeaker = activePeer;
const activeVideoProducer = activePeer.producers
.find((producer) => producer.kind === 'video');
for (const peer of this._mediaRoom.peers)
{
for (const consumer of peer.consumers)
{
if (consumer.kind !== 'video')
continue;
if (consumer.source === activeVideoProducer)
{
consumer.setPreferredProfile('high');
}
else
{
consumer.setPreferredProfile('low');
}
}
}
}
else
{
logger.info('no active speaker');
this._currentActiveSpeaker = null;
for (const peer of this._mediaRoom.peers)
{
for (const consumer of peer.consumers)
{
if (consumer.kind !== 'video')
continue;
consumer.setPreferredProfile('low');
}
}
}
// Spread to others via protoo.
this._protooRoom.spread(
'active-speaker',
{
peerName : activePeer ? activePeer.name : null
});
});
}
_handleProtooPeer(protooPeer)
{
logger.debug('_handleProtooPeer() [peerId:"%s"]', protooPeer.id);
logger.debug('_handleProtooPeer() [peer:"%s"]', protooPeer.id);
let mediaPeer = this._mediaRoom.Peer(protooPeer.id);
let peerconnection;
protooPeer.on('request', (request, accept, reject) =>
{
logger.debug(
'protoo "request" event [method:%s, peer:"%s"]',
request.method, protooPeer.id);
protooPeer.data.msids = [];
switch (request.method)
{
case 'mediasoup-request':
{
const mediasoupRequest = request.data;
this._handleMediasoupClientRequest(
protooPeer, mediasoupRequest, accept, reject);
break;
}
case 'mediasoup-notification':
{
accept();
const mediasoupNotification = request.data;
this._handleMediasoupClientNotification(
protooPeer, mediasoupNotification);
break;
}
case 'change-display-name':
{
accept();
const { displayName } = request.data;
const { mediaPeer } = protooPeer.data;
const oldDisplayName = mediaPeer.appData.displayName;
mediaPeer.appData.displayName = displayName;
// Spread to others via protoo.
this._protooRoom.spread(
'display-name-changed',
{
peerName : protooPeer.id,
displayName : displayName,
oldDisplayName : oldDisplayName
},
[ protooPeer ]);
break;
}
default:
{
logger.error('unknown request.method "%s"', request.method);
reject(400, `unknown request.method "${request.method}"`);
}
}
});
protooPeer.on('close', () =>
{
logger.debug('protoo Peer "close" event [peerId:"%s"]', protooPeer.id);
logger.debug('protoo Peer "close" event [peer:"%s"]', protooPeer.id);
this._protooRoom.spread(
'removepeer',
{
peer :
{
id : protooPeer.id,
msids : protooPeer.data.msids
}
});
const { mediaPeer } = protooPeer.data;
// Close the media stuff.
if (peerconnection)
peerconnection.close();
else
if (mediaPeer && !mediaPeer.closed)
mediaPeer.close();
// If this is the latest peer in the room, close the room.
// However, wait a bit (for reconnections).
// However wait a bit (for reconnections).
setTimeout(() =>
{
if (this._mediaRoom && this._mediaRoom.closed)
return;
if (this._protooRoom.peers.length === 0)
if (this._mediaRoom.peers.length === 0)
{
logger.log(
logger.info(
'last peer in the room left, closing the room [roomId:"%s"]',
this._roomId);
this.close();
}
}, 10000);
}, 5000);
});
Promise.resolve()
// Send 'join' request to the new peer.
.then(() =>
{
return protooPeer.send(
'joinme',
{
peerId : protooPeer.id,
roomId : this.id
});
})
// Create a RTCPeerConnection instance and set media capabilities.
.then((data) =>
{
peerconnection = new webrtc.RTCPeerConnection(
{
peer : mediaPeer,
usePlanB : !!data.usePlanB,
transportOptions : config.mediasoup.peerTransport,
maxBitrate : this._maxBitrate
});
// Store the RTCPeerConnection instance within the protoo Peer.
protooPeer.data.peerconnection = peerconnection;
mediaPeer.on('newtransport', (transport) =>
{
transport.on('iceselectedtuplechange', (data) =>
{
logger.log('"iceselectedtuplechange" event [peerId:"%s", protocol:%s, remoteIP:%s, remotePort:%s]',
protooPeer.id, data.protocol, data.remoteIP, data.remotePort);
});
});
// Set RTCPeerConnection capabilities.
return peerconnection.setCapabilities(data.capabilities);
})
// Send 'peers' request for the new peer to know about the existing peers.
.then(() =>
{
return protooPeer.send(
'peers',
{
peers : this._protooRoom.peers
// Filter this protoo Peer.
.filter((peer) =>
{
return peer !== protooPeer;
})
.map((peer) =>
{
return {
id : peer.id,
msids : peer.data.msids
};
})
});
})
// Tell all the other peers about the new peer.
.then(() =>
{
this._protooRoom.spread(
'addpeer',
{
peer :
{
id : protooPeer.id,
msids : protooPeer.data.msids
}
},
[ protooPeer ]);
})
.then(() =>
{
// Send initial SDP offer.
return this._sendOffer(protooPeer,
{
offerToReceiveAudio : 1,
offerToReceiveVideo : 1
});
})
.then(() =>
{
// Handle PeerConnection events.
peerconnection.on('negotiationneeded', () =>
{
logger.debug('"negotiationneeded" event [peerId:"%s"]', protooPeer.id);
// Send SDP re-offer.
this._sendOffer(protooPeer);
});
peerconnection.on('signalingstatechange', () =>
{
logger.debug('"signalingstatechange" event [peerId:"%s", signalingState:%s]',
protooPeer.id, peerconnection.signalingState);
});
})
.then(() =>
{
protooPeer.on('request', (request, accept, reject) =>
{
logger.debug('protoo Peer "request" event [method:%s]', request.method);
switch(request.method)
{
case 'reofferme':
{
accept();
this._sendOffer(protooPeer);
break;
}
case 'restartice':
{
peerconnection.restartIce()
.then(() =>
{
accept();
})
.catch((error) =>
{
logger.error('"restartice" request failed: %s', error);
logger.error('stack:\n' + error.stack);
reject(500, `"restartice" failed: ${error.message}`);
});
break;
}
case 'disableremotevideo':
{
let videoMsid = request.data.msid;
let disable = request.data.disable;
let videoRtpSender;
for (let rtpSender of mediaPeer.rtpSenders)
{
if (rtpSender.kind !== 'video')
continue;
let msid = rtpSender.rtpParameters.userParameters.msid.split(/\s/)[0];
if (msid === videoMsid)
{
videoRtpSender = rtpSender;
break;
}
}
if (videoRtpSender)
{
return Promise.resolve()
.then(() =>
{
if (disable)
return videoRtpSender.disable({ emit: false });
else
return videoRtpSender.enable({ emit: false });
})
.then(() =>
{
logger.log('"disableremotevideo" request succeed [disable:%s]',
!!disable);
accept();
})
.catch((error) =>
{
logger.error('"disableremotevideo" request failed: %s', error);
logger.error('stack:\n' + error.stack);
reject(500, `"disableremotevideo" failed: ${error.message}`);
});
}
else
{
reject(404, 'msid not found');
}
break;
}
default:
{
logger.error('unknown method');
reject(404, 'unknown method');
}
}
});
})
.catch((error) =>
{
logger.error('_handleProtooPeer() failed: %s', error.message);
logger.error('stack:\n' + error.stack);
protooPeer.close();
});
}
_sendOffer(protooPeer, options)
_handleMediaPeer(protooPeer, mediaPeer)
{
logger.debug('_sendOffer() [peerId:"%s"]', protooPeer.id);
mediaPeer.on('notify', (notification) =>
{
protooPeer.send('mediasoup-notification', notification)
.catch(() => {});
});
let peerconnection = protooPeer.data.peerconnection;
let mediaPeer = peerconnection.peer;
mediaPeer.on('newtransport', (transport) =>
{
logger.info(
'mediaPeer "newtransport" event [id:%s, direction:%s]',
transport.id, transport.direction);
return Promise.resolve()
.then(() =>
// Update peers max sending bitrate.
if (transport.direction === 'send')
{
return peerconnection.createOffer(options);
})
.then((desc) =>
{
return peerconnection.setLocalDescription(desc);
})
// Send the SDP offer to the peer.
.then(() =>
{
return protooPeer.send(
'offer',
{
offer : peerconnection.localDescription.serialize()
});
})
// Process the SDP answer from the peer.
.then((data) =>
{
let answer = data.answer;
this._updateMaxBitrate();
return peerconnection.setRemoteDescription(answer);
})
.then(() =>
{
let oldMsids = protooPeer.data.msids;
// Reset peer's msids.
protooPeer.data.msids = [];
let setMsids = new Set();
// Update peer's msids information.
for (let rtpReceiver of mediaPeer.rtpReceivers)
transport.on('close', () =>
{
let msid = rtpReceiver.rtpParameters.userParameters.msid.split(/\s/)[0];
this._updateMaxBitrate();
});
}
setMsids.add(msid);
this._handleMediaTransport(transport);
});
mediaPeer.on('newproducer', (producer) =>
{
logger.info('mediaPeer "newproducer" event [id:%s]', producer.id);
this._handleMediaProducer(producer);
});
mediaPeer.on('newconsumer', (consumer) =>
{
logger.info('mediaPeer "newconsumer" event [id:%s]', consumer.id);
this._handleMediaConsumer(consumer);
});
// Also handle already existing Consumers.
for (const consumer of mediaPeer.consumers)
{
logger.info('mediaPeer existing "consumer" [id:%s]', consumer.id);
this._handleMediaConsumer(consumer);
}
// Notify about the existing active speaker.
if (this._currentActiveSpeaker)
{
protooPeer.send(
'active-speaker',
{
peerName : this._currentActiveSpeaker.name
})
.catch(() => {});
}
}
_handleMediaTransport(transport)
{
transport.on('close', (originator) =>
{
logger.info(
'Transport "close" event [originator:%s]', originator);
});
}
_handleMediaProducer(producer)
{
producer.on('close', (originator) =>
{
logger.info(
'Producer "close" event [originator:%s]', originator);
});
producer.on('pause', (originator) =>
{
logger.info(
'Producer "pause" event [originator:%s]', originator);
});
producer.on('resume', (originator) =>
{
logger.info(
'Producer "resume" event [originator:%s]', originator);
});
}
_handleMediaConsumer(consumer)
{
consumer.on('close', (originator) =>
{
logger.info(
'Consumer "close" event [originator:%s]', originator);
});
consumer.on('pause', (originator) =>
{
logger.info(
'Consumer "pause" event [originator:%s]', originator);
});
consumer.on('resume', (originator) =>
{
logger.info(
'Consumer "resume" event [originator:%s]', originator);
});
consumer.on('effectiveprofilechange', (profile) =>
{
logger.info(
'Consumer "effectiveprofilechange" event [profile:%s]', profile);
});
// If video, initially make it 'low' profile unless this is for the current
// active speaker.
if (consumer.kind === 'video' && consumer.peer !== this._currentActiveSpeaker)
consumer.setPreferredProfile('low');
}
_handleMediasoupClientRequest(protooPeer, request, accept, reject)
{
logger.debug(
'mediasoup-client request [method:%s, peer:"%s"]',
request.method, protooPeer.id);
switch (request.method)
{
case 'queryRoom':
{
this._mediaRoom.receiveRequest(request)
.then((response) => accept(response))
.catch((error) => reject(500, error.toString()));
break;
}
case 'join':
{
// TODO: Handle appData. Yes?
const { peerName } = request;
if (peerName !== protooPeer.id)
{
reject(403, 'that is not your corresponding mediasoup Peer name');
break;
}
else if (protooPeer.data.mediaPeer)
{
reject(500, 'already have a mediasoup Peer');
break;
}
protooPeer.data.msids = Array.from(setMsids);
// If msids changed, notify.
let sameValues = (
oldMsids.length == protooPeer.data.msids.length) &&
oldMsids.every((element, index) =>
this._mediaRoom.receiveRequest(request)
.then((response) =>
{
return element === protooPeer.data.msids[index];
accept(response);
// Get the newly created mediasoup Peer.
const mediaPeer = this._mediaRoom.getPeerByName(peerName);
protooPeer.data.mediaPeer = mediaPeer;
this._handleMediaPeer(protooPeer, mediaPeer);
})
.catch((error) =>
{
reject(500, error.toString());
});
if (!sameValues)
{
this._protooRoom.spread(
'updatepeer',
{
peer :
{
id : protooPeer.id,
msids : protooPeer.data.msids
}
},
[ protooPeer ]);
}
})
.catch((error) =>
{
logger.error('_sendOffer() failed: %s', error);
logger.error('stack:\n' + error.stack);
break;
}
logger.warn('resetting peerconnection');
peerconnection.reset();
});
default:
{
const { mediaPeer } = protooPeer.data;
if (!mediaPeer)
{
logger.error(
'cannot handle mediasoup request, no mediasoup Peer [method:"%s"]',
request.method);
reject(400, 'no mediasoup Peer');
}
mediaPeer.receiveRequest(request)
.then((response) => accept(response))
.catch((error) => reject(500, error.toString()));
}
}
}
_handleMediasoupClientNotification(protooPeer, notification)
{
logger.debug(
'mediasoup-client notification [method:%s, peer:"%s"]',
notification.method, protooPeer.id);
// NOTE: mediasoup-client just sends notifications with target 'peer',
// so first of all, get the mediasoup Peer.
const { mediaPeer } = protooPeer.data;
if (!mediaPeer)
{
logger.error(
'cannot handle mediasoup notification, no mediasoup Peer [method:"%s"]',
notification.method);
return;
}
mediaPeer.receiveNotification(notification);
}
_updateMaxBitrate()
@@ -514,8 +489,8 @@ class Room extends EventEmitter
if (this._mediaRoom.closed)
return;
let numPeers = this._mediaRoom.peers.length;
let previousMaxBitrate = this._maxBitrate;
const numPeers = this._mediaRoom.peers.length;
const previousMaxBitrate = this._maxBitrate;
let newMaxBitrate;
if (numPeers <= 2)
@@ -530,29 +505,28 @@ class Room extends EventEmitter
newMaxBitrate = MIN_BITRATE;
}
if (newMaxBitrate === previousMaxBitrate)
return;
this._maxBitrate = newMaxBitrate;
for (let peer of this._mediaRoom.peers)
for (const peer of this._mediaRoom.peers)
{
if (!peer.capabilities || peer.closed)
continue;
for (let transport of peer.transports)
for (const transport of peer.transports)
{
if (transport.closed)
continue;
transport.setMaxBitrate(newMaxBitrate);
if (transport.direction === 'send')
{
transport.setMaxBitrate(newMaxBitrate)
.catch((error) =>
{
logger.error('transport.setMaxBitrate() failed: %s', String(error));
});
}
}
}
logger.log('_updateMaxBitrate() [num peers:%s, before:%skbps, now:%skbps]',
logger.info(
'_updateMaxBitrate() [num peers:%s, before:%skbps, now:%skbps]',
numPeers,
Math.round(previousMaxBitrate / 1000),
Math.round(newMaxBitrate / 1000));
this._maxBitrate = newMaxBitrate;
}
}
+15 -16
View File
@@ -2,7 +2,7 @@
const debug = require('debug');
const NAMESPACE = 'mediasoup-demo-server';
const APP_NAME = 'mediasoup-demo-server';
class Logger
{
@@ -10,23 +10,25 @@ class Logger
{
if (prefix)
{
this._debug = debug(NAMESPACE + ':' + prefix);
this._log = debug(NAMESPACE + ':LOG:' + prefix);
this._warn = debug(NAMESPACE + ':WARN:' + prefix);
this._error = debug(NAMESPACE + ':ERROR:' + prefix);
this._debug = debug(`${APP_NAME}:${prefix}`);
this._info = debug(`${APP_NAME}:INFO:${prefix}`);
this._warn = debug(`${APP_NAME}:WARN:${prefix}`);
this._error = debug(`${APP_NAME}:ERROR:${prefix}`);
}
else
{
this._debug = debug(NAMESPACE);
this._log = debug(NAMESPACE + ':LOG');
this._warn = debug(NAMESPACE + ':WARN');
this._error = debug(NAMESPACE + ':ERROR');
this._debug = debug(APP_NAME);
this._info = debug(`${APP_NAME}:INFO`);
this._warn = debug(`${APP_NAME}:WARN`);
this._error = debug(`${APP_NAME}:ERROR`);
}
/* eslint-disable no-console */
this._debug.log = console.info.bind(console);
this._log.log = console.info.bind(console);
this._info.log = console.info.bind(console);
this._warn.log = console.warn.bind(console);
this._error.log = console.error.bind(console);
/* eslint-enable no-console */
}
get debug()
@@ -34,9 +36,9 @@ class Logger
return this._debug;
}
get log()
get info()
{
return this._log;
return this._info;
}
get warn()
@@ -50,7 +52,4 @@ class Logger
}
}
module.exports = function(prefix)
{
return new Logger(prefix);
};
module.exports = Logger;