From 263671b4a257a20c6f4e3c1eab2a5065ca27b3a9 Mon Sep 17 00:00:00 2001
From: zengh <123456>
Date: Thu, 12 Aug 2021 16:24:20 +0800
Subject: [PATCH] 视频通话
---
src/main/java/org/springblade/modules/webscoket/service/IPushMsgService.java | 35 +++
src/main/java/org/springblade/modules/webscoket/WebSocketHandler.java | 193 +++++++++++++++++++++
src/main/java/org/springblade/modules/webscoket/index.html | 57 ++++++
src/main/java/org/springblade/modules/webscoket/ChannelSupervise.java | 38 ++++
src/main/java/org/springblade/modules/webscoket/service/impl/PushMsgServiceImpl.java | 65 +++++++
src/main/java/org/springblade/modules/webscoket/WebSocketServer.java | 59 ++++++
src/main/java/org/springblade/modules/webscoket/controller/PushMsgController.java | 59 ++++++
7 files changed, 506 insertions(+), 0 deletions(-)
diff --git a/src/main/java/org/springblade/modules/webscoket/ChannelSupervise.java b/src/main/java/org/springblade/modules/webscoket/ChannelSupervise.java
new file mode 100644
index 0000000..141cff9
--- /dev/null
+++ b/src/main/java/org/springblade/modules/webscoket/ChannelSupervise.java
@@ -0,0 +1,38 @@
+package org.springblade.modules.webscoket;
+import io.netty.channel.Channel;
+import io.netty.channel.ChannelId;
+import io.netty.channel.group.ChannelGroup;
+import io.netty.channel.group.DefaultChannelGroup;
+import io.netty.handler.codec.http.websocketx.TextWebSocketFrame;
+import io.netty.util.concurrent.GlobalEventExecutor;
+
+import java.util.HashMap;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentMap;
+
+public class ChannelSupervise {
+ private static ChannelGroup GlobalGroup=new DefaultChannelGroup(GlobalEventExecutor.INSTANCE);
+ private static ConcurrentMap<String, ChannelId> ChannelMap=new ConcurrentHashMap();
+ private static Map<String, String> map = new HashMap<String, String>();;
+ public static void addChannel(Channel channel,String name){
+ GlobalGroup.add(channel);
+ ChannelMap.put(channel.id().asShortText(),channel.id());
+ map.put(channel.id().asShortText(),name);
+ }
+ public static void removeChannel(Channel channel){
+ GlobalGroup.remove(channel);
+ ChannelMap.remove(channel.id().asShortText());
+ map.remove(channel.id().asShortText());
+ }
+ public static Channel findChannel(String id){
+ return GlobalGroup.find(ChannelMap.get(id));
+ }
+ public static String findName(String id){
+ return map.get(id);
+ }
+ public static void send2All(TextWebSocketFrame tws){
+ GlobalGroup.writeAndFlush(tws);
+ }
+
+}
diff --git a/src/main/java/org/springblade/modules/webscoket/WebSocketHandler.java b/src/main/java/org/springblade/modules/webscoket/WebSocketHandler.java
new file mode 100644
index 0000000..2e4796d
--- /dev/null
+++ b/src/main/java/org/springblade/modules/webscoket/WebSocketHandler.java
@@ -0,0 +1,193 @@
+package org.springblade.modules.webscoket;
+
+import com.alibaba.fastjson.JSON;
+import com.alibaba.fastjson.JSONObject;
+import io.netty.buffer.ByteBuf;
+import io.netty.buffer.Unpooled;
+import io.netty.channel.ChannelFuture;
+import io.netty.channel.ChannelHandlerContext;
+import io.netty.channel.SimpleChannelInboundHandler;
+import io.netty.handler.codec.http.DefaultFullHttpResponse;
+import io.netty.handler.codec.http.FullHttpRequest;
+import io.netty.handler.codec.http.HttpResponseStatus;
+import io.netty.handler.codec.http.HttpVersion;
+import io.netty.handler.codec.http.websocketx.*;
+import io.netty.util.AttributeKey;
+import io.netty.util.CharsetUtil;
+import org.springblade.modules.nettyServer.NettyConfig;
+import org.springblade.modules.suser.service.ISuserService;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Component;
+
+import javax.annotation.PostConstruct;
+
+@Component
+public class WebSocketHandler extends SimpleChannelInboundHandler<Object> {
+
+ private WebSocketServerHandshaker handshaker;
+
+ private String on=null;
+
+ @Autowired
+ private ISuserService suserService;
+
+ private static WebSocketHandler webSocketHandler;
+
+ @PostConstruct
+ public void init() {
+ webSocketHandler = this;
+ }
+
+ @Override
+ protected void channelRead0(ChannelHandlerContext ctx, Object msg) throws Exception {
+ if (msg instanceof FullHttpRequest) {
+ //以http请求形式接入,但是走的是websocket
+ handleHttpRequest(ctx, (FullHttpRequest) msg);
+ } else if (msg instanceof WebSocketFrame) {
+ //处理websocket客户端的消息
+ handlerWebSocketFrame(ctx, (WebSocketFrame) msg);
+ }
+ }
+
+ @Override
+ public void channelActive(ChannelHandlerContext ctx) throws Exception {
+ //添加连接
+ System.out.println("客户端加入连接:" + ctx.channel());
+ //ChannelSupervise.addChannel(ctx.channel());
+ }
+
+ @Override
+ public void channelInactive(ChannelHandlerContext ctx) throws Exception {
+ //断开连接
+ System.out.println("客户端断开连接:" + ctx.channel());
+ //用户离线状态
+ String name = ChannelSupervise.findName(ctx.channel().id().asShortText());
+ System.out.println(name);
+ if ( name != null &&!name.equals("ping") ){
+ NettyConfig.getChannelGroup().remove(ctx.channel());
+ removeUserId(ctx);
+ }
+ }
+
+ @Override
+ public void channelReadComplete(ChannelHandlerContext ctx) throws Exception {
+ ctx.flush();
+ }
+
+ private void handlerWebSocketFrame(ChannelHandlerContext ctx, WebSocketFrame frame) throws NumberFormatException, Exception {
+ // 判断是否关闭链路的指令
+ if (frame instanceof CloseWebSocketFrame) {
+ handshaker.close(ctx.channel(), (CloseWebSocketFrame) frame.retain());
+ return;
+ }
+ // 判断是否ping消息
+ if (frame instanceof PingWebSocketFrame) {
+ ctx.channel().write(
+ new PongWebSocketFrame(frame.content().retain()));
+ return;
+ }
+ // 本例程仅支持文本消息,不支持二进制消息
+ if (!(frame instanceof TextWebSocketFrame)) {
+ System.out.println("本例程仅支持文本消息,不支持二进制消息");
+ throw new UnsupportedOperationException(String.format(
+ "%s frame types not supported", frame.getClass().getName()));
+ }
+ // 返回应答消息
+ String request = ((TextWebSocketFrame) frame).text();
+
+ if (!request.equals("ping")){
+
+ JSONObject jsonObj = JSON.parseObject(request);
+ String type = jsonObj.get("type").toString();
+
+ if (type != null && type.equals("login")){
+ //登录链接
+ String id = jsonObj.get("id").toString();
+ NettyConfig.getUserChannelMap().put(id,ctx.channel());
+
+ System.out.println(jsonObj.get("type"));
+ System.out.println(jsonObj.get("id"));
+
+ //将用户id作为自定义属性加入到channel 中,方便随时channel中获取用户id
+ AttributeKey<String> key = AttributeKey.valueOf("userId");
+ ctx.channel().attr(key).setIfAbsent(id);
+ //把用户信息添加到通道里
+ ChannelSupervise.addChannel(ctx.channel(),id);
+ }
+
+
+
+
+
+
+ }
+
+// TextWebSocketFrame tws = new TextWebSocketFrame(new Date().toString()
+// + ctx.channel().id() + ":" + ctx.channel());
+// // 返回【谁发的发给谁】
+//
+// ctx.channel().writeAndFlush(tws);
+ // 群发
+ // ChannelSupervise.send2All(tws);
+ }
+
+ /**
+ * 唯一的一次http请求,用于创建websocket
+ *
+ * @throws InterruptedException
+ */
+ private void handleHttpRequest(final ChannelHandlerContext ctx,
+ FullHttpRequest req) throws InterruptedException {
+ //要求Upgrade为websocket,过滤掉get/Post
+ if (!req.decoderResult().isSuccess() || (!"websocket".equals(req.headers().get("Upgrade")))) {
+ //若不是websocket方式,则创建BAD_REQUEST的req,返回给客户端
+ sendHttpResponse(ctx, req, new DefaultFullHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.BAD_REQUEST));
+ return;
+ }
+ //握手
+ WebSocketServerHandshakerFactory wsFactory = new WebSocketServerHandshakerFactory(
+ "ws://localhost:9034/websocket", null, false);
+ handshaker = wsFactory.newHandshaker(req);
+ if (handshaker == null) {
+ WebSocketServerHandshakerFactory
+ .sendUnsupportedVersionResponse(ctx.channel());
+ } else {
+ handshaker.handshake(ctx.channel(), req);
+ }
+ }
+
+ /**
+ * 拒绝不合法的请求,并返回错误信息
+ */
+ private static void sendHttpResponse(ChannelHandlerContext ctx,
+ FullHttpRequest req, DefaultFullHttpResponse res) {
+ // 返回应答给客户端
+ if (res.status().code() != 200) {
+ ByteBuf buf = Unpooled.copiedBuffer(res.status().toString(),
+ CharsetUtil.UTF_8);
+ res.content().writeBytes(buf);
+ buf.release();
+ }
+ ChannelFuture f = ctx.channel().writeAndFlush(res);
+ // 如果是非Keep-Alive,关闭连接
+// if (!isKeepAlive(req) || res.status().code() != 200) {
+// f.addListener(ChannelFutureListener.CLOSE);
+// }
+ }
+ @Override
+ public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
+ cause.printStackTrace();
+ ctx.close();
+ }
+
+ /**
+ * 删除用户与channel 对应关系
+ * @param ctx
+ */
+ private void removeUserId(ChannelHandlerContext ctx){
+ AttributeKey<String> key = AttributeKey.valueOf("userId");
+ String userId = ctx.channel().attr(key).get();
+ NettyConfig.getUserChannelMap().remove(userId);
+ }
+
+}
diff --git a/src/main/java/org/springblade/modules/webscoket/WebSocketServer.java b/src/main/java/org/springblade/modules/webscoket/WebSocketServer.java
new file mode 100644
index 0000000..a47840a
--- /dev/null
+++ b/src/main/java/org/springblade/modules/webscoket/WebSocketServer.java
@@ -0,0 +1,59 @@
+package org.springblade.modules.webscoket;
+import io.netty.bootstrap.ServerBootstrap;
+import io.netty.channel.ChannelFuture;
+import io.netty.channel.ChannelInitializer;
+import io.netty.channel.ChannelOption;
+import io.netty.channel.EventLoopGroup;
+import io.netty.channel.nio.NioEventLoopGroup;
+import io.netty.channel.socket.SocketChannel;
+import io.netty.channel.socket.nio.NioServerSocketChannel;
+import io.netty.handler.codec.http.HttpObjectAggregator;
+import io.netty.handler.codec.http.HttpServerCodec;
+import io.netty.handler.logging.LogLevel;
+import io.netty.handler.logging.LoggingHandler;
+import io.netty.handler.stream.ChunkedWriteHandler;
+
+public class WebSocketServer {
+ private int port = 9034;
+
+ public WebSocketServer(int port) {
+ bind(port);
+ }
+
+ public void bind(int port) {
+ Thread thread = new Thread(new Runnable() {
+ @Override
+ public void run() {
+ EventLoopGroup bossGroup = new NioEventLoopGroup();
+ EventLoopGroup workerGroup = new NioEventLoopGroup();
+ try {
+ ServerBootstrap serverBootstrap = new ServerBootstrap();
+ serverBootstrap.group(bossGroup, workerGroup).channel(NioServerSocketChannel.class)
+ .handler(new LoggingHandler(LogLevel.INFO))
+ //保持连接
+ .childOption(ChannelOption.SO_KEEPALIVE, true)
+ .childHandler(new ChannelInitializer<SocketChannel>() {
+ @Override
+ protected void initChannel(SocketChannel ch) {
+// ch.pipeline().addLast("logging",new LoggingHandler("DEBUG"));//设置log监听器,并且日志级别为debug,方便观察运行流程
+ ch.pipeline().addLast("http-codec", new HttpServerCodec());//设置解码器
+ ch.pipeline().addLast("aggregator", new HttpObjectAggregator(65536));//聚合器,使用websocket会用到
+ ch.pipeline().addLast("http-chunked", new ChunkedWriteHandler());//用于大数据的分区传输
+ ch.pipeline().addLast("handler", new WebSocketHandler());//自定义的业务handler
+ }
+ });
+
+ ChannelFuture channelFuture = serverBootstrap.bind(port).sync();
+ System.out.println("WebSocketServer启动成功");
+ channelFuture.channel().closeFuture().sync();
+ } catch (Exception e) {
+ e.printStackTrace();
+ } finally {
+ bossGroup.shutdownGracefully();
+ workerGroup.shutdownGracefully();
+ }
+ }
+ });
+ thread.start();
+ }
+}
diff --git a/src/main/java/org/springblade/modules/webscoket/controller/PushMsgController.java b/src/main/java/org/springblade/modules/webscoket/controller/PushMsgController.java
new file mode 100644
index 0000000..2a3e9a6
--- /dev/null
+++ b/src/main/java/org/springblade/modules/webscoket/controller/PushMsgController.java
@@ -0,0 +1,59 @@
+package org.springblade.modules.webscoket.controller;
+
+import org.springblade.core.tool.api.R;
+import org.springblade.modules.webscoket.service.IPushMsgService;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.web.bind.annotation.GetMapping;
+import org.springframework.web.bind.annotation.PostMapping;
+import org.springframework.web.bind.annotation.RequestMapping;
+import org.springframework.web.bind.annotation.RestController;
+
+import java.util.Date;
+import java.util.HashMap;
+import java.util.Map;
+
+/**
+ * @author lq
+ * @date 2020/4/1 11:22
+ */
+@RestController
+@RequestMapping("pushMsg")
+public class PushMsgController {
+
+ @Autowired
+ private IPushMsgService pushMsgService;
+
+ @PostMapping("/pushUser")
+ public String pushUser(String userId, String msg) {
+ pushMsgService.pushMsg(userId, msg);
+ return "消息发送成功:" + msg;
+ }
+
+ @PostMapping("/pushAll")
+ public String pushAll(String msg) {
+ pushMsgService.pushMsg(msg);
+ return "消息发送成功:" + msg;
+ }
+
+ @GetMapping("/inviteVideoCall")
+ public R<Map> inviteVideoCall(String userId, String type, String name,String faqiid) {
+ //获取当前时间戳作为房间号
+ String roomId = "";
+ Map<String, Object> map = new HashMap<String, Object>();
+ String time = String.valueOf(new Date().getTime());
+ int msg = pushMsgService.inviteVideoCall(userId, time, type, name,faqiid);
+ map.put("type", type);
+ map.put("roomId", time);
+ map.put("res", msg);
+
+ return R.data(map);
+ }
+
+ @GetMapping("/closeVideoCall")
+ public void closeVideoCall(String sentId, String acceptId) {
+
+ pushMsgService.closeVideoCall(sentId, acceptId);
+
+ }
+
+}
diff --git a/src/main/java/org/springblade/modules/webscoket/index.html b/src/main/java/org/springblade/modules/webscoket/index.html
new file mode 100644
index 0000000..aed5d0a
--- /dev/null
+++ b/src/main/java/org/springblade/modules/webscoket/index.html
@@ -0,0 +1,57 @@
+<html>
+<head>
+ <meta http-equiv="Content-Type" content="text/html; charset = utf-8"/>
+ <title>WebSocket客户端</title>
+ <script type="text/javascript">
+ var socket;
+ if(!window.WebSocket){
+ window.WebSocket = window.MozWebSocket;
+ }
+
+ if(window.WebSocket){
+ socket = new WebSocket("ws://localhost:9034/websocket");
+ socket.onmessage = function(event){
+ var ta = document.getElementById('responseContent');
+ ta.value += event.data + "\r\n";
+ };
+
+ socket.onopen = function(event){
+ var ta = document.getElementById('responseContent');
+ ta.value = "你当前的浏览器支持WebSocket,请进行后续操作\r\n";
+ };
+
+ socket.onclose = function(event){
+ var ta = document.getElementById('responseContent');
+ ta.value = "";
+
+ ta.value = "WebSocket连接已经关闭\r\n";
+ };
+ }else{
+ alert("您的浏览器不支持WebSocket");
+ }
+
+
+
+ function send(message){
+ if(!window.WebSocket){
+ return;
+ }
+ if(socket.readyState == WebSocket.OPEN){
+ socket.send(message);
+ }else{
+ alert("WebSocket连接没有建立成功!!");
+ }
+ }
+ </script>
+</head>
+<body>
+<form onSubmit="return false;">
+ <input type = "text" name = "message" value = ""/>
+ <br/><br/>
+ <input type = "button" value = "发送WebSocket请求消息" onClick = "send(this.form.message.value)"/>
+ <hr color="red"/>
+ <h2>客户端接收到服务端返回的应答消息</h2>
+ <textarea id = "responseContent" style = "width:1024px; height:300px"></textarea>
+</form>
+</body>
+</html>
diff --git a/src/main/java/org/springblade/modules/webscoket/service/IPushMsgService.java b/src/main/java/org/springblade/modules/webscoket/service/IPushMsgService.java
new file mode 100644
index 0000000..57c2a3b
--- /dev/null
+++ b/src/main/java/org/springblade/modules/webscoket/service/IPushMsgService.java
@@ -0,0 +1,35 @@
+package org.springblade.modules.webscoket.service;
+
+/**
+ * @author 123456
+ */
+public interface IPushMsgService {
+
+ /**
+ * 给指定用户发送消息
+ *
+ * @param userId
+ * @param msg
+ */
+ void pushMsg(String userId, String msg);
+
+ /**
+ * 给所有用户发送消息
+ *
+ * @param msg
+ */
+ void pushMsg(String msg);
+
+ /**
+ * 给指定用户发起视频邀请
+ */
+ int inviteVideoCall(String userid, String time, String type, String name,String faqiid);
+
+ /**
+ * 关闭视频请求
+ *
+ * @return
+ */
+ void closeVideoCall(String sentId, String acceptId);
+
+}
diff --git a/src/main/java/org/springblade/modules/webscoket/service/impl/PushMsgServiceImpl.java b/src/main/java/org/springblade/modules/webscoket/service/impl/PushMsgServiceImpl.java
new file mode 100644
index 0000000..1d408ad
--- /dev/null
+++ b/src/main/java/org/springblade/modules/webscoket/service/impl/PushMsgServiceImpl.java
@@ -0,0 +1,65 @@
+package org.springblade.modules.webscoket.service.impl;
+
+import com.alibaba.fastjson.JSONObject;
+import io.netty.channel.Channel;
+import io.netty.channel.group.ChannelGroup;
+import io.netty.handler.codec.http.websocketx.TextWebSocketFrame;
+import org.springblade.modules.nettyServer.NettyConfig;
+import org.springblade.modules.webscoket.service.IPushMsgService;
+import org.springframework.stereotype.Service;
+
+/**
+ * @author lq
+ * @date 2020/4/1 11:20
+ */
+@Service
+public class PushMsgServiceImpl implements IPushMsgService {
+ @Override
+ public void pushMsg(String userId, String msg) {
+ Channel channel = NettyConfig.getUserChannelMap().get(userId);
+ if (channel != null) {
+ channel.writeAndFlush(new TextWebSocketFrame(msg));
+ }
+
+ }
+
+ @Override
+ public void pushMsg(String msg) {
+ ChannelGroup group = NettyConfig.getChannelGroup();
+ String name = group.name();
+ System.out.println("空间大小:" + group.size() + ",名字:" + name);
+ group.writeAndFlush(new TextWebSocketFrame(msg));
+ }
+
+ @Override
+ public int inviteVideoCall(String userId, String time, String type, String name,String faqiid) {
+ //返回值
+ int res = 0;
+ Channel channel = NettyConfig.getUserChannelMap().get(userId);
+ JSONObject jsonObject = new JSONObject();
+ jsonObject.put("type", type);
+ jsonObject.put("roomId", time);
+ jsonObject.put("name", name);
+ jsonObject.put("faqiid", faqiid);
+ if (channel != null) {
+ channel.writeAndFlush(new TextWebSocketFrame(String.valueOf(jsonObject)));
+ res = 1;
+ }
+ return res;
+ }
+
+ @Override
+ public void closeVideoCall(String sentId, String acceptId) {
+ Channel channel1 = NettyConfig.getUserChannelMap().get(sentId);
+ Channel channel2 = NettyConfig.getUserChannelMap().get(acceptId);
+ JSONObject jsonObject = new JSONObject();
+ jsonObject.put("type", "close");
+ if (channel1 != null) {
+ channel1.writeAndFlush(new TextWebSocketFrame(String.valueOf(jsonObject)));
+ }
+ if (channel2 != null) {
+ channel2.writeAndFlush(new TextWebSocketFrame(String.valueOf(jsonObject)));
+ }
+
+ }
+}
--
Gitblit v1.9.3