webrtcP2P音视频通话(二)
·
文章目录
学习链接
webrtc实现视频会议,web多人视频通话,websocket通信
- 源码:https://gitee.com/zzhua195/video-meeting-basic
关联:
本地简单示例
客户端代码
client.html
<!DOCTYPE html>
<html lang="en">
<head>
<meta charset="UTF-8">
<meta name="viewport" content="width=device-width, initial-scale=1.0">
<title>Document</title>
<style>
body {
margin: 0;
}
.video-wrapper {
video {
width: 400px;
height: 300px;
border: 1px dashed black;
}
}
</style>
</head>
<body>
<div class="video-wrapper">
<video id="localVideo"></video>
<video id="remoteVideo"></video>
</div>
<div>
目标用户id: <input type="text" id="targetIdIpt" value="2">
<br />
<button id="startCallBtn" onclick="startCall">发起通话</button></button>
</div>
</body>
<script>
let localVideo = document.querySelector('#localVideo')
let remoteVideo = document.querySelector('#remoteVideo')
let targetIdIpt = document.querySelector('#targetIdIpt')
const params = location.href.substring(location.href.indexOf('?') + 1)
const _params = {}
params.split('&').forEach(item => {
const t = item.split('=')
_params[t[0]] = t[1]
})
const { id: currUserId } = _params
startCallBtn.addEventListener('click', startCall)
let localStream = null
let socket = null
let peer = null
async function startCamera() {
localStream = await navigator.mediaDevices.getUserMedia({ video:true, /* audio:true */ })
localVideo.srcObject = localStream
localVideo.play()
}
// 开启摄像头
startCamera()
function initWs() {
socket = new WebSocket('ws://localhost:9095/ws/' + currUserId)
socket.onopen = () => {
console.log('连接建立');
}
socket.onmessage = (e) => {
console.log('收到消息: ', e.data);
let {code, data} = JSON.parse(e.data)
console.log('收到消息的code', code);
if (code == 'connect_success') {
console.log('连接成功');
} else if (code == 'offer') {
acceptCall(data)
} else if (code == 'icecandidate') {
console.log('收到icecandidate消息',peer,data)
let {candidate} = data
peer.addIceCandidate(candidate)
} else if (code == 'answer') {
console.log('对方已同意', data);
acceptAnswer(data)
}
}
socket.onerror = (err) => {
console.log('连接发生错误');
}
socket.onclose = () => {
console.log('连接关闭');
}
}
// 初始化websocket
initWs()
// 发起通话
async function startCall() {
peer = new RTCPeerConnection({})
localStream.getTracks().forEach(track => {
peer.addTrack(track, localStream)
})
peer.addEventListener('track', e => {
console.log('发起方ontrack触发', e);
remoteVideo.srcObject = e.streams[0]
remoteVideo.play()
})
// 这个事件要被触发的前提是, 要先执行上面的把音视频流的轨道给到peer。然后 调用peer.setLocalDescription时会多次连续触发icecandidate事件
peer.addEventListener('icecandidate', e => {
// RTCPeerConnectionIceEvent {isTrusted: true, candidate: RTCIceCandidate, type: 'icecandidate', target: RTCPeerConnection, currentTarget: RTCPeerConnection, …}
console.log('发起方 icecandidate触发', e);
if (e.candidate) {
// e.candidate RTCIceCandidate {candidate: 'candidate:1764731264 1 udp 2122260223 192.168.134.…760 typ host generation 0 ufrag KPfd network-id 1', sdpMid: '1', sdpMLineIndex: 1, foundation: '1764731264', component: 'rtp', …}
console.log('发起方 e.candidate', e.candidate);
// 内网穿透
const message = {
code: 'icecandidate',
data: {
targetId: targetIdIpt.value,
candidate: e.candidate
}
}
socket.send(JSON.stringify(message))
}
})
sendOffer()
}
async function sendOffer() {
// RTCSessionDescription {type: 'offer', sdp: 'v=0\r\no=- 8833742509318497131 2 IN IP4 127.0.0.1\r\ns…770 msid:- 8e3b686c-7779-48a1-8714-8a7a0caf0136\r\n'}
let offer = await peer.createOffer()
console.log('创建了offer', offer);
// 设置offer, 会连续触发多次 peer 的 icecandidat 事件, 最后1次的candidate是null
peer.setLocalDescription(offer)
const message = {
code: 'offer',
data: {
targetId: targetIdIpt.value,
offer: offer
}
}
socket.send(JSON.stringify(message))
}
function acceptAnswer({fromId, answer}) {
console.log('设置answer', answer);
peer.setRemoteDescription(answer)
}
function acceptCall({fromId, offer}) {
console.log('接受通话请求', localStream);
peer = new RTCPeerConnection({})
localStream.getTracks().forEach(track => {
peer.addTrack(track, localStream)
})
peer.addEventListener('track', e => {
console.log('接收方 ontrack触发', e);
remoteVideo.srcObject = e.streams[0]
remoteVideo.play()
})
// 接收方获得offer {type: 'offer', sdp: 'v=0\r\no=- 1381205837170471750 2 IN IP4 127.0.0.1\r\ns…8adcbdb976 87647395-dbcb-4db4-b105-83145f853af1\r\n'}
console.log('接收方获得offer', offer);
peer.addEventListener('icecandidate', e => {
console.log('接收方 icecandidate触发', e);
if (e.candidate) {
console.log('接收方 e.candidate', e.candidate);
// 内网穿透
const message = {
code: 'icecandidate',
data: {
targetId: fromId,
candidate: e.candidate
}
}
socket.send(JSON.stringify(message))
}
})
// 设置远程发过来的offer, 会触发 peer 的 track事件 监听函数
peer.setRemoteDescription(offer)
sendAnswer(fromId)
}
async function sendAnswer(fromId) {
let answer = await peer.createAnswer()
peer.setLocalDescription(answer)
const message = {
code: 'answer',
data: {
targetId: fromId,
answer: answer
}
}
socket.send(JSON.stringify(message))
}
</script>
</html>
服务端代码
SignalWsServer
就是将消息转发给目标对象
@Slf4j
@Component
@ServerEndpoint("/ws/{id}")
public class SignalWsServer {
private static ConcurrentHashMap<String, Session> SESSION_MAP = new ConcurrentHashMap<>();
private Session session;
private String id;
@OnOpen
public void onOpen(@PathParam("id") String id, Session session) {
log.info("连接::onOpen->{}, {}", id, session);
SESSION_MAP.put(id, session);
this.id = id;
this.session = session;
this.session.getUserProperties().put("id", id);
try {
this.session.getBasicRemote().sendText(
JsonUtil.obj2Json(MapBuilder.newHashMap().put("code", "connect_success").build())
);
} catch (IOException e) {
e.printStackTrace();
}
}
@OnMessage
public void onMessage(String msg) throws Exception {
log.info("收到客户端 【{}】 的消息::{}", id, msg);
JSONObject jsonObject = JSON.parseObject(msg);
String code = String.valueOf(jsonObject.get("code"));
if ("offer".equals(code)) {
JSONObject data = jsonObject.getJSONObject("data");
String targetId = data.getString("targetId");
Object offer = data.get("offer");
Session session = SESSION_MAP.get(targetId);
if (session == null) {
log.error("目标用户不存在");
return;
}
Map<String, Object> obj = MapBuilder.newHashMap()
.put("code", "offer")
.put("data", MapBuilder.newHashMap().put("fromId", id).put("offer", offer).build())
.build();
sendMsg(session, obj);
} else if ("icecandidate".equals(code)) {
JSONObject data = jsonObject.getJSONObject("data");
String targetId = data.getString("targetId");
Object candidate = data.get("candidate");
Session session = SESSION_MAP.get(targetId);
if (session == null) {
log.error("目标用户不存在");
return;
}
Map<String, Object> obj = MapBuilder.newHashMap()
.put("code", "icecandidate")
.put("data", MapBuilder.newHashMap().put("fromId", id).put("candidate", candidate).build())
.build();
sendMsg(session, obj);
} else if ("answer".equals(code)) {
JSONObject data = jsonObject.getJSONObject("data");
String targetId = data.getString("targetId");
Object answer = data.get("answer");
Session session = SESSION_MAP.get(targetId);
if (session == null) {
log.error("目标用户不存在");
return;
}
Map<String, Object> obj = MapBuilder.newHashMap()
.put("code", "answer")
.put("data", MapBuilder.newHashMap().put("fromId", id).put("answer", answer).build())
.build();
sendMsg(session, obj);
}
}
private void sendMsg(Session session, Object obj) {
try {
session.getBasicRemote().sendText(JsonUtil.obj2Json(obj));
} catch (IOException e) {
e.printStackTrace();
}
}
@OnClose
public void onClose() {
log.info("关闭::onClose->{}",session);
String id = session.getUserProperties().get("id").toString();
SESSION_MAP.remove(id);
}
@OnError
public void onError(Throwable ex) throws IOException {
log.info("发生错误::onError->{}", ex);
if (this.session.isOpen()) {
session.close();
}
}
}
MapBuilder
public class MapBuilder<K, V> {
private HashMap<K, V> map = new HashMap();
public MapBuilder put(K k, V v) {
map.put(k, v);
return this;
}
public static MapBuilder newHashMap() {
MapBuilder mapBuilder = new MapBuilder();
return mapBuilder;
}
public Map build() {
return map;
}
}
JsonUtil
public class JsonUtil {
static ObjectMapper mapper = new ObjectMapper();
public static String obj2Json(Object obj) {
try {
return mapper.writeValueAsString(obj);
} catch (JsonProcessingException e) {
e.printStackTrace();
return null;
}
}
}
WsConfig
@Configuration
public class WsConfig {
@Bean
public ServerEndpointExporter serverEndpointExporter() {
return new ServerEndpointExporter();
}
}
效果

源代码
把https://gitee.com/zzhua195/video-meeting-basic源代码也贴一下,
图解(先看这个)
先看这个,了解过程,再看代码

前端代码
Single.vue
一对一通信
<template>
<div class="container">
<div>
<video ref="localVideo" autoplay></video>
<video ref="remoteVideo" autoplay></video>
</div>
<div>
<button @click="handleStart">start</button>
</div>
<div>
<input v-model="targetId" type="text">
<button @click="handleCall">call</button>
</div>
</div>
</template>
<script setup>
import { ref } from 'vue'
const params = location.href.substring(location.href.indexOf('?') + 1)
const _params = {}
params.split('&').forEach(item => {
const t = item.split('=')
_params[t[0]] = t[1]
})
const { id: currentUserId } = _params
const localVideo = ref()
const remoteVideo = ref()
let localStream = null
let socket = null
const targetId = ref(null)
let pc = null
async function handleStart(){
localStream = await navigator.mediaDevices.getUserMedia({video: true, audio: true})
localVideo.value.srcObject = localStream
}
function initWebsocket(){
socket = new WebSocket(`wss://192.168.1.8:8080?id=${currentUserId}`)
socket.onopen = e => {
console.log('open_success')
}
socket.onmessage = e => {
const message = e.data
const { code, data } = JSON.parse(message)
if(code === 'connect_success'){
console.log('connect success')
}else if(code === 'offer'){
// 接收offer
getOffer(data)
}else if(code === 'answer'){
// 接收answer
getAnswer(data)
}else if(code === 'icecandidate'){
const { fromId, candidate } = data
pc.addIceCandidate(candidate)
}
}
socket.onerror = e => {
console.error(e.data)
}
}
initWebsocket()
// ====== 发送方 ======= //
function handleCall(){
pc = new RTCPeerConnection({})
localStream.getTracks().forEach(track => pc.addTrack(track, localStream))
pc.addEventListener('track', e => {
remoteVideo.value.srcObject = e.streams[0]
})
pc.addEventListener('icecandidate', e => {
if(e.candidate){
// 内网穿透
const message = {
code: 'icecandidate',
data: {
targetId: targetId.value,
candidate: e.candidate
}
}
socket.send(JSON.stringify(message))
}
})
sendOffer()
}
async function sendOffer(){
const description = await pc.createOffer()
pc.setLocalDescription(description)
const message = {
code: 'offer',
data: {
targetId: targetId.value,
offer: description
}
}
socket.send(JSON.stringify(message))
}
function getAnswer({fromId, answer}){
pc.setRemoteDescription(answer)
}
// ====== 接收方 ======= //
function getOffer({fromId, offer}){
pc = new RTCPeerConnection({})
localStream.getTracks().forEach(track => pc.addTrack(track, localStream))
pc.addEventListener('track', e => {
remoteVideo.value.srcObject = e.streams[0]
})
pc.addEventListener('icecandidate', e => {
if(e.candidate){
// 内网穿透
const message = {
code: 'icecandidate',
data: {
targetId: fromId,
candidate: e.candidate
}
}
socket.send(JSON.stringify(message))
}
})
pc.setRemoteDescription(offer)
sendAnswer(fromId)
}
async function sendAnswer(targetId){
const description = await pc.createAnswer()
pc.setLocalDescription(description)
const message = {
code: 'answer',
data: {
targetId,
answer: description
}
}
socket.send(JSON.stringify(message))
}
</script>
<style scoped>
video {
width: 500px;
height: 400px;
border: 1px dashed black;
}
</style>
Multiple.vue
一对多通信
<template>
<div class="container">
<div ref="videoContainer">
<video ref="localVideo" src="" autoplay></video>
</div>
<div>
<button @click="handleJoin">join</button>
</div>
</div>
</template>
<script setup>
import {ref} from 'vue'
const params = location.href.substring(location.href.indexOf('?') + 1)
const _params = {}
params.split('&').forEach(item => {
const t = item.split('=')
_params[t[0]] = t[1]
})
const {id: currentUserId} = _params
const videoContainer = ref()
let socket = null
let pcMap = new Map()
let localStream = null
const localVideo = ref()
function initWebsocket() {
socket = new WebSocket(`wss://192.168.1.8:8080?id=${currentUserId}`)
socket.onopen = e => {
console.log('open_success')
}
socket.onmessage = e => {
const message = e.data
const {code, data} = JSON.parse(message)
if (code === 'connect_success') {
console.log('connect success')
} else if (code === 'offer') {
// 接收offer
getOffer(data)
} else if (code === 'answer') {
// 接收answer
getAnswer(data)
} else if (code === 'icecandidate') {
const {fromId, candidate} = data
const pc = pcMap.get(fromId)
pc.addIceCandidate(candidate)
} else if (code === 'join_group') {
const {fromId} = data
sendOffer(fromId)
}
}
socket.onerror = e => {
console.error(e.data)
}
}
initWebsocket()
async function handleJoin() {
localStream = await navigator.mediaDevices.getUserMedia({video: true, audio: true})
localVideo.value.srcObject = localStream
const message = {
code: 'join_group',
data: {
groupId: 1
}
}
socket.send(JSON.stringify(message))
}
async function sendOffer(targetId) {
const pc = new RTCPeerConnection({})
pcMap.set(targetId, pc)
localStream.getTracks().forEach(track => pc.addTrack(track, localStream))
const video = document.createElement('video')
video.autoplay = true
videoContainer.value.appendChild(video)
pc.addEventListener('track', e => {
// 更新远程的视频
video.srcObject = e.streams[0]
})
pc.addEventListener('icecandidate', e => {
if (e.candidate) {
const message = {
code: 'icecandidate',
data: {
targetId,
candidate: e.candidate
}
}
socket.send(JSON.stringify(message))
}
})
const description = await pc.createOffer()
await pc.setLocalDescription(description)
const message = {
code: 'offer',
data: {
targetId,
offer: description
}
}
socket.send(JSON.stringify(message))
}
// 接收offer
async function getOffer({fromId: targetId, offer}) {
const pc = new RTCPeerConnection({})
pcMap.set(targetId, pc)
localStream.getTracks().forEach(track => pc.addTrack(track, localStream))
const video = document.createElement('video')
video.autoplay = true
videoContainer.value.appendChild(video)
pc.addEventListener('track', e => {
// 更新远程的视频
video.srcObject = e.streams[0]
})
pc.addEventListener('icecandidate', e => {
if (e.candidate) {
const message = {
code: 'icecandidate',
data: {
targetId,
candidate: e.candidate
}
}
socket.send(JSON.stringify(message))
}
})
await pc.setRemoteDescription(offer)
const description = await pc.createAnswer()
await pc.setLocalDescription(description)
const message = {
code: 'answer',
data: {
targetId,
answer: description
}
}
socket.send(JSON.stringify(message))
}
async function getAnswer({fromId, answer}){
const pc = pcMap.get(fromId)
await pc.setRemoteDescription(answer)
}
</script>
<style>
video {
width: 500px;
height: 400px;
border: 1px dashed black;
}
</style>
后端代码
websocket.js - node
// https://www.npmjs.com/package/ws?activeTab=readme
const {WebSocketServer} = require('ws')
const {createServer} = require('https')
const {readFileSync} = require('fs')
const server = createServer({
cert: readFileSync('./cert/server.pem'),
key: readFileSync('./cert/server.key')
});
const wss = new WebSocketServer({server});
// 在线人员
const groupOnlineMembers = new Set()
wss.on('connection', function connection(ws, request) {
const _params = {}
const params = request.url.substring(request.url.indexOf('?') + 1)
params.split('&').forEach(item => {
const t = item.split('=')
_params[t[0]] = t[1]
})
const {id: currentUserId} = _params
ws._userId = currentUserId
ws.on('error', console.error);
ws.on('close', e => {
groupOnlineMembers.delete(currentUserId)
})
ws.on('message', function message(res) {
const message = res.toString()
const {code, data} = JSON.parse(message)
if (code === 'offer') {
const {targetId, offer} = data
wss.clients.forEach(client => {
if (client._userId === targetId) {
const message = {
code: 'offer',
data: {
fromId: currentUserId,
offer
}
}
client.send(JSON.stringify(message))
}
})
} else if (code === 'answer') {
// 转发answer
const {targetId, answer} = data
wss.clients.forEach(client => {
if (client._userId === targetId) {
const message = {
code: 'answer',
data: {
fromId: currentUserId,
answer
}
}
client.send(JSON.stringify(message))
}
})
} else if (code === 'icecandidate') {
const {targetId, candidate} = data
wss.clients.forEach(client => {
if (client._userId === targetId) {
const message = {
code: 'icecandidate',
data: {
fromId: currentUserId,
candidate
}
}
client.send(JSON.stringify(message))
}
})
} else if (code === 'join_group') {
const { groupId } = data
groupOnlineMembers.add(currentUserId)
wss.clients.forEach(client => {
if (groupOnlineMembers.has(client._userId) && client._userId !== currentUserId) {
// 该用户在线
const message = {
code: 'join_group',
data: {
fromId: currentUserId
}
}
client.send(JSON.stringify(message))
}
})
}
});
const message = {
code: 'connect_success',
data: ''
}
ws.send(JSON.stringify(message));
});
server.listen(8080);
WebSocket - java
@Component
@ServerEndpoint("/websocket/{appId}")
@Slf4j
public class WebSocket {
private String appId;
private static final ConcurrentHashMap<String, Session> sessionPool = new ConcurrentHashMap<String, Session>();
private static final Set<String> groupList = new HashSet<>();
/**
* 链接成功调用的方法
*/
@OnOpen
public void onOpen(Session session, @PathParam(value = "appId") String appId) {
try {
log.debug(appId);
this.appId = appId;
sessionPool.put(appId, session);
session.getId();
log.debug("【websocket消息】有新的连接,总数为:" + sessionPool.size());
JSONObject result = new JSONObject();
result.put("code", "connect_success");
result.put("data", "success");
this.sendMessage(appId, result.toJSONString());
} catch (Exception e) {
e.printStackTrace();
log.error(e.getMessage());
}
}
/**
* 链接关闭调用的方法
*/
@OnClose
public void onClose(Session session, CloseReason closeReason) {
try {
sessionPool.remove(this.appId);
groupList.remove(this.appId);
log.debug("【websocket消息】连接断开,总数为:" + sessionPool.size());
} catch (Exception e) {
e.printStackTrace();
log.error(e.getMessage());
}
}
/**
* 收到客户端消息后调用的方法
*/
@OnMessage
public void onMessage(String message, Session session) {
log.debug(appId);
log.debug("【websocket消息】收到客户端消息:" + message);
if (ObjectUtil.isNotEmpty(message)) {
JSONObject json = JSON.parseObject(message);
String code = json.getString("code");
if ("offer".equals(code)) {
// 发送offer
JSONObject data = json.getJSONObject("data");
String targetId = data.getString("targetId");
JSONObject description = data.getJSONObject("description");
JSONObject response = new JSONObject();
response.put("code", "offer");
JSONObject result = new JSONObject();
result.put("fromId", this.appId);
result.put("description", description);
response.put("data", result);
this.sendMessage(targetId, response.toJSONString());
} else if ("answer".equals(code)) {
// 发送answer
JSONObject data = json.getJSONObject("data");
String targetId = data.getString("targetId");
JSONObject description = data.getJSONObject("description");
JSONObject response = new JSONObject();
response.put("code", "answer");
response.put("errcode", "0");
JSONObject result = new JSONObject();
result.put("fromId", this.appId);
result.put("description", description);
response.put("data", result);
this.sendMessage(targetId, response.toJSONString());
} else if ("candidate".equals(code)) {
// 发送candidate
JSONObject data = json.getJSONObject("data");
String targetId = data.getString("targetId");
JSONObject candidate = data.getJSONObject("candidate");
JSONObject response = new JSONObject();
response.put("code", "candidate");
JSONObject result = new JSONObject();
result.put("fromId", this.appId);
result.put("candidate", candidate);
response.put("data", result);
this.sendMessage(targetId, response.toJSONString());
} else if ("join_group_video".equals(code)) {
groupList.add(this.appId);
JSONObject data = json.getJSONObject("data");
String groupId = data.getString("targetId");
JSONObject response = new JSONObject();
response.put("code", "join_group_video");
JSONObject result = new JSONObject();
result.put("fromId", this.appId);
response.put("data", result);
this.sendMoreMessage(groupList, response.toJSONString());
}
}
}
/**
* 发送错误时的处理
*/
@OnError
public void onError(Session session, Throwable error) {
log.info("用户错误,原因:" + error.getMessage());
error.printStackTrace();
}
// 此为单点消息
public void sendMessage(String appId, String message) {
Session session = sessionPool.get(appId);
if(session == null){
return;
}
synchronized(session){
if (session.isOpen()) {
try {
log.info("【websocket消息】 单点消息:" + message);
session.getAsyncRemote().sendText(message);
} catch (Exception e) {
e.printStackTrace();
log.error(e.getMessage());
}
}
}
}
// 此为单点消息(多人)
public void sendMoreMessage(Set<String> appIds, String message) {
for (String appId : appIds) {
if(appId.equals(this.appId)){
continue;
}
Session session = sessionPool.get(appId);
if (session != null && session.isOpen()) {
try {
log.info(appId + "【websocket消息】 单点消息:" + message);
session.getAsyncRemote().sendText(message);
} catch (Exception e) {
e.printStackTrace();
}
}
}
}
}
更多推荐



所有评论(0)