From eb53489e9035e71ba8cec446751303b864160346 Mon Sep 17 00:00:00 2001 From: qiminbao Date: Tue, 21 Oct 2025 14:21:54 +0800 Subject: [PATCH] =?UTF-8?q?v1.1=E5=88=9D=E5=A7=8B=E5=8C=96?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../agent/server/controller/TcpService.java | 171 +++++----- .../agent/server/core/enums/MsgEnum.java | 7 + .../agent/server/netty/AgentNettyConfig.java | 22 ++ .../agent/server/netty/AgentNettyServer.java | 56 ++++ .../server/netty/MultiTargetNettyClient.java | 301 ------------------ .../server/netty/enpoint/AgentEndpoint.java | 75 ++++- .../netty/handler/TCPListenHandler.java | 66 ++++ .../netty/handler/UDPListenHandler.java | 24 ++ .../listener/AgentConsumerDownListener.java | 128 +------- .../server/service/impl/AgentServiceImpl.java | 38 ++- .../agent/server/utils/ClientExample.java | 2 +- .../utils/EnhancedConnectionManager.java | 73 +++-- src/main/resources/application-dev.yml | 21 +- src/main/resources/application.yml | 8 +- 14 files changed, 388 insertions(+), 604 deletions(-) create mode 100644 src/main/java/com/tongran/agent/server/netty/AgentNettyConfig.java create mode 100644 src/main/java/com/tongran/agent/server/netty/AgentNettyServer.java delete mode 100644 src/main/java/com/tongran/agent/server/netty/MultiTargetNettyClient.java create mode 100644 src/main/java/com/tongran/agent/server/netty/handler/TCPListenHandler.java create mode 100644 src/main/java/com/tongran/agent/server/netty/handler/UDPListenHandler.java diff --git a/src/main/java/com/tongran/agent/server/controller/TcpService.java b/src/main/java/com/tongran/agent/server/controller/TcpService.java index be24582..ccc8642 100644 --- a/src/main/java/com/tongran/agent/server/controller/TcpService.java +++ b/src/main/java/com/tongran/agent/server/controller/TcpService.java @@ -1,24 +1,13 @@ package com.tongran.agent.server.controller; -import cn.hutool.core.util.ObjectUtil; -import com.alibaba.fastjson2.JSON; -import com.alibaba.fastjson2.JSONObject; -import com.tongran.agent.server.core.enums.MsgEnum; import com.tongran.agent.server.core.session.SessionManager; -import com.tongran.agent.server.netty.MultiTargetNettyClient; -import com.tongran.agent.server.netty.annotation.AgentDispatcher; import com.tongran.agent.server.netty.basics.AgentDispatcherManager; -import com.tongran.agent.server.netty.basics.AgentHandler; -import com.tongran.agent.server.netty.model.Message; import com.tongran.agent.server.rockermq.AgentProducer; -import com.tongran.agent.server.rockermq.model.MqMsg; import com.tongran.agent.server.utils.AssertLog; import lombok.extern.slf4j.Slf4j; -import org.apache.commons.lang3.StringUtils; import org.springframework.stereotype.Service; import javax.annotation.Resource; -import java.util.Objects; @Slf4j @Service @@ -28,8 +17,8 @@ public class TcpService { @Resource private AgentDispatcherManager agentDispatcherManager; - @Resource - private MultiTargetNettyClient client; +// @Resource +// private MultiTargetNettyClient client; @Resource private AgentProducer agentProducer; @@ -41,88 +30,88 @@ public class TcpService { public void create(String clientId){ boolean success1 = false; // 动态创建多个不同IP和端口的连接 - if(StringUtils.equals(clientId, "clientId007")){ - success1 = client.createConnection(clientId, "127.0.0.1", 6610, 5); - }else{ - success1 = client.createConnection(clientId, "127.0.0.1", 6710, 5); - } - System.out.println("当前活跃连接数: " + client.getActiveConnections()); - if(success1){ - // 连接建立后可以发送登录认证消息 - long timestamp = System.currentTimeMillis(); - JSONObject object = new JSONObject(); - object.put("uuid",clientId); - object.put("timestamp",timestamp); - Message message = Message.builder().clientId(clientId).dataType("CREATE").data(object.toString()).build(); - // 将对象转为 JSON 字符串 - String json = JSON.toJSONString(message); - System.out.println("连接建立成功,发送登录认证="+json); - client.sendMessages(clientId, message); - } +// if(StringUtils.equals(clientId, "clientId007")){ +// success1 = client.createConnection(clientId, "127.0.0.1", 6610, 5); +// }else{ +// success1 = client.createConnection(clientId, "127.0.0.1", 6710, 5); +// } +// System.out.println("当前活跃连接数: " + client.getActiveConnections()); +// if(success1){ +// // 连接建立后可以发送登录认证消息 +// long timestamp = System.currentTimeMillis(); +// JSONObject object = new JSONObject(); +// object.put("uuid",clientId); +// object.put("timestamp",timestamp); +// Message message = Message.builder().clientId(clientId).dataType("CREATE").data(object.toString()).build(); +// // 将对象转为 JSON 字符串 +// String json = JSON.toJSONString(message); +// System.out.println("连接建立成功,发送登录认证="+json); +// client.sendMessages(clientId, message); +// } } public void sendMessage(String message){ AssertLog.info("consumer==> received down message: {}", message); - try { - Message dto = JSON.parseObject(message, Message.class); - if(StringUtils.equals(dto.getDataType(), MsgEnum.注册.getValue())){ - JSONObject jsonObject = JSONObject.parseObject(dto.getData()); - String clientIp = ""; - int clientPort = 0; - String switchBoard = ""; - long timestamp = 0; - if(jsonObject.containsKey("clientIp")){ - clientIp = jsonObject.getString("clientIp"); - } - if(jsonObject.containsKey("clientPort")){ - clientPort = jsonObject.getInteger("clientPort"); - } - if(jsonObject.containsKey("switchBoard")){ - switchBoard = jsonObject.getString("switchBoard"); - } - if(jsonObject.containsKey("timestamp")){ - timestamp = jsonObject.getLongValue("timestamp"); - } - if(StringUtils.isBlank(clientIp) || clientPort == 0){ - //返回提醒ip 端口 空 - JSONObject json = new JSONObject(); - json.put("clientId",dto.getClientId()); - json.put("dataType",MsgEnum.注册应答.getValue()); - JSONObject dataJson = new JSONObject(); - dataJson.put("resCode",0); - dataJson.put("resMag","客户端IP、端口不能为空"); - dataJson.put("timestamp", System.currentTimeMillis()); - json.put("data",dataJson.toString()); - MqMsg mqMsg = MqMsg.builder().topic("tongran_agent_down").content(json.toString()).build(); - agentProducer.asyncSend(mqMsg); - return; - } - boolean success = client.createConnection(dto.getClientId(), clientIp, clientPort, 5); - if(success){ - //连接成功 - // 连接建立后可以发送登录认证消息 - JSONObject object = new JSONObject(); - object.put("timestamp", timestamp); - if(StringUtils.isNotBlank(switchBoard)){ - object.put("switchBoard",switchBoard); - } - Message msg = Message.builder().clientId(dto.getClientId()).dataType(MsgEnum.注册.getValue()).data(object.toString()).build(); - // 将对象转为 JSON 字符串 - client.sendMessages(dto.getClientId(), msg); - } - } else { - AgentHandler msgHandler = agentDispatcherManager.getHandler(dto.getDataType() + "&" + AgentDispatcher.VersionEnum.V1.value); - if (ObjectUtil.isNotEmpty(msgHandler)) { - String msg = msgHandler.downHandle(dto.getData(), dto.getClientId()); - Message chargeMessage = Message.builder().clientId(dto.getClientId()).dataType(dto.getDataType()).data(msg).build(); - if (Objects.nonNull(sessionManager.getSessionById(dto.getClientId()))) { - sessionManager.writeAndFlush(sessionManager.getSessionById(dto.getClientId()).getChannel(), chargeMessage); - } - } - } - } catch (Exception e) { - AssertLog.error("=====rocketmq-exception{}=====" + e.getMessage()); - } +// try { +// Message dto = JSON.parseObject(message, Message.class); +// if(StringUtils.equals(dto.getDataType(), MsgEnum.注册.getValue())){ +// JSONObject jsonObject = JSONObject.parseObject(dto.getData()); +// String clientIp = ""; +// int clientPort = 0; +// String switchBoard = ""; +// long timestamp = 0; +// if(jsonObject.containsKey("clientIp")){ +// clientIp = jsonObject.getString("clientIp"); +// } +// if(jsonObject.containsKey("clientPort")){ +// clientPort = jsonObject.getInteger("clientPort"); +// } +// if(jsonObject.containsKey("switchBoard")){ +// switchBoard = jsonObject.getString("switchBoard"); +// } +// if(jsonObject.containsKey("timestamp")){ +// timestamp = jsonObject.getLongValue("timestamp"); +// } +// if(StringUtils.isBlank(clientIp) || clientPort == 0){ +// //返回提醒ip 端口 空 +// JSONObject json = new JSONObject(); +// json.put("clientId",dto.getClientId()); +// json.put("dataType",MsgEnum.注册应答.getValue()); +// JSONObject dataJson = new JSONObject(); +// dataJson.put("resCode",0); +// dataJson.put("resMag","客户端IP、端口不能为空"); +// dataJson.put("timestamp", System.currentTimeMillis()); +// json.put("data",dataJson.toString()); +// MqMsg mqMsg = MqMsg.builder().topic("tongran_agent_down").content(json.toString()).build(); +// agentProducer.asyncSend(mqMsg); +// return; +// } +// boolean success = client.createConnection(dto.getClientId(), clientIp, clientPort, 5); +// if(success){ +// //连接成功 +// // 连接建立后可以发送登录认证消息 +// JSONObject object = new JSONObject(); +// object.put("timestamp", timestamp); +// if(StringUtils.isNotBlank(switchBoard)){ +// object.put("switchBoard",switchBoard); +// } +// Message msg = Message.builder().clientId(dto.getClientId()).dataType(MsgEnum.注册.getValue()).data(object.toString()).build(); +// // 将对象转为 JSON 字符串 +// client.sendMessages(dto.getClientId(), msg); +// } +// } else { +// AgentHandler msgHandler = agentDispatcherManager.getHandler(dto.getDataType() + "&" + AgentDispatcher.VersionEnum.V1.value); +// if (ObjectUtil.isNotEmpty(msgHandler)) { +// String msg = msgHandler.downHandle(dto.getData(), dto.getClientId()); +// Message chargeMessage = Message.builder().clientId(dto.getClientId()).dataType(dto.getDataType()).data(msg).build(); +// if (Objects.nonNull(sessionManager.getSessionById(dto.getClientId()))) { +// sessionManager.writeAndFlush(sessionManager.getSessionById(dto.getClientId()).getChannel(), chargeMessage); +// } +// } +// } +// } catch (Exception e) { +// AssertLog.error("=====rocketmq-exception{}=====" + e.getMessage()); +// } } diff --git a/src/main/java/com/tongran/agent/server/core/enums/MsgEnum.java b/src/main/java/com/tongran/agent/server/core/enums/MsgEnum.java index ba39c1b..f254fad 100644 --- a/src/main/java/com/tongran/agent/server/core/enums/MsgEnum.java +++ b/src/main/java/com/tongran/agent/server/core/enums/MsgEnum.java @@ -11,11 +11,18 @@ import lombok.NoArgsConstructor; @AllArgsConstructor @NoArgsConstructor public enum MsgEnum { + 建立连接("CONNECT"), + + 建立连接应答("CONNECT_RSP"), 注册("REGISTER"), 注册应答("REGISTER_RSP"), + 获取最新策略("GET_POLICY"), + + 获取最新策略应答("GET_POLICY_RSP"), + 断开("DISCONNECT"), 断开应答("DISCONNECT_RSP"), diff --git a/src/main/java/com/tongran/agent/server/netty/AgentNettyConfig.java b/src/main/java/com/tongran/agent/server/netty/AgentNettyConfig.java new file mode 100644 index 0000000..208b159 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/netty/AgentNettyConfig.java @@ -0,0 +1,22 @@ +package com.tongran.agent.server.netty; + +import com.tongran.agent.server.netty.config.BaseNettyConfig; +import lombok.Data; +import lombok.EqualsAndHashCode; +import lombok.ToString; +import org.springframework.boot.context.properties.ConfigurationProperties; +import org.springframework.context.annotation.Configuration; + +/** + * Netty配置 + */ +@Data +@Configuration +@ToString(callSuper = true) +@EqualsAndHashCode(callSuper = true) +@ConfigurationProperties(prefix = "tcp.netty.charge") +public class AgentNettyConfig extends BaseNettyConfig { + + private static final long serialVersionUID = -4284214721383107291L; + +} diff --git a/src/main/java/com/tongran/agent/server/netty/AgentNettyServer.java b/src/main/java/com/tongran/agent/server/netty/AgentNettyServer.java new file mode 100644 index 0000000..2403486 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/netty/AgentNettyServer.java @@ -0,0 +1,56 @@ +package com.tongran.agent.server.netty; + +import com.tongran.agent.server.netty.handler.AgentDecoderHandler; +import com.tongran.agent.server.netty.handler.AgentDispatcherHandler; +import com.tongran.agent.server.netty.handler.AgentEncoderHandler; +import com.tongran.agent.server.netty.handler.TCPListenHandler; +import io.netty.channel.ChannelInitializer; +import io.netty.channel.socket.nio.NioSocketChannel; +import io.netty.handler.timeout.IdleStateHandler; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; + +import javax.annotation.Resource; + +@Configuration +public class AgentNettyServer extends BaseNettyServer { + + @Resource + private AgentNettyConfig config; + + @Resource + private TCPListenHandler tcpListenHandler; + + @Resource + private AgentDecoderHandler decoderHandler; + + @Resource + private AgentEncoderHandler encoderHandler; + + @Resource + private AgentDispatcherHandler dispatcherHandler; + + + protected AgentNettyServer(AgentNettyConfig config) { + super(config); + } + + @Bean(name = "NettyCharge", initMethod = "start", destroyMethod = "stop") + BaseNettyServer run() { + config.setHander(new ChannelInitializer() { + @Override + public void initChannel(NioSocketChannel ch) throws Exception { + ch.pipeline().addLast(new IdleStateHandler(config.readerIdleTime, config.writerIdleTime, config.allIdleTime));// 心跳 + ch.pipeline().addLast(tcpListenHandler);// 监听器 + //入栈 + ch.pipeline().addLast(decoderHandler);//解码器 + //出栈 + ch.pipeline().addLast(encoderHandler);//加码器 + // 业务分发是入站最后一个处理器, 出站第一个处理器, 位置要放在最后 + ch.pipeline().addLast(businessGroup, dispatcherHandler);//业务分发 这里使用业务线程组 + } + }); + return new AgentNettyServer(config); + } + +} \ No newline at end of file diff --git a/src/main/java/com/tongran/agent/server/netty/MultiTargetNettyClient.java b/src/main/java/com/tongran/agent/server/netty/MultiTargetNettyClient.java deleted file mode 100644 index 4cf8c3f..0000000 --- a/src/main/java/com/tongran/agent/server/netty/MultiTargetNettyClient.java +++ /dev/null @@ -1,301 +0,0 @@ -package com.tongran.agent.server.netty; - -import com.tongran.agent.server.netty.handler.AgentDecoderHandler; -import com.tongran.agent.server.netty.handler.AgentDispatcherHandler; -import com.tongran.agent.server.netty.handler.AgentEncoderHandler; -import com.tongran.agent.server.netty.model.Message; -import io.netty.bootstrap.Bootstrap; -import io.netty.channel.*; -import io.netty.channel.nio.NioEventLoopGroup; -import io.netty.channel.socket.SocketChannel; -import io.netty.channel.socket.nio.NioSocketChannel; -import io.netty.util.AttributeKey; -import org.springframework.stereotype.Component; - -import javax.annotation.Resource; -import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.TimeUnit; - -@Component -public class MultiTargetNettyClient { - // 连接标识属性键 - private static final AttributeKey CONNECTION_KEY = AttributeKey.valueOf("connectionKey"); - - // 连接管理器 - private final Map connectionMap = new ConcurrentHashMap<>(); - private final Map reverseConnectionMap = new ConcurrentHashMap<>(); - private EventLoopGroup workerGroup; - - @Resource - private AgentDecoderHandler decoderHandler; - - @Resource - private AgentEncoderHandler encoderHandler; - - @Resource - private AgentDispatcherHandler dispatcherHandler; - - /** - * 动态创建TCP连接 - * @param connectionKey 连接唯一标识(格式:ip:port 或自定义) - * @param host 目标主机 - * @param port 目标端口 - * @param timeout 超时时间(秒) - * @return 是否创建成功 - */ - public boolean createConnection(String connectionKey, String host, int port, int timeout) { - if (workerGroup == null) { - workerGroup = new NioEventLoopGroup(); - } - - if (connectionMap.containsKey(connectionKey)) { - System.out.println("连接已存在: " + connectionKey); - return true; - } - - Bootstrap bootstrap = new Bootstrap(); - bootstrap.group(workerGroup) - .channel(NioSocketChannel.class) - .option(ChannelOption.TCP_NODELAY, true) - .option(ChannelOption.CONNECT_TIMEOUT_MILLIS, timeout * 1000) - .handler(new ChannelInitializer() { - @Override - protected void initChannel(SocketChannel ch) { - ChannelPipeline pipeline = ch.pipeline(); - // 添加编解码器 - pipeline.addLast(decoderHandler); - pipeline.addLast(encoderHandler); - // 添加心跳机制 -// pipeline.addLast(new IdleStateHandler(0, 30, 0, TimeUnit.SECONDS)); - // 添加自定义处理器 - pipeline.addLast(dispatcherHandler); - } - }); - - try { - CountDownLatch latch = new CountDownLatch(1); - final boolean[] success = {false}; - - bootstrap.connect(host, port).addListener((ChannelFuture future) -> { - if (future.isSuccess()) { - Channel channel = future.channel(); - connectionMap.put(connectionKey, channel); - reverseConnectionMap.put(channel, connectionKey); - success[0] = true; - System.out.println("连接建立成功: " + connectionKey + " -> " + host + ":" + port); - } else { - System.err.println("连接建立失败: " + connectionKey + " - " + future.cause().getMessage()); - success[0] = false; - } - latch.countDown(); - }); - - // 等待连接结果 - return latch.await(timeout, TimeUnit.SECONDS) && success[0]; - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - return false; - } - } - - /** - * 向指定连接发送消息 - * @param connectionKey 连接标识 - * @param message 消息内容 - * @return 是否发送成功 - */ - public boolean sendMessage(String connectionKey, String message) { - Channel channel = connectionMap.get(connectionKey); - if (channel == null) { - System.err.println("连接不存在: " + connectionKey); - return false; - } - - if (!channel.isActive()) { - System.err.println("连接已断开: " + connectionKey); - removeConnection(connectionKey); - return false; - } - - try { - ChannelFuture future = channel.writeAndFlush(message).sync(); - System.out.println("消息发送成功到: " + connectionKey + " - " + message); - return future.isSuccess(); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - return false; - } - } - - public boolean sendMessages(String connectionKey, Message message) { - Channel channel = connectionMap.get(connectionKey); - if (channel == null) { - System.err.println("连接不存在: " + connectionKey); - return false; - } - - if (!channel.isActive()) { - System.err.println("连接已断开: " + connectionKey); - removeConnection(connectionKey); - return false; - } - - try { - ChannelFuture future = channel.writeAndFlush(message).sync(); - System.out.println("消息发送成功到: " + connectionKey + " - " + message); - return future.isSuccess(); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - return false; - } - } - - /** - * 批量发送消息到不同连接 - * @param messages Map<连接标识, 消息内容> - */ - public void sendMessages(Map messages) { - messages.forEach((connectionKey, message) -> { - new Thread(() -> { - if (sendMessage(connectionKey, message)) { - System.out.println("批量发送成功: " + connectionKey); - } else { - System.err.println("批量发送失败: " + connectionKey); - } - }).start(); - }); - } - - /** - * 关闭指定连接 - * @param connectionKey 连接标识 - */ - public void closeConnection(String connectionKey) { - Channel channel = connectionMap.get(connectionKey); - if (channel != null) { - channel.close(); - removeConnection(connectionKey); - System.out.println("连接已关闭: " + connectionKey); - } - } - - /** - * 移除连接记录 - * @param connectionKey 连接标识 - */ - private void removeConnection(String connectionKey) { - Channel channel = connectionMap.remove(connectionKey); - if (channel != null) { - reverseConnectionMap.remove(channel); - } - } - - /** - * 关闭所有连接 - */ - public void shutdown() { - // 关闭所有连接 - connectionMap.forEach((key, channel) -> { - if (channel.isActive()) { - channel.close(); - } - }); - - connectionMap.clear(); - reverseConnectionMap.clear(); - - // 关闭EventLoopGroup - if (workerGroup != null) { - workerGroup.shutdownGracefully(); - } - System.out.println("所有连接已关闭"); - } - - /** - * 获取连接状态 - * @return 当前活跃连接数 - */ - public int getActiveConnections() { - return (int) connectionMap.values().stream() - .filter(Channel::isActive) - .count(); - } - -// /** -// * 客户端处理器 -// */ -// private class ClientHandler extends SimpleChannelInboundHandler { -// @Override -// protected void channelRead0(ChannelHandlerContext ctx, String msg) { -// String connectionKey = ctx.channel().attr(CONNECTION_KEY).get(); -// System.out.println("收到来自 " + connectionKey + " 的消息: " + msg); -// } -// -// @Override -// public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { -// String connectionKey = ctx.channel().attr(CONNECTION_KEY).get(); -// System.err.println("连接 " + connectionKey + " 发生异常: " + cause.getMessage()); -// ctx.close(); -// } -// -// @Override -// public void channelInactive(ChannelHandlerContext ctx) { -// String connectionKey = ctx.channel().attr(CONNECTION_KEY).get(); -// System.out.println("连接断开: " + connectionKey); -// removeConnection(connectionKey); -// } -// } - - /** - * 使用示例 - */ - public static void main(String[] args) throws InterruptedException { - MultiTargetNettyClient client = new MultiTargetNettyClient(); - - try { - // 动态创建多个不同IP和端口的连接 - boolean success1 = client.createConnection("server1", "127.0.0.1", 8080, 5); - boolean success2 = client.createConnection("server2", "192.168.1.100", 8081, 5); - boolean success3 = client.createConnection("server3", "10.0.0.50", 8082, 5); - - // 等待连接建立 - Thread.sleep(2000); - - if (success1) { - client.sendMessage("server1", "Hello Server 1 from 8080"); - } - - if (success2) { - client.sendMessage("server2", "Hello Server 2 from 8081"); - } - - if (success3) { - client.sendMessage("server3", "Hello Server 3 from 8082"); - } - -// // 批量发送示例 -// Map batchMessages = Map.of( -// "server1", "Batch message to server 1", -// "server2", "Batch message to server 2", -// "server3", "Batch message to server 3" -// ); -// -// client.sendMessages(batchMessages); - - // 保持运行一段时间 - Thread.sleep(5000); - - System.out.println("当前活跃连接数: " + client.getActiveConnections()); - - } finally { - client.shutdown(); - } - } - - public Channel get(String connectionKey){ - Channel channel = connectionMap.get(connectionKey); - return channel; - } -} \ No newline at end of file diff --git a/src/main/java/com/tongran/agent/server/netty/enpoint/AgentEndpoint.java b/src/main/java/com/tongran/agent/server/netty/enpoint/AgentEndpoint.java index e63bfe9..d767d9b 100644 --- a/src/main/java/com/tongran/agent/server/netty/enpoint/AgentEndpoint.java +++ b/src/main/java/com/tongran/agent/server/netty/enpoint/AgentEndpoint.java @@ -2,7 +2,6 @@ package com.tongran.agent.server.netty.enpoint; import com.alibaba.fastjson2.JSONObject; import com.tongran.agent.server.core.enums.MsgEnum; -import com.tongran.agent.server.netty.MultiTargetNettyClient; import com.tongran.agent.server.netty.annotation.AgentDispatcher; import com.tongran.agent.server.netty.basics.AgentHandler; import com.tongran.agent.server.netty.model.UpMsgResponse; @@ -15,8 +14,6 @@ import javax.annotation.Resource; @Component public class AgentEndpoint { - @Resource - private MultiTargetNettyClient client; @Resource private AgentMqConfig agentMqConfig; @@ -24,20 +21,66 @@ public class AgentEndpoint { @Resource private AgentProducer agentProducer; - @AgentDispatcher(msgId = MsgEnum.注册应答) - public class RegisterRspHandler implements AgentHandler { + @AgentDispatcher(msgId = MsgEnum.建立连接) + public class ConnectHandler implements AgentHandler { + @Override + public UpMsgResponse upHandle(String data, String clientId) { + JSONObject json = new JSONObject(); + json.put("resCode",1); + return UpMsgResponse.builder().clientId(clientId).dataType(MsgEnum.建立连接应答.getValue()) + .content(json.toString()).build(); + } + } + + @AgentDispatcher(msgId = MsgEnum.建立连接应答) + public class ConnectRspHandler implements AgentHandler { + @Override + public UpMsgResponse upHandle(String data, String clientId) { + return null; + } + } + + @AgentDispatcher(msgId = MsgEnum.注册) + public class RegisterHandler implements AgentHandler { @Override public UpMsgResponse upHandle(String data, String clientId) { JSONObject jsonObject = new JSONObject(); jsonObject.put("clientId", clientId); - jsonObject.put("dataType", MsgEnum.注册应答.getValue()); + jsonObject.put("dataType", MsgEnum.注册.getValue()); jsonObject.put("data", data); MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(jsonObject.toString()).build(); agentProducer.asyncSend(mqMsg); - JSONObject json = new JSONObject(); - json.put("resCode",1); - return UpMsgResponse.builder().clientId(clientId).dataType(MsgEnum.注册应答.getValue()) - .content(json.toString()).build(); + return null; + } + } + + @AgentDispatcher(msgId = MsgEnum.注册应答) + public class RegisterRspHandler implements AgentHandler { + @Override + public String downHandle(String data, String clientId) { + return data; + } + } + + @AgentDispatcher(msgId = MsgEnum.获取最新策略) + public class GetPolicyHandler implements AgentHandler { + @Override + public UpMsgResponse upHandle(String data, String clientId) { + JSONObject jsonObject = new JSONObject(); + jsonObject.put("clientId", clientId); + jsonObject.put("dataType", MsgEnum.获取最新策略.getValue()); + jsonObject.put("data", data); + MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(jsonObject.toString()).build(); + agentProducer.asyncSend(mqMsg); + return null; + } + } + + @AgentDispatcher(msgId = MsgEnum.获取最新策略应答) + public class GetPolicyRspHandler implements AgentHandler { + @Override + public String downHandle(String data, String clientId) { + return data; } } @@ -64,9 +107,9 @@ public class AgentEndpoint { if(json.containsKey("resCode")){ resCode = json.getInteger("resCode"); } - if(resCode == 1){ - client.closeConnection(clientId); - } +// if(resCode == 1){ +// client.closeConnection(clientId); +// } return UpMsgResponse.builder().clientId(clientId).build(); } } @@ -350,9 +393,9 @@ public class AgentEndpoint { if(json.containsKey("resCode")){ resCode = json.getInteger("resCode"); } - if(resCode == 1){ - client.closeConnection(clientId); - } +// if(resCode == 1){ +// client.closeConnection(clientId); +// } return null; } } diff --git a/src/main/java/com/tongran/agent/server/netty/handler/TCPListenHandler.java b/src/main/java/com/tongran/agent/server/netty/handler/TCPListenHandler.java new file mode 100644 index 0000000..cd228e7 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/netty/handler/TCPListenHandler.java @@ -0,0 +1,66 @@ +package com.tongran.agent.server.netty.handler; + + +import com.tongran.agent.server.core.session.SessionManager; +import com.tongran.agent.server.service.AgentService; +import com.tongran.agent.server.utils.AssertLog; +import io.netty.channel.ChannelHandler; +import io.netty.channel.ChannelHandlerContext; +import io.netty.channel.ChannelInboundHandlerAdapter; +import io.netty.handler.timeout.IdleState; +import io.netty.handler.timeout.IdleStateEvent; +import org.springframework.stereotype.Component; + +import javax.annotation.Resource; +import java.io.IOException; + +/** + * @Description: TCP消息适配器 + */ +@Component +@ChannelHandler.Sharable +public class TCPListenHandler extends ChannelInboundHandlerAdapter { + + private final SessionManager sessionManager; + + @Resource + private AgentService agentService; + + public TCPListenHandler() { + this.sessionManager = SessionManager.getInstance(); + } + + @Override + public void channelActive(ChannelHandlerContext ctx) { + AssertLog.info("<<<<<<[终端连接]{}", ctx.channel().remoteAddress()); + } + + @Override + public void channelInactive(ChannelHandlerContext ctx) { + AssertLog.info(">>>>>>[断开连接]{}", sessionManager.client(ctx)); + sessionManager.remove(ctx.channel()); + } + + @Override + public void exceptionCaught(ChannelHandlerContext ctx, Throwable e) { + if (e instanceof IOException) { + AssertLog.info("<<<<<<[终端断开连接]{} {}", sessionManager.client(ctx), e.getMessage()); + } else { + AssertLog.info(">>>>>>[消息处理异常]" + sessionManager.client(ctx), e); + } + } + + @Override + public void userEventTriggered(ChannelHandlerContext ctx, Object evt) { + if (evt instanceof IdleStateEvent) { + IdleStateEvent event = (IdleStateEvent) evt; + IdleState state = event.state(); + if (state == IdleState.READER_IDLE || state == IdleState.WRITER_IDLE || state == IdleState.ALL_IDLE) { + AssertLog.warn(">>>>>>[终端心跳超时]{} {}", state, sessionManager.client(ctx)); + sessionManager.remove(ctx.channel()); + ctx.close(); + } + } + } + +} diff --git a/src/main/java/com/tongran/agent/server/netty/handler/UDPListenHandler.java b/src/main/java/com/tongran/agent/server/netty/handler/UDPListenHandler.java new file mode 100644 index 0000000..e39de19 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/netty/handler/UDPListenHandler.java @@ -0,0 +1,24 @@ +package com.tongran.agent.server.netty.handler; + +import io.netty.buffer.ByteBuf; +import io.netty.channel.ChannelHandler; +import io.netty.channel.ChannelHandlerContext; +import io.netty.channel.ChannelInboundHandlerAdapter; +import io.netty.channel.socket.DatagramPacket; +import org.springframework.stereotype.Component; + +/** + * @Description: UDP消息适配器 + */ +@Component +@ChannelHandler.Sharable +public class UDPListenHandler extends ChannelInboundHandlerAdapter { + + @Override + public void channelRead(ChannelHandlerContext ctx, Object msg) { + DatagramPacket packet = (DatagramPacket) msg; + ByteBuf buf = packet.content(); + buf.clear();// 暂未实现 + } + +} diff --git a/src/main/java/com/tongran/agent/server/rockermq/listener/AgentConsumerDownListener.java b/src/main/java/com/tongran/agent/server/rockermq/listener/AgentConsumerDownListener.java index 527b18d..e62b1ca 100644 --- a/src/main/java/com/tongran/agent/server/rockermq/listener/AgentConsumerDownListener.java +++ b/src/main/java/com/tongran/agent/server/rockermq/listener/AgentConsumerDownListener.java @@ -2,19 +2,12 @@ package com.tongran.agent.server.rockermq.listener; import cn.hutool.core.util.ObjectUtil; import com.alibaba.fastjson2.JSON; -import com.alibaba.fastjson2.JSONObject; -import com.tongran.agent.server.core.enums.MsgEnum; import com.tongran.agent.server.core.session.SessionManager; -import com.tongran.agent.server.netty.MultiTargetNettyClient; import com.tongran.agent.server.netty.annotation.AgentDispatcher; import com.tongran.agent.server.netty.basics.AgentDispatcherManager; import com.tongran.agent.server.netty.basics.AgentHandler; import com.tongran.agent.server.netty.model.Message; -import com.tongran.agent.server.rockermq.AgentProducer; -import com.tongran.agent.server.rockermq.config.AgentMqConfig; -import com.tongran.agent.server.rockermq.model.MqMsg; import com.tongran.agent.server.utils.AssertLog; -import org.apache.commons.lang3.StringUtils; import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; import org.apache.rocketmq.spring.core.RocketMQListener; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; @@ -33,15 +26,6 @@ public class AgentConsumerDownListener implements RocketMQListener { @Resource private AgentDispatcherManager agentDispatcherManager; - @Resource - private MultiTargetNettyClient client; - - @Resource - private AgentMqConfig agentMqConfig; - - @Resource - private AgentProducer agentProducer; - public AgentConsumerDownListener() { this.sessionManager = SessionManager.getInstance(); } @@ -51,112 +35,12 @@ public class AgentConsumerDownListener implements RocketMQListener { AssertLog.info("consumer==> received down message: {}", message); try { Message dto = JSON.parseObject(message, Message.class); - if(StringUtils.equals(dto.getDataType(), MsgEnum.注册.getValue())){ - JSONObject jsonObject = JSONObject.parseObject(dto.getData()); - String clientIp = ""; - int clientPort = 0; - String switchBoard = ""; - long timestamp = 0; - if(jsonObject.containsKey("clientIp")){ - clientIp = jsonObject.getString("clientIp"); - } - if(jsonObject.containsKey("clientPort")){ - clientPort = jsonObject.getInteger("clientPort"); - } - if(jsonObject.containsKey("switchBoard")){ - switchBoard = jsonObject.getString("switchBoard"); - } - if(jsonObject.containsKey("timestamp")){ - timestamp = jsonObject.getLongValue("timestamp"); - } - if(StringUtils.isBlank(clientIp) || clientPort == 0){ - //返回提醒ip 端口 空 - JSONObject json = new JSONObject(); - json.put("clientId",dto.getClientId()); - json.put("dataType",MsgEnum.注册应答.getValue()); - JSONObject dataJson = new JSONObject(); - dataJson.put("resCode",0); - dataJson.put("resMag","客户端IP、端口不能为空"); - dataJson.put("timestamp", System.currentTimeMillis()); - json.put("data",dataJson.toString()); - MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(json.toString()).build(); - agentProducer.asyncSend(mqMsg); - return; - } - boolean success = client.createConnection(dto.getClientId(), clientIp, clientPort, 5); - if(success){ - //连接成功 - // 连接建立后可以发送登录认证消息 - JSONObject object = new JSONObject(); - object.put("timestamp", timestamp); - if(StringUtils.isNotBlank(switchBoard)){ - String community = ""; - String switchIp = ""; - int switchPort = 0; - String netOID = ""; - String moduleOID = ""; - String mpuOID = ""; - String pwrOID = ""; - String fanOID = ""; - String otherOID = ""; - String filters = ""; - JSONObject json = JSONObject.parseObject(switchBoard); - if(json.containsKey("community")){ - community = json.getString("community"); - } - if(json.containsKey("ip")){ - switchIp = json.getString("ip"); - } - if(json.containsKey("port")){ - switchPort = json.getInteger("port"); - } - if(json.containsKey("netOID")){ - netOID = json.getString("netOID"); - } - if(json.containsKey("moduleOID")){ - moduleOID = json.getString("moduleOID"); - } - if(json.containsKey("mpuOID")){ - mpuOID = json.getString("mpuOID"); - } - if(json.containsKey("pwrOID")){ - pwrOID = json.getString("pwrOID"); - } - if(json.containsKey("fanOID")){ - fanOID = json.getString("fanOID"); - } - if(json.containsKey("otherOID")){ - otherOID = json.getString("otherOID"); - } - if(json.containsKey("filters")){ - filters = json.getString("filters"); - } - JSONObject switchBoardJson = new JSONObject(); - switchBoardJson.put("community",community); - switchBoardJson.put("ip",switchIp); - switchBoardJson.put("port",switchPort); - switchBoardJson.put("netOID",netOID); - switchBoardJson.put("moduleOID",moduleOID); - switchBoardJson.put("mpuOID",mpuOID); - switchBoardJson.put("pwrOID",pwrOID); - switchBoardJson.put("fanOID",fanOID); - switchBoardJson.put("otherOID",otherOID); - switchBoardJson.put("filters",filters); - object.put("switchBoard",switchBoardJson.toString()); - AssertLog.info("交换机配置信息={}",switchBoardJson.toString()); - } - Message msg = Message.builder().clientId(dto.getClientId()).dataType(MsgEnum.注册.getValue()).data(object.toString()).build(); - // 将对象转为 JSON 字符串 - client.sendMessages(dto.getClientId(), msg); - } - } else { - AgentHandler msgHandler = agentDispatcherManager.getHandler(dto.getDataType() + "&" + AgentDispatcher.VersionEnum.V1.value); - if (ObjectUtil.isNotEmpty(msgHandler)) { - String msg = msgHandler.downHandle(dto.getData(), dto.getClientId()); - Message chargeMessage = Message.builder().clientId(dto.getClientId()).dataType(dto.getDataType()).data(msg).build(); - if (Objects.nonNull(sessionManager.getSessionById(dto.getClientId()))) { - sessionManager.writeAndFlush(sessionManager.getSessionById(dto.getClientId()).getChannel(), chargeMessage); - } + AgentHandler msgHandler = agentDispatcherManager.getHandler(dto.getDataType() + "&" + AgentDispatcher.VersionEnum.V1.value); + if (ObjectUtil.isNotEmpty(msgHandler)) { + String msg = msgHandler.downHandle(dto.getData(), dto.getClientId()); + Message chargeMessage = Message.builder().clientId(dto.getClientId()).dataType(dto.getDataType()).data(msg).build(); + if (Objects.nonNull(sessionManager.getSessionById(dto.getClientId()))) { + sessionManager.writeAndFlush(sessionManager.getSessionById(dto.getClientId()).getChannel(), chargeMessage); } } } catch (Exception e) { diff --git a/src/main/java/com/tongran/agent/server/service/impl/AgentServiceImpl.java b/src/main/java/com/tongran/agent/server/service/impl/AgentServiceImpl.java index 787ffb5..fd289fc 100644 --- a/src/main/java/com/tongran/agent/server/service/impl/AgentServiceImpl.java +++ b/src/main/java/com/tongran/agent/server/service/impl/AgentServiceImpl.java @@ -3,37 +3,35 @@ package com.tongran.agent.server.service.impl; import com.alibaba.fastjson2.JSON; import com.alibaba.fastjson2.JSONObject; import com.tongran.agent.server.core.vo.RegisterVO; -import com.tongran.agent.server.netty.MultiTargetNettyClient; import com.tongran.agent.server.netty.model.Message; import com.tongran.agent.server.service.AgentService; import org.springframework.stereotype.Service; -import javax.annotation.Resource; - @Service public class AgentServiceImpl implements AgentService { - @Resource - private MultiTargetNettyClient client; +// @Resource +// private MultiTargetNettyClient client; @Override public boolean createConn(String clientId, String host, int port, int timeout) { // 动态创建不同IP和端口的连接 - boolean success = client.createConnection(clientId, host, port, timeout); - if(success){ - // 连接建立后可以发送建立消息 - long timestamp = System.currentTimeMillis(); - JSONObject object = new JSONObject(); - object.put("clientId",clientId); - object.put("timestamp",timestamp); - Message message = Message.builder().clientId(clientId).dataType("CREATE").data(object.toString()).build(); - // 将对象转为 JSON 字符串 - String json = JSON.toJSONString(message); - System.out.println("连接建立成功,发送建立消息="+json); - client.sendMessages(clientId, message); - } - return success; +// boolean success = client.createConnection(clientId, host, port, timeout); +// if(success){ +// // 连接建立后可以发送建立消息 +// long timestamp = System.currentTimeMillis(); +// JSONObject object = new JSONObject(); +// object.put("clientId",clientId); +// object.put("timestamp",timestamp); +// Message message = Message.builder().clientId(clientId).dataType("CREATE").data(object.toString()).build(); +// // 将对象转为 JSON 字符串 +// String json = JSON.toJSONString(message); +// System.out.println("连接建立成功,发送建立消息="+json); +// client.sendMessages(clientId, message); +// } +// return success; + return false; } @Override @@ -52,6 +50,6 @@ public class AgentServiceImpl implements AgentService { // 将对象转为 JSON 字符串 String json = JSON.toJSONString(message); System.out.println("发送注册消息="+json); - client.sendMessages(registerVO.getClientId(), message); +// client.sendMessages(registerVO.getClientId(), message); } } diff --git a/src/main/java/com/tongran/agent/server/utils/ClientExample.java b/src/main/java/com/tongran/agent/server/utils/ClientExample.java index ef75204..84eab55 100644 --- a/src/main/java/com/tongran/agent/server/utils/ClientExample.java +++ b/src/main/java/com/tongran/agent/server/utils/ClientExample.java @@ -40,7 +40,7 @@ public class ClientExample { } catch (InterruptedException e) { Thread.currentThread().interrupt(); } finally { - manager.shutdown(); +// manager.shutdown(); } } } \ No newline at end of file diff --git a/src/main/java/com/tongran/agent/server/utils/EnhancedConnectionManager.java b/src/main/java/com/tongran/agent/server/utils/EnhancedConnectionManager.java index 223b4e4..d7060d8 100644 --- a/src/main/java/com/tongran/agent/server/utils/EnhancedConnectionManager.java +++ b/src/main/java/com/tongran/agent/server/utils/EnhancedConnectionManager.java @@ -1,6 +1,5 @@ package com.tongran.agent.server.utils; -import com.tongran.agent.server.netty.MultiTargetNettyClient; import com.tongran.agent.server.netty.config.ConnectionConfig; import java.util.HashMap; @@ -12,11 +11,11 @@ import java.util.concurrent.ConcurrentHashMap; * 增强版连接管理器 */ public class EnhancedConnectionManager { - private final MultiTargetNettyClient client; +// private final MultiTargetNettyClient client; private final Map connectionConfigs; public EnhancedConnectionManager() { - this.client = new MultiTargetNettyClient(); +// this.client = new MultiTargetNettyClient(); this.connectionConfigs = new ConcurrentHashMap<>(); } @@ -26,18 +25,18 @@ public class EnhancedConnectionManager { public Map createConnections(List configs) { Map results = new HashMap<>(); - configs.forEach(config -> { - boolean success = client.createConnection( - config.getConnectionKey(), - config.getHost(), - config.getPort(), - config.getTimeout() - ); - if (success) { - connectionConfigs.put(config.getConnectionKey(), config); - } - results.put(config.getConnectionKey(), success); - }); +// configs.forEach(config -> { +// boolean success = client.createConnection( +// config.getConnectionKey(), +// config.getHost(), +// config.getPort(), +// config.getTimeout() +// ); +// if (success) { +// connectionConfigs.put(config.getConnectionKey(), config); +// } +// results.put(config.getConnectionKey(), success); +// }); return results; } @@ -47,21 +46,21 @@ public class EnhancedConnectionManager { */ public void routeMessageByType(String messageType, String message) { // 这里可以根据业务逻辑决定发送到哪个连接 - switch (messageType) { - case "TYPE_A": - client.sendMessage("server1", "[TYPE_A] " + message); - break; - case "TYPE_B": - client.sendMessage("server2", "[TYPE_B] " + message); - break; - case "TYPE_C": - client.sendMessage("server3", "[TYPE_C] " + message); - break; - default: - // 广播到所有连接 - connectionConfigs.keySet().forEach(key -> - client.sendMessage(key, "[BROADCAST] " + message)); - } +// switch (messageType) { +// case "TYPE_A": +// client.sendMessage("server1", "[TYPE_A] " + message); +// break; +// case "TYPE_B": +// client.sendMessage("server2", "[TYPE_B] " + message); +// break; +// case "TYPE_C": +// client.sendMessage("server3", "[TYPE_C] " + message); +// break; +// default: +// // 广播到所有连接 +// connectionConfigs.keySet().forEach(key -> +// client.sendMessage(key, "[BROADCAST] " + message)); +// } } /** @@ -69,14 +68,14 @@ public class EnhancedConnectionManager { */ public Map getConnectionStatus() { Map status = new HashMap<>(); - connectionConfigs.forEach((key, config) -> { - // 这里可以添加更详细的状态检查 - status.put(key, client.sendMessage(key, "PING")); - }); +// connectionConfigs.forEach((key, config) -> { +// // 这里可以添加更详细的状态检查 +// status.put(key, client.sendMessage(key, "PING")); +// }); return status; } - public void shutdown() { - client.shutdown(); - } +// public void shutdown() { +// client.shutdown(); +// } } \ No newline at end of file diff --git a/src/main/resources/application-dev.yml b/src/main/resources/application-dev.yml index 3da7cc2..7d213ba 100644 --- a/src/main/resources/application-dev.yml +++ b/src/main/resources/application-dev.yml @@ -12,21 +12,18 @@ knife4j: logging: file: # path: /usr/local/tongran/logs - path: /data/tr-agent-server/logs + path: /usr/local/tongran_server/logs # path: D:/job/agent-logs/logs rocketmq: name-server: 172.16.15.103:9876 enabled: true -netty: - server: -# host: 172.16.15.103 - host: 127.0.0.1 - port: 6610 - client: - client-id: client-001 - reconnect-interval: 5 - maxReconnectAttempts: 10 # 最大重连次数 - initialReconnectDelay: 1000 # 初始重连延迟(毫秒) - maxReconnectDelay: 30000 # 最大重连延迟(毫秒) \ No newline at end of file +tcp: + netty: + charge: + enable: true + name: AGENT-SERVER-服务 + port: 6620 + readerIdleTime: 300 + diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml index 3e7f3cb..b504121 100644 --- a/src/main/resources/application.yml +++ b/src/main/resources/application.yml @@ -34,10 +34,10 @@ rocketmq: producer: group: tongran-group topic: default_topic - agent-group: tongran_agent_up_group - agent-topic: tongran_agent_up + agent-group: tr_agent_up_group + agent-topic: tr_agent_up consumer: group: tongran-group topic: default_topic - agent-group: tongran_agent_down_group - agent-topic: tongran_agent_down \ No newline at end of file + agent-group: tr_agent_down_group + agent-topic: tr_agent_down \ No newline at end of file