[+] 信令调整

This commit is contained in:
acgist
2023-03-03 08:09:35 +08:00
parent b53f50adf7
commit c3dbe52d5c
36 changed files with 588 additions and 306 deletions

View File

@@ -18,6 +18,8 @@ module.exports = {
version: "1.0.0",
// 终端标识
clientId: "taoyao-client-media",
// 终端名称
name: "桃夭媒体服务",
// 地址
host: "127.0.0.1",
// 端口

View File

@@ -2,7 +2,7 @@
const config = require("./Config");
const mediasoup = require("mediasoup");
const { Signal, signalChannel } = require("./Signal");
const { Taoyao, signalChannel } = require("./Taoyao");
// 线程名称
process.title = config.name;
@@ -11,8 +11,8 @@ process.env.NODE_TLS_REJECT_UNAUTHORIZED = "0";
// Mediasoup Worker列表
const mediasoupWorkers = [];
// 信令服务
const signal = new Signal(mediasoupWorkers);
// 桃夭
const taoyao = new Taoyao(mediasoupWorkers);
/**
* 创建Mediasoup Worker列表
@@ -56,12 +56,8 @@ async function buildMediasoupWorkers() {
* 连接信令服务
*/
async function connectSignalServer() {
await signalChannel.connect(
`wss://${config.signal.host}:${config.signal.port}/websocket.signal`,
function (message) {
signal.on(message);
}
);
signalChannel.taoyao = taoyao;
await signalChannel.connect(`wss://${config.signal.host}:${config.signal.port}/websocket.signal`);
}
/**

View File

@@ -20,19 +20,19 @@ const protocol = {
this.index = 0;
}
const date = new Date();
return 100000000000000 * date.getDate() +
1000000000000 * date.getHours() +
10000000000 * date.getMinutes() +
100000000 * date.getSeconds() +
1000 * this.clientIndex +
this.index;
return (
100000000000000 * date.getDate() +
1000000000000 * date.getHours() +
10000000000 * date.getMinutes() +
100000000 * date.getSeconds() +
1000 * this.clientIndex +
this.index
);
},
/**
* 生成信令消息
*
* @param {*} signal 信令标识
* @param {*} body 信令消息
* @param {*} id ID
* @param {*} body 消息主体
* @param {*} id 消息ID
*
* @returns 信令消息
*/
@@ -53,12 +53,12 @@ const protocol = {
* 信令通道
*/
const signalChannel = {
// 桃夭
taoyao: null,
// 通道
channel: null,
// 地址
address: null,
// 回调
callback: null,
// 心跳时间
heartbeatTime: 30 * 1000,
// 心跳定时器
@@ -70,7 +70,7 @@ const signalChannel = {
// 防止重复重连
lockReconnect: false,
// 当前重连时间
connectionTimeout: 5 * 1000,
reconnectionTimeout: 5 * 1000,
// 最小重连时间
minReconnectionDelay: 5 * 1000,
// 最大重连时间
@@ -79,92 +79,87 @@ const signalChannel = {
* 心跳
*/
heartbeat() {
const self = this;
if (self.heartbeatTimer) {
clearTimeout(self.heartbeatTimer);
const me = this;
if (me.heartbeatTimer) {
clearTimeout(me.heartbeatTimer);
}
self.heartbeatTimer = setTimeout(async function () {
if (self.channel && self.channel.readyState === WebSocket.OPEN) {
// TODO信号强度、电池信息
self.push(
me.heartbeatTimer = setTimeout(async function () {
if (me.channel && me.channel.readyState === WebSocket.OPEN) {
me.push(
// TODO电池信息
protocol.buildMessage("client::heartbeat", {
signal: 100,
battery: 100,
charging: true,
})
);
self.heartbeat();
me.heartbeat();
} else {
console.warn("发送心跳失败:", self.address);
console.warn("心跳失败:", me.address);
}
}, self.heartbeatTime);
}, me.heartbeatTime);
},
/**
* 连接
*
* @param {*} address 地址
* @param {*} callback 回调
* @param {*} reconnection 是否重连
*
* @returns Promise
*/
async connect(address, callback, reconnection = true) {
const self = this;
if (self.channel && self.channel.readyState === WebSocket.OPEN) {
async connect(address, reconnection = true) {
const me = this;
if (me.channel && me.channel.readyState === WebSocket.OPEN) {
return new Promise((resolve, reject) => {
resolve(self.channel);
resolve(me.channel);
});
}
self.address = address;
self.callback = callback;
self.reconnection = reconnection;
me.address = address;
me.reconnection = reconnection;
return new Promise((resolve, reject) => {
console.debug("连接信令通道:", self.address);
self.channel = new WebSocket(self.address, { handshakeTimeout: 5000 });
self.channel.on("open", async function () {
console.info("打开信令通道:", self.address);
// 注册终端
// TODO信号强度、电池信息
self.push(
console.debug("连接信令通道:", me.address);
me.channel = new WebSocket(me.address, { handshakeTimeout: 5000 });
me.channel.on("open", async function () {
console.info("打开信令通道:", me.address);
// TODO电池信息
me.push(
protocol.buildMessage("client::register", {
name: "桃夭媒体服务",
clientId: config.signal.clientId,
name: config.signal.name,
clientType: "MEDIA",
signal: 100,
battery: 100,
charging: true,
username: config.signal.username,
password: config.signal.password,
})
);
// 重置时间
self.connectionTimeout = self.minReconnectionDelay;
// 开始心跳
self.heartbeat();
// 成功回调
resolve(self.channel);
me.reconnectionTimeout = me.minReconnectionDelay;
me.taoyao.connect = true;
me.heartbeat();
resolve(me.channel);
});
self.channel.on("close", async function () {
console.warn("信令通道关闭:", self.address);
if (self.reconnection) {
self.reconnect();
me.channel.on("close", async function () {
console.warn("信令通道关闭:", me.address);
me.taoyao.connect = false;
if (me.reconnection) {
me.reconnect();
}
// 不要失败回调
});
self.channel.on("error", async function (e) {
console.error("信令通道异常:", self.address, e);
if (self.reconnection) {
self.reconnect();
me.channel.on("error", async function (e) {
console.error("信令通道异常:", me.address, e);
me.taoyao.connect = false;
if (me.reconnection) {
me.reconnect();
}
// 不要失败回调
});
self.channel.on("message", async function (data) {
me.channel.on("message", async function (data) {
try {
const content = data.toString();
console.debug("信令通道消息:", content);
self.callback(JSON.parse(content));
me.taoyao.on(JSON.parse(content));
} catch (error) {
console.error("处理信令消息异常:", content, error);
console.error("处理信令消息异常:", data, error);
}
});
});
@@ -173,26 +168,26 @@ const signalChannel = {
* 重连
*/
reconnect() {
const self = this;
const me = this;
if (
self.lockReconnect ||
(self.channel && self.channel.readyState === WebSocket.OPEN)
me.lockReconnect ||
(me.channel && me.channel.readyState === WebSocket.OPEN)
) {
return;
}
self.lockReconnect = true;
if (self.reconnectTimer) {
clearTimeout(self.reconnectTimer);
me.lockReconnect = true;
if (me.reconnectTimer) {
clearTimeout(me.reconnectTimer);
}
// 定时重连
self.reconnectTimer = setTimeout(function () {
console.info("信令通道重连", self.address);
self.connect(self.address, self.callback, self.reconnection);
self.lockReconnect = false;
}, self.connectionTimeout);
self.connectionTimeout = Math.min(
self.connectionTimeout + self.minReconnectionDelay,
self.maxReconnectionDelay
me.reconnectTimer = setTimeout(function () {
console.info("重连信令通道:", me.address);
me.connect(me.address, me.reconnection);
me.lockReconnect = false;
}, me.reconnectionTimeout);
me.reconnectionTimeout = Math.min(
me.reconnectionTimeout + me.minReconnectionDelay,
me.maxReconnectionDelay
);
},
/**
@@ -200,7 +195,7 @@ const signalChannel = {
*
* @param {*} message 消息
*/
push(message) {
push(message) {
try {
this.channel.send(JSON.stringify(message));
} catch (error) {
@@ -211,11 +206,11 @@ const signalChannel = {
* 关闭通道
*/
close() {
const self = this;
self.reconnection = false;
self.channel.close();
clearTimeout(self.heartbeatTimer);
clearTimeout(self.reconnectTimer);
const me = this;
clearTimeout(me.heartbeatTimer);
clearTimeout(me.reconnectTimer);
me.reconnection = false;
me.channel.close();
},
};
@@ -227,8 +222,8 @@ class Room {
close = false;
// 房间ID
roomId = null;
// 信令
signal = null;
// 桃夭
taoyao = null;
// WebRTCServer
webRtcServer = null;
// 路由
@@ -252,7 +247,7 @@ class Room {
constructor({
roomId,
signal,
taoyao,
webRtcServer,
mediasoupRouter,
audioLevelObserver,
@@ -261,7 +256,7 @@ class Room {
this.close = false;
this.roomId = roomId;
this.networkThrottled = false;
this.signal = signal;
this.taoyao = taoyao;
this.webRtcServer = webRtcServer;
this.mediasoupRouter = mediasoupRouter;
this.audioLevelObserver = audioLevelObserver;
@@ -340,9 +335,11 @@ class Room {
}
/**
* 信令服务
* 桃夭
*/
class Signal {
class Taoyao {
// 是否连接
connect = false;
// 房间列表
rooms = new Map();
// 回调事件
@@ -588,7 +585,15 @@ class Signal {
}
async mediaConsume(message, body) {
const { roomId, clientId, sourceId, streamId, producerId, transportId, rtpCapabilities } = body;
const {
roomId,
clientId,
sourceId,
streamId,
producerId,
transportId,
rtpCapabilities,
} = body;
const room = this.rooms.get(roomId);
const producer = room.producers.get(producerId);
const transport = room.transports.get(transportId);
@@ -691,20 +696,22 @@ class Signal {
);
});
// TODO改为同步
this.push(protocol.buildMessage("media::consume", {
//await this.request("media::consume", {
kind: consumer.kind,
type: consumer.type,
roomId: roomId,
clientId: clientId,
sourceId: sourceId,
streamId: streamId,
producerId: producerId,
consumerId: consumer.id,
rtpParameters: consumer.rtpParameters,
appData: producer.appData,
producerPaused: consumer.producerPaused,
}));
this.push(
protocol.buildMessage("media::consume", {
//await this.request("media::consume", {
kind: consumer.kind,
type: consumer.type,
roomId: roomId,
clientId: clientId,
sourceId: sourceId,
streamId: streamId,
producerId: producerId,
consumerId: consumer.id,
rtpParameters: consumer.rtpParameters,
appData: producer.appData,
producerPaused: consumer.producerPaused,
})
);
await consumer.resume();
this.push(
protocol.buildMessage("media::consumer::score", {
@@ -883,4 +890,4 @@ class Signal {
}
}
module.exports = { Signal, signalChannel };
module.exports = { Taoyao, signalChannel };