From 22acd693dc26fc195498b40009a455ab4685d3ec Mon Sep 17 00:00:00 2001 From: gaoyutao Date: Mon, 17 Nov 2025 11:36:46 +0800 Subject: [PATCH] =?UTF-8?q?mtr-agent-server=20v1.0.0=E5=88=9D=E5=A7=8B?= =?UTF-8?q?=E5=8C=96?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pom.xml | 107 +++++ .../server/TrAgentServerApplication.java | 13 + .../server/controller/TcpController.java | 29 ++ .../agent/server/controller/TcpService.java | 118 +++++ .../server/core/config/GlobalConfig.java | 17 + .../agent/server/core/enums/MsgEnum.java | 86 ++++ .../agent/server/core/session/Session.java | 85 ++++ .../server/core/session/SessionManager.java | 137 ++++++ .../agent/server/core/vo/RegisterVO.java | 63 +++ .../server/exception/ServerException.java | 27 ++ .../server/exception/base/BaseException.java | 42 ++ .../server/exception/code/ErrorCode.java | 16 + .../exception/code/GlobalErrorCode.java | 20 + .../agent/server/netty/AgentNettyConfig.java | 22 + .../agent/server/netty/AgentNettyServer.java | 56 +++ .../agent/server/netty/BaseNettyServer.java | 117 +++++ .../netty/annotation/AgentDispatcher.java | 47 ++ .../netty/basics/AgentDispatcherManager.java | 48 ++ .../server/netty/basics/AgentHandler.java | 21 + .../server/netty/config/BaseNettyConfig.java | 58 +++ .../server/netty/config/ConnectionConfig.java | 39 ++ .../server/netty/enpoint/AgentEndpoint.java | 420 ++++++++++++++++++ .../netty/handler/AgentDecoderHandler.java | 174 ++++++++ .../netty/handler/AgentDispatcherHandler.java | 54 +++ .../netty/handler/AgentEncoderHandler.java | 49 ++ .../netty/handler/TCPListenHandler.java | 66 +++ .../netty/handler/UDPListenHandler.java | 24 + .../agent/server/netty/model/Message.java | 37 ++ .../server/netty/model/UpMsgResponse.java | 22 + .../agent/server/rockermq/AgentProducer.java | 54 +++ .../server/rockermq/RocketMqService.java | 45 ++ .../server/rockermq/config/AgentMqConfig.java | 14 + .../listener/AgentConsumerDownListener.java | 51 +++ .../agent/server/rockermq/model/MqMsg.java | 38 ++ .../agent/server/service/AgentService.java | 10 + .../server/service/impl/AgentServiceImpl.java | 55 +++ .../tongran/agent/server/utils/AssertLog.java | 73 +++ .../agent/server/utils/ClientExample.java | 46 ++ .../utils/EnhancedConnectionManager.java | 81 ++++ .../com/tongran/agent/server/utils/R.java | 76 ++++ src/main/resources/application-dev.yml | 29 ++ src/main/resources/application.yml | 43 ++ src/main/resources/logback/logback-dev.xml | 116 +++++ src/main/resources/logback/logback-prod.xml | 115 +++++ 44 files changed, 2860 insertions(+) create mode 100644 pom.xml create mode 100644 src/main/java/com/tongran/agent/server/TrAgentServerApplication.java create mode 100644 src/main/java/com/tongran/agent/server/controller/TcpController.java create mode 100644 src/main/java/com/tongran/agent/server/controller/TcpService.java create mode 100644 src/main/java/com/tongran/agent/server/core/config/GlobalConfig.java create mode 100644 src/main/java/com/tongran/agent/server/core/enums/MsgEnum.java create mode 100644 src/main/java/com/tongran/agent/server/core/session/Session.java create mode 100644 src/main/java/com/tongran/agent/server/core/session/SessionManager.java create mode 100644 src/main/java/com/tongran/agent/server/core/vo/RegisterVO.java create mode 100644 src/main/java/com/tongran/agent/server/exception/ServerException.java create mode 100644 src/main/java/com/tongran/agent/server/exception/base/BaseException.java create mode 100644 src/main/java/com/tongran/agent/server/exception/code/ErrorCode.java create mode 100644 src/main/java/com/tongran/agent/server/exception/code/GlobalErrorCode.java 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 create mode 100644 src/main/java/com/tongran/agent/server/netty/BaseNettyServer.java create mode 100644 src/main/java/com/tongran/agent/server/netty/annotation/AgentDispatcher.java create mode 100644 src/main/java/com/tongran/agent/server/netty/basics/AgentDispatcherManager.java create mode 100644 src/main/java/com/tongran/agent/server/netty/basics/AgentHandler.java create mode 100644 src/main/java/com/tongran/agent/server/netty/config/BaseNettyConfig.java create mode 100644 src/main/java/com/tongran/agent/server/netty/config/ConnectionConfig.java create mode 100644 src/main/java/com/tongran/agent/server/netty/enpoint/AgentEndpoint.java create mode 100644 src/main/java/com/tongran/agent/server/netty/handler/AgentDecoderHandler.java create mode 100644 src/main/java/com/tongran/agent/server/netty/handler/AgentDispatcherHandler.java create mode 100644 src/main/java/com/tongran/agent/server/netty/handler/AgentEncoderHandler.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 create mode 100644 src/main/java/com/tongran/agent/server/netty/model/Message.java create mode 100644 src/main/java/com/tongran/agent/server/netty/model/UpMsgResponse.java create mode 100644 src/main/java/com/tongran/agent/server/rockermq/AgentProducer.java create mode 100644 src/main/java/com/tongran/agent/server/rockermq/RocketMqService.java create mode 100644 src/main/java/com/tongran/agent/server/rockermq/config/AgentMqConfig.java create mode 100644 src/main/java/com/tongran/agent/server/rockermq/listener/AgentConsumerDownListener.java create mode 100644 src/main/java/com/tongran/agent/server/rockermq/model/MqMsg.java create mode 100644 src/main/java/com/tongran/agent/server/service/AgentService.java create mode 100644 src/main/java/com/tongran/agent/server/service/impl/AgentServiceImpl.java create mode 100644 src/main/java/com/tongran/agent/server/utils/AssertLog.java create mode 100644 src/main/java/com/tongran/agent/server/utils/ClientExample.java create mode 100644 src/main/java/com/tongran/agent/server/utils/EnhancedConnectionManager.java create mode 100644 src/main/java/com/tongran/agent/server/utils/R.java create mode 100644 src/main/resources/application-dev.yml create mode 100644 src/main/resources/application.yml create mode 100644 src/main/resources/logback/logback-dev.xml create mode 100644 src/main/resources/logback/logback-prod.xml diff --git a/pom.xml b/pom.xml new file mode 100644 index 0000000..701f994 --- /dev/null +++ b/pom.xml @@ -0,0 +1,107 @@ + + + 4.0.0 + + org.springframework.boot + spring-boot-starter-parent + 2.5.6 + + + com.tongran.agent + tr-agent-server + 0.0.1-SNAPSHOT + tr-agent-server + tr-agent-server + + 1.8 + + + + org.springframework.boot + spring-boot-starter-web + + + + org.springframework.boot + spring-boot-starter-test + test + + + + org.springframework.boot + spring-boot-configuration-processor + + + + org.springframework.boot + spring-boot-starter-aop + + + + org.springframework.boot + spring-boot-starter-validation + + + + org.projectlombok + lombok + true + + + + io.swagger.core.v3 + swagger-annotations + 2.2.19 + + + + cn.hutool + hutool-all + 5.7.11 + + + + org.apache.commons + commons-lang3 + + + + io.netty + netty-all + 4.1.86.Final + + + + com.alibaba.fastjson2 + fastjson2 + 2.0.31 + + + + org.apache.rocketmq + rocketmq-spring-boot-starter + 2.2.3 + + + + + + + + org.springframework.boot + spring-boot-maven-plugin + + + + org.apache.maven.plugins + maven-surefire-plugin + ${maven-surefire-plugin.version} + + true + + + + + + diff --git a/src/main/java/com/tongran/agent/server/TrAgentServerApplication.java b/src/main/java/com/tongran/agent/server/TrAgentServerApplication.java new file mode 100644 index 0000000..90e6784 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/TrAgentServerApplication.java @@ -0,0 +1,13 @@ +package com.tongran.agent.server; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; + +@SpringBootApplication +public class TrAgentServerApplication { + + public static void main(String[] args) { + SpringApplication.run(TrAgentServerApplication.class, args); + } + +} diff --git a/src/main/java/com/tongran/agent/server/controller/TcpController.java b/src/main/java/com/tongran/agent/server/controller/TcpController.java new file mode 100644 index 0000000..04233ec --- /dev/null +++ b/src/main/java/com/tongran/agent/server/controller/TcpController.java @@ -0,0 +1,29 @@ +package com.tongran.agent.server.controller; + +import com.alibaba.fastjson2.JSONObject; +import com.tongran.agent.server.netty.model.Message; +import org.springframework.validation.annotation.Validated; +import org.springframework.web.bind.annotation.PostMapping; +import org.springframework.web.bind.annotation.RequestBody; +import org.springframework.web.bind.annotation.RestController; + +import javax.annotation.Resource; + +@RestController +public class TcpController { + + @Resource + private TcpService tcpService; + + @PostMapping("/send") + public void sendMessage(@Validated @RequestBody Message dto) { + JSONObject jsonObject = new JSONObject(); + jsonObject.put("clientId", dto.getClientId()); + jsonObject.put("dataType", dto.getDataType()); + jsonObject.put("data", dto.getData()); + tcpService.sendMessage(jsonObject.toString()); + } + + + +} \ No newline at end of file diff --git a/src/main/java/com/tongran/agent/server/controller/TcpService.java b/src/main/java/com/tongran/agent/server/controller/TcpService.java new file mode 100644 index 0000000..ccc8642 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/controller/TcpService.java @@ -0,0 +1,118 @@ +package com.tongran.agent.server.controller; + +import com.tongran.agent.server.core.session.SessionManager; +import com.tongran.agent.server.netty.basics.AgentDispatcherManager; +import com.tongran.agent.server.rockermq.AgentProducer; +import com.tongran.agent.server.utils.AssertLog; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; + +import javax.annotation.Resource; + +@Slf4j +@Service +public class TcpService { + protected final SessionManager sessionManager; + + @Resource + private AgentDispatcherManager agentDispatcherManager; + +// @Resource +// private MultiTargetNettyClient client; + + @Resource + private AgentProducer agentProducer; + + public TcpService() { + this.sessionManager = SessionManager.getInstance(); + } + + 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); +// } + } + + 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()); +// } + } + + +} diff --git a/src/main/java/com/tongran/agent/server/core/config/GlobalConfig.java b/src/main/java/com/tongran/agent/server/core/config/GlobalConfig.java new file mode 100644 index 0000000..af17931 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/core/config/GlobalConfig.java @@ -0,0 +1,17 @@ +package com.tongran.agent.server.core.config; + +/** + * 全局配置类 + */ +public class GlobalConfig { + + /** + * 交换机 + */ + public static String SWITCH_COMMUNITY; //团名 + public static String SWITCH_IP; //交换机IP + public static int SWITCH_PORT; //交换机端口 + public static String[] SWITCH_OID; //交换机OID + + +} 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 new file mode 100644 index 0000000..1bbb9d7 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/core/enums/MsgEnum.java @@ -0,0 +1,86 @@ +package com.tongran.agent.server.core.enums; + +import lombok.AllArgsConstructor; +import lombok.Getter; +import lombok.NoArgsConstructor; + +/** + * 控制码 + */ +@Getter +@AllArgsConstructor +@NoArgsConstructor +public enum MsgEnum { + 建立连接("CONNECT"), + + 建立连接应答("CONNECT_RSP"), + + 注册("REGISTER"), + + 注册应答("REGISTER_RSP"), + + 获取最新策略("GET_POLICY"), + + 获取最新策略应答("GET_POLICY_RSP"), + + 断开("DISCONNECT"), + + 断开应答("DISCONNECT_RSP"), + + 心跳上报("HEARTBEAT"), + + 心跳上报应答("HEARTBEAT_RSP"), + + CPU上报("CPU"), + + 磁盘上报("DISK"), + + 容器上报("DOCKER"), + + 内存上报("MEMORY"), + + 网络上报("NET"), + + 挂载上报("POINT"), + + 系统其他上报("OTHER_SYSTEM"), + + 交换机上报("SWITCHBOARD"), + + 采集上报应答("COLLECT_RSP"), + + 开启或更新系统采集("SYSTEM_COLLECT_START"), + + 开启或更新系统采集应答("SYSTEM_COLLECT_START_RSP"), + + 关闭所有系统采集("SYSTEM_COLLECT_STOP"), + + 关闭所有系统采集应答("SYSTEM_COLLECT_STOP_RSP"), + + 开启或更新交换机采集("SWITCH_COLLECT_START"), + + 开启或更新交换机采集应答("SWITCH_COLLECT_START_RSP"), + + 关闭所有交换机采集("SWITCH_COLLECT_STOP"), + + 关闭所有交换机采集应答("SWITCH_COLLECT_STOP_RSP"), + + 告警设置("ALARM_SET"), + + 告警设置应答("ALARM_SET_RSP"), + + 执行脚本策略("SCRIPT_POLICY"), + + 执行脚本策略应答("SCRIPT_POLICY_RSP"), + + Agent版本更新("AGENT_VERSION_UPDATE"), + + Agent版本更新应答("AGENT_VERSION_UPDATE_RSP"), + + 多网IP探测上报("NETWORK_DETECT"); + + private String value; + + + +} diff --git a/src/main/java/com/tongran/agent/server/core/session/Session.java b/src/main/java/com/tongran/agent/server/core/session/Session.java new file mode 100644 index 0000000..04f3cbf --- /dev/null +++ b/src/main/java/com/tongran/agent/server/core/session/Session.java @@ -0,0 +1,85 @@ +package com.tongran.agent.server.core.session; + +import io.netty.channel.ChannelHandlerContext; +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; +import lombok.experimental.Accessors; + +import java.io.Serializable; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.atomic.AtomicInteger; + +/** + * SESSION + * + */ +@Data +@Builder +@Accessors(chain = true) +@AllArgsConstructor +@NoArgsConstructor +public class Session implements Serializable { + + private static final long serialVersionUID = 1L; + + /** + * sessionId + */ + private String sessionId; + + /** + * 客户端唯一标识 + */ + private String clientId; + + /** + * Seesion额外参数 + */ + private ConcurrentHashMap extraMap; + + /** + * 是否鉴权 + */ + @Builder.Default + private boolean isAuthenticated = false; + + @Builder.Default + private long createTime = System.currentTimeMillis(); + + @Builder.Default + private long updateTime = System.currentTimeMillis(); + + /** + * 消息渠道 + */ + private ChannelHandlerContext channel; + + @Builder.Default + private AtomicInteger serialNo = new AtomicInteger(0); + + public static String buildSessionId(ChannelHandlerContext channel) { + return channel.channel().id().asLongText(); + } + + public static Session buildSession(ChannelHandlerContext channel, String clientId) { + return Session.builder().channel(channel).sessionId(buildSessionId(channel)).clientId(clientId).build(); + } + + /** + * 自生成流水号 + * + * @return + */ + public int nextSerialNo() { + int current; + int next; + do { + current = serialNo.get(); + next = current > 0xffff ? 0 : current; + } while (!serialNo.compareAndSet(current, next + 1)); + return next; + } + +} \ No newline at end of file diff --git a/src/main/java/com/tongran/agent/server/core/session/SessionManager.java b/src/main/java/com/tongran/agent/server/core/session/SessionManager.java new file mode 100644 index 0000000..41dcc58 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/core/session/SessionManager.java @@ -0,0 +1,137 @@ +package com.tongran.agent.server.core.session; + + +import cn.hutool.cache.Cache; +import cn.hutool.cache.CacheUtil; +import cn.hutool.core.bean.BeanUtil; +import cn.hutool.core.util.ObjectUtil; +import com.tongran.agent.server.exception.ServerException; +import com.tongran.agent.server.utils.AssertLog; +import io.netty.channel.Channel; +import io.netty.channel.ChannelHandlerContext; +import io.netty.util.Attribute; +import io.netty.util.AttributeKey; +import lombok.Data; +import org.apache.commons.lang3.StringUtils; + +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; + +@Data +public class SessionManager { + + private static volatile SessionManager instance = null; + + // clientId-Session + private final Map sessionMap = new ConcurrentHashMap<>(); + + private final Cache sessionCache = CacheUtil.newTimedCache(10 * 60 * 1000); + + public static SessionManager getInstance() { + if (instance == null) { + synchronized (SessionManager.class) { + if (instance == null) { + instance = new SessionManager(); + } + } + } + return instance; + } + + public synchronized void put(String clientId, Session session) { + if (StringUtils.isNotBlank(session.getClientId())) { + sessionMap.put(session.getClientId(), session); + } + } + + public synchronized void remove(String clientId) { + if (StringUtils.isNotBlank(clientId)) { + sessionMap.remove(clientId); + } + } + + public synchronized void remove(Session session) { + if (session != null && StringUtils.isNotBlank(session.getClientId())) { + sessionMap.remove(session.getClientId(), session); + } + } + + public synchronized void remove(Channel channel) { + remove(getSessionByChannel(channel)); + } + + public Session getSessionById(String clientId) { + return sessionMap.get(clientId); + } + + public Session getSessionByChannel(Channel channel) { + Session session = new Session(); + sessionMap.values().forEach(s -> { + if (s.getChannel().channel() == channel) { + BeanUtil.copyProperties(s, session, true); + } + }); + return session; + } + + public boolean containsSession(String clientId) { + return sessionMap.containsKey(clientId); + } + + public boolean containsSession(Session session) { + return sessionMap.containsValue(session); + } + + public void setSessionCache(String clientId, Object value) { + sessionCache.put(clientId, value); + } + + public Object getSessionCache(String clientId) { + return sessionCache.get(clientId); + } + + public String client(ChannelHandlerContext ctx) { + Channel channel = ctx.channel(); + Session session = this.getSessionByChannel(channel); + if (ObjectUtil.isNotNull(session) && StringUtils.isNotBlank(session.getClientId())) { + return channel.remoteAddress().toString() + "/" + session.getClientId(); + } + return channel.remoteAddress().toString(); + } + + /** + * 根据channel生成流水号 + * + * @param channel + * @return + */ + public short getSerialNumber(Channel channel, AttributeKey serialNumber) { + Attribute flowIdAttr = channel.attr(serialNumber); + Short flowId = flowIdAttr.get(); + if (flowId == null) { + flowId = 0; + } else { + flowId++; + } + flowIdAttr.set(flowId); + return flowId; + } + + public void writeAndFlush(String clientId, Object msg) { + Session session = this.getSessionById(clientId); + if (ObjectUtil.isNotNull(session)) { + this.writeAndFlush(session.getChannel(), msg); + } else { + throw new ServerException(500, "终端:" + clientId + "离线或者不存在"); + } + } + + public void writeAndFlush(ChannelHandlerContext ctx, Object msg) { + ctx.writeAndFlush(msg).addListener(future -> { + if (!future.isSuccess()) { + AssertLog.error("消息发送失败:{}", future.cause()); + } + }); + } + +} \ No newline at end of file diff --git a/src/main/java/com/tongran/agent/server/core/vo/RegisterVO.java b/src/main/java/com/tongran/agent/server/core/vo/RegisterVO.java new file mode 100644 index 0000000..9c4cd35 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/core/vo/RegisterVO.java @@ -0,0 +1,63 @@ +package com.tongran.agent.server.core.vo; + +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; + +@Data +@Builder +@AllArgsConstructor +@NoArgsConstructor +public class RegisterVO { + + //客户端ID(对应平台唯一SN) + private String clientId; + + //客户端IP + private String clientIp; + + //客户端端口 + private int clientPort; + + //交换机配置 + private SwitchBoard switchBoard; + + //时间戳 + private long timestamp; + + @Data + @Builder + @AllArgsConstructor + @NoArgsConstructor + public static class SwitchBoard{ + //交换机团名 + private String community; + + //交换机IP + private String ip; + + //交换机端口 + private int port; + + //交换机网络端口发现OID + private String netOID; + + //交换机光模块发现OID + private String moduleOID; + + //交换机MPU发现OID + private String mpuOID; + + //交换机电源发现OID + private String pwrOID; + + //交换机风扇发现OID + private String fanOID; + + //交换机系统其他OID + private String otherOID; + + } + +} diff --git a/src/main/java/com/tongran/agent/server/exception/ServerException.java b/src/main/java/com/tongran/agent/server/exception/ServerException.java new file mode 100644 index 0000000..ae1ac62 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/exception/ServerException.java @@ -0,0 +1,27 @@ +package com.tongran.agent.server.exception; + +import cn.hutool.core.util.StrUtil; +import com.tongran.agent.server.exception.base.BaseException; +import com.tongran.agent.server.exception.code.ErrorCode; + +public class ServerException extends BaseException { + + private static final long serialVersionUID = -519977195881081739L; + + public ServerException(ErrorCode code) { + super(code.getCode(), code.getMsg()); + } + + public ServerException(Integer code, String msg) { + super(code, msg); + } + + public ServerException(String msg) { + super(500, msg); + } + + public ServerException(String msg, Object... arguments) { + super(500, StrUtil.format(msg, arguments)); + } + +} diff --git a/src/main/java/com/tongran/agent/server/exception/base/BaseException.java b/src/main/java/com/tongran/agent/server/exception/base/BaseException.java new file mode 100644 index 0000000..c27f4a2 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/exception/base/BaseException.java @@ -0,0 +1,42 @@ +package com.tongran.agent.server.exception.base; + + +import com.tongran.agent.server.exception.code.ErrorCode; + +public class BaseException extends RuntimeException implements ErrorCode { + + private static final long serialVersionUID = 1966249840643379123L; + + private Integer code; + + private String msg; + + public BaseException() { + } + + public BaseException(String msg, Throwable cause) { + super(msg, cause); + } + + public BaseException(Integer code, String msg) { + super(msg); + this.code = code; + this.msg = msg; + } + + public BaseException(Integer code, String msg, Throwable cause) { + super(msg, cause); + this.code = code; + this.msg = msg; + } + + @Override + public Integer getCode() { + return code; + } + + @Override + public String getMsg() { + return msg; + } +} \ No newline at end of file diff --git a/src/main/java/com/tongran/agent/server/exception/code/ErrorCode.java b/src/main/java/com/tongran/agent/server/exception/code/ErrorCode.java new file mode 100644 index 0000000..f59044e --- /dev/null +++ b/src/main/java/com/tongran/agent/server/exception/code/ErrorCode.java @@ -0,0 +1,16 @@ +package com.tongran.agent.server.exception.code; + + +import com.tongran.agent.server.utils.R; + +public interface ErrorCode { + + Integer getCode(); + + String getMsg(); + + default R toResult() { + return R.error(getMsg(), getCode()); + } + +} diff --git a/src/main/java/com/tongran/agent/server/exception/code/GlobalErrorCode.java b/src/main/java/com/tongran/agent/server/exception/code/GlobalErrorCode.java new file mode 100644 index 0000000..ba0498b --- /dev/null +++ b/src/main/java/com/tongran/agent/server/exception/code/GlobalErrorCode.java @@ -0,0 +1,20 @@ +package com.tongran.agent.server.exception.code; + +import lombok.AllArgsConstructor; +import lombok.Getter; + +@Getter +@AllArgsConstructor +public enum GlobalErrorCode implements ErrorCode { + + SUCCESS(200, "成功"), + + NOT_FOUND(404, "未找到相关资源"), + + ERROR(500, "系统错误"); + + private final Integer code; + + private final String msg; + +} 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/BaseNettyServer.java b/src/main/java/com/tongran/agent/server/netty/BaseNettyServer.java new file mode 100644 index 0000000..50e5d72 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/netty/BaseNettyServer.java @@ -0,0 +1,117 @@ +package com.tongran.agent.server.netty; + +import com.tongran.agent.server.netty.config.BaseNettyConfig; +import com.tongran.agent.server.utils.AssertLog; +import io.netty.bootstrap.AbstractBootstrap; +import io.netty.bootstrap.Bootstrap; +import io.netty.bootstrap.ServerBootstrap; +import io.netty.channel.ChannelFuture; +import io.netty.channel.ChannelOption; +import io.netty.channel.EventLoopGroup; +import io.netty.channel.nio.NioEventLoopGroup; +import io.netty.channel.socket.nio.NioChannelOption; +import io.netty.channel.socket.nio.NioDatagramChannel; +import io.netty.channel.socket.nio.NioServerSocketChannel; +import io.netty.util.ResourceLeakDetector; +import io.netty.util.concurrent.DefaultEventExecutorGroup; +import io.netty.util.concurrent.DefaultThreadFactory; +import io.netty.util.concurrent.EventExecutorGroup; +import io.netty.util.concurrent.Future; + +/** + * 基础Netty服务 + */ +public abstract class BaseNettyServer { + + protected boolean isRunning; + + protected BaseNettyConfig config; + + protected EventLoopGroup bossGroup; + + protected EventLoopGroup workerGroup; + + protected EventExecutorGroup businessGroup; + + protected BaseNettyServer(BaseNettyConfig config) { + this.config = config; + } + + protected AbstractBootstrap initializeTcp() { + bossGroup = new NioEventLoopGroup(1, new DefaultThreadFactory(config.name, Thread.MAX_PRIORITY)); + workerGroup = new NioEventLoopGroup(config.workerCore, new DefaultThreadFactory(config.name, Thread.MAX_PRIORITY)); + if (config.businessCore > 0) { + businessGroup = new DefaultEventExecutorGroup(config.businessCore); + } + ServerBootstrap serverBootstrap = new ServerBootstrap(); + serverBootstrap.group(bossGroup, workerGroup) + .channel(NioServerSocketChannel.class) + .childOption(ChannelOption.SO_REUSEADDR, true) + .option(ChannelOption.SO_BACKLOG, 1024) + .childOption(NioChannelOption.TCP_NODELAY, true) + .childHandler(config.hander); + //内存泄漏检测 开发推荐PARANOID 线上SIMPLE + ResourceLeakDetector.setLevel(ResourceLeakDetector.Level.SIMPLE); + return serverBootstrap; + } + + protected AbstractBootstrap initializeUdp() { + bossGroup = new NioEventLoopGroup(1, new DefaultThreadFactory(config.name, Thread.MAX_PRIORITY)); + if (config.businessCore > 0) { + businessGroup = new DefaultEventExecutorGroup(config.businessCore); + } + return new Bootstrap() + .group(bossGroup).channel(NioDatagramChannel.class) + .option(NioChannelOption.SO_REUSEADDR, true) + .option(NioChannelOption.SO_RCVBUF, 1024 * 1024 * 50) + .handler(config.hander); + } + + public synchronized boolean start() { + if (!config.enable) { + return false; + } + if (isRunning) { + AssertLog.info("======{}已经启动,port:{}======", config.name, config.port); + return isRunning; + } + AbstractBootstrap bootstrap = config.isTcp ? initializeTcp() : initializeUdp(); + ChannelFuture future = bootstrap.bind(config.port).awaitUninterruptibly(); + future.channel().closeFuture().addListener(f -> { + if (isRunning) { + stop(); + } + }); + if (future.cause() != null) { + AssertLog.error("===启动失败===", future.cause()); + } + if (isRunning = future.isSuccess()) { + AssertLog.info("\n\n\t\t\t\t\t\t\t\t======{}启动成功,port:{}======\n", config.name, config.port); + } + return isRunning; + } + + public synchronized void stop() { + if (!config.enable) { + return; + } + isRunning = false; + try { + Future future = this.workerGroup.shutdownGracefully().await(); + if (!future.isSuccess()) { + AssertLog.error("workerGroup 无法正常停止:{}", future.cause()); + } + future = this.bossGroup.shutdownGracefully().await(); + if (!future.isSuccess()) { + AssertLog.error("bossGroup 无法正常停止:{}", future.cause()); + } + future = this.businessGroup.shutdownGracefully().await(); + if (!future.isSuccess()) { + AssertLog.error("businessGroup 无法正常停止:{}", future.cause()); + } + } catch (InterruptedException e) { + e.printStackTrace(); + } + AssertLog.info("\n\n\t\t\t\t\t\t\t\t======{} 已经停止,port:{}======\n", config.name, config.port); + } +} diff --git a/src/main/java/com/tongran/agent/server/netty/annotation/AgentDispatcher.java b/src/main/java/com/tongran/agent/server/netty/annotation/AgentDispatcher.java new file mode 100644 index 0000000..fe4764a --- /dev/null +++ b/src/main/java/com/tongran/agent/server/netty/annotation/AgentDispatcher.java @@ -0,0 +1,47 @@ +package com.tongran.agent.server.netty.annotation; + +import com.tongran.agent.server.core.enums.MsgEnum; +import lombok.AllArgsConstructor; +import lombok.Getter; +import org.springframework.stereotype.Component; + +import java.lang.annotation.*; + +/** + * 指明类为Agent消息 + */ +@Component +@Documented +@Target({ElementType.TYPE}) +@Retention(RetentionPolicy.RUNTIME) +public @interface AgentDispatcher { + + @Getter + @AllArgsConstructor + enum VersionEnum { + V1("2026"); + public final String value; + } + + /** + * 消息ID + * + * @return + */ + MsgEnum msgId(); + + /** + * 消息版本 默认版本2025 + * + * @return + */ + VersionEnum version() default VersionEnum.V1; + + /** + * 描述 + * + * @return + */ + String desc() default ""; + +} diff --git a/src/main/java/com/tongran/agent/server/netty/basics/AgentDispatcherManager.java b/src/main/java/com/tongran/agent/server/netty/basics/AgentDispatcherManager.java new file mode 100644 index 0000000..04014ad --- /dev/null +++ b/src/main/java/com/tongran/agent/server/netty/basics/AgentDispatcherManager.java @@ -0,0 +1,48 @@ +package com.tongran.agent.server.netty.basics; + +import com.tongran.agent.server.netty.annotation.AgentDispatcher; +import org.springframework.beans.BeansException; +import org.springframework.context.ApplicationContext; +import org.springframework.context.ApplicationContextAware; +import org.springframework.core.annotation.Order; +import org.springframework.stereotype.Component; +import org.springframework.util.CollectionUtils; + +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; + +@Component +@Order(1) +public class AgentDispatcherManager implements ApplicationContextAware { + + /** + * 所有实现的包处理器 + */ + private static Map MSG_HANDLER_MAP; + + /** + * 唤醒时 初始化 packHandlerMap + */ + @Override + public void setApplicationContext(ApplicationContext applicationContext) throws BeansException { + // 仅一次性初始化完成 + if (MSG_HANDLER_MAP == null) { + MSG_HANDLER_MAP = new ConcurrentHashMap<>(); + Map handlers = applicationContext.getBeansWithAnnotation(AgentDispatcher.class); + if (!CollectionUtils.isEmpty(handlers)) { + handlers.values().forEach(tempHandler -> { + boolean result = tempHandler.getClass().isAnnotationPresent(AgentDispatcher.class); + if (result) { + AgentDispatcher annotation = tempHandler.getClass().getAnnotation(AgentDispatcher.class); + MSG_HANDLER_MAP.put(annotation.msgId().getValue() + "&" + annotation.version().value, (AgentHandler) tempHandler); + } + }); + } + } + } + + public AgentHandler getHandler(String msgIdVersion) { + return MSG_HANDLER_MAP == null ? null : MSG_HANDLER_MAP.get(msgIdVersion); + } + +} \ No newline at end of file diff --git a/src/main/java/com/tongran/agent/server/netty/basics/AgentHandler.java b/src/main/java/com/tongran/agent/server/netty/basics/AgentHandler.java new file mode 100644 index 0000000..67c3e0e --- /dev/null +++ b/src/main/java/com/tongran/agent/server/netty/basics/AgentHandler.java @@ -0,0 +1,21 @@ +package com.tongran.agent.server.netty.basics; + + +import com.tongran.agent.server.netty.model.UpMsgResponse; + +public interface AgentHandler { + + /** + * 处理终端传入的消息, 然后进行返回 + * + * @return 需要发送给终端的消息 + */ + + default UpMsgResponse upHandle(String data, String clientId) { + return null; + } + + default String downHandle(String data, String clientId) { + return ""; + } +} \ No newline at end of file diff --git a/src/main/java/com/tongran/agent/server/netty/config/BaseNettyConfig.java b/src/main/java/com/tongran/agent/server/netty/config/BaseNettyConfig.java new file mode 100644 index 0000000..83875c4 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/netty/config/BaseNettyConfig.java @@ -0,0 +1,58 @@ +package com.tongran.agent.server.netty.config; + +import com.tongran.agent.server.core.session.SessionManager; +import io.netty.channel.ChannelHandler; +import io.netty.util.NettyRuntime; +import lombok.Data; + +import java.io.Serializable; + +/** + * 基础Netty配置 + */ +@Data +public class BaseNettyConfig implements Serializable { + + private static final long serialVersionUID = -2736718499972710895L; + + /** + * 是否开启 + */ + public Boolean enable = false; + + /** + * 是否TCP + */ + public Boolean isTcp = true; + + /** + * 服务名称 + */ + public String name = "Netty"; + + /** + * 端口号 + */ + public Integer port = 11111; + + public Integer workerCore = NettyRuntime.availableProcessors() * 2; + + public Integer businessCore = Math.max(1, NettyRuntime.availableProcessors() >> 1); + + public Integer readerIdleTime = 240; + + public Integer writerIdleTime = 0; + + public Integer allIdleTime = 0; + + /** + * netty Session管理 + */ + public SessionManager sessionManager; + + /** + * netty Channel处理列表 + */ + public ChannelHandler hander; + +} diff --git a/src/main/java/com/tongran/agent/server/netty/config/ConnectionConfig.java b/src/main/java/com/tongran/agent/server/netty/config/ConnectionConfig.java new file mode 100644 index 0000000..fa2c49c --- /dev/null +++ b/src/main/java/com/tongran/agent/server/netty/config/ConnectionConfig.java @@ -0,0 +1,39 @@ +package com.tongran.agent.server.netty.config; + +import java.util.Objects; + +/** + * 连接配置信息 + */ +public class ConnectionConfig { + private final String connectionKey; + private final String host; + private final int port; + private final int timeout; + + public ConnectionConfig(String connectionKey, String host, int port, int timeout) { + this.connectionKey = connectionKey; + this.host = host; + this.port = port; + this.timeout = timeout; + } + + // Getter方法 + public String getConnectionKey() { return connectionKey; } + public String getHost() { return host; } + public int getPort() { return port; } + public int getTimeout() { return timeout; } + + @Override + public boolean equals(Object o) { + if (this == o) return true; + if (o == null || getClass() != o.getClass()) return false; + ConnectionConfig that = (ConnectionConfig) o; + return Objects.equals(connectionKey, that.connectionKey); + } + + @Override + public int hashCode() { + return Objects.hash(connectionKey); + } +} \ 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 new file mode 100644 index 0000000..ea5b04b --- /dev/null +++ b/src/main/java/com/tongran/agent/server/netty/enpoint/AgentEndpoint.java @@ -0,0 +1,420 @@ +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.annotation.AgentDispatcher; +import com.tongran.agent.server.netty.basics.AgentHandler; +import com.tongran.agent.server.netty.model.UpMsgResponse; +import com.tongran.agent.server.rockermq.AgentProducer; +import com.tongran.agent.server.rockermq.config.AgentMqConfig; +import com.tongran.agent.server.rockermq.model.MqMsg; +import org.springframework.stereotype.Component; + +import javax.annotation.Resource; + +@Component +public class AgentEndpoint { + + @Resource + private AgentMqConfig agentMqConfig; + + @Resource + private AgentProducer agentProducer; + + @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("data", data); + MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(jsonObject.toString()).build(); + agentProducer.asyncSend(mqMsg); + 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; + } + } + + @AgentDispatcher(msgId = MsgEnum.断开) + public class DisconnectHandler implements AgentHandler { + @Override + public String downHandle(String data, String clientId) { + return data; + } + } + + @AgentDispatcher(msgId = MsgEnum.断开应答) + public class DisconnectRspHandler 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); + JSONObject json = JSONObject.parseObject(data); + int resCode = 0; + if(json.containsKey("resCode")){ + resCode = json.getInteger("resCode"); + } +// if(resCode == 1){ +// client.closeConnection(clientId); +// } + return UpMsgResponse.builder().clientId(clientId).build(); + } + } + + @AgentDispatcher(msgId = MsgEnum.心跳上报) + public class HeartBeatHandler 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); + JSONObject json = new JSONObject(); + json.put("resCode",1); + return UpMsgResponse.builder().clientId(clientId).dataType(MsgEnum.心跳上报应答.getValue()).content(json.toString()).build(); + } + } + + @AgentDispatcher(msgId = MsgEnum.CPU上报) + public class CpuHandler implements AgentHandler { + @Override + public UpMsgResponse upHandle(String data, String clientId) { + JSONObject jsonObject = new JSONObject(); + jsonObject.put("clientId", clientId); + jsonObject.put("dataType", MsgEnum.CPU上报.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(); + } + } + + @AgentDispatcher(msgId = MsgEnum.磁盘上报) + public class DiskHandler 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); + 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 DockerHandler 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); + 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 MemoryHandler 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); + 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 NetHandler 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); + 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 PointHandler 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); + 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 SystemHandler 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); + 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 SwitchBoardHandler 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); + 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 SystemCollectStartHandler implements AgentHandler { + @Override + public String downHandle(String data, String clientId) { + return data; + } + } + + @AgentDispatcher(msgId = MsgEnum.开启或更新系统采集应答) + public class SystemCollectStartRspHandler 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 SystemCollectClosetHandler implements AgentHandler { + @Override + public String downHandle(String data, String clientId) { + return data; + } + } + + @AgentDispatcher(msgId = MsgEnum.开启或更新交换机采集) + public class SwitchCollectStartHandler implements AgentHandler { + @Override + public String downHandle(String data, String clientId) { + return data; + } + } + + @AgentDispatcher(msgId = MsgEnum.开启或更新交换机采集应答) + public class SwitchCollectStartRspHandler 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 SwitchCollectClosetHandler implements AgentHandler { + @Override + public String downHandle(String data, String clientId) { + return data; + } + } + + @AgentDispatcher(msgId = MsgEnum.告警设置) + public class AlarmSetHandler implements AgentHandler { + @Override + public String downHandle(String data, String clientId) { + return data; + } + } + + @AgentDispatcher(msgId = MsgEnum.告警设置应答) + public class AlarmSetRspHandler 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 ScriptPolicyHandler implements AgentHandler { + @Override + public String downHandle(String data, String clientId) { + return data; + } + } + + @AgentDispatcher(msgId = MsgEnum.执行脚本策略应答) + public class ScriptPolicyRspHandler implements AgentHandler { + @Override + public UpMsgResponse upHandle(String data, String clientId) { + JSONObject jsonObject = JSONObject.parseObject(data); + 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.Agent版本更新) + public class AgentVersionUpdateHandler implements AgentHandler { + @Override + public String downHandle(String data, String clientId) { + return data; + } + } + + @AgentDispatcher(msgId = MsgEnum.Agent版本更新应答) + public class AgentVersionUpdateRspHandler implements AgentHandler { + @Override + public UpMsgResponse upHandle(String data, String clientId) { + JSONObject jsonObject = new JSONObject(); + jsonObject.put("clientId", clientId); + jsonObject.put("dataType", MsgEnum.Agent版本更新应答.getValue()); + jsonObject.put("data", data); + MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(jsonObject.toString()).build(); + agentProducer.asyncSend(mqMsg); + JSONObject json = JSONObject.parseObject(data); + int resCode = 0; + if(json.containsKey("resCode")){ + resCode = json.getInteger("resCode"); + } +// if(resCode == 1){ +// client.closeConnection(clientId); +// } + return null; + } + } + + @AgentDispatcher(msgId = MsgEnum.多网IP探测上报) + public class DetectNetworkHandler implements AgentHandler { + @Override + public UpMsgResponse upHandle(String data, String clientId) { + JSONObject jsonObject = new JSONObject(); + jsonObject.put("clientId", clientId); + jsonObject.put("dataType", MsgEnum.多网IP探测上报.getValue()); + jsonObject.put("data", data); + MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(jsonObject.toString()).build(); + agentProducer.asyncSend(mqMsg); + return null; + } + } + + + + +} \ No newline at end of file diff --git a/src/main/java/com/tongran/agent/server/netty/handler/AgentDecoderHandler.java b/src/main/java/com/tongran/agent/server/netty/handler/AgentDecoderHandler.java new file mode 100644 index 0000000..cba59b1 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/netty/handler/AgentDecoderHandler.java @@ -0,0 +1,174 @@ +package com.tongran.agent.server.netty.handler; + +import cn.hutool.cache.Cache; +import cn.hutool.cache.CacheUtil; +import cn.hutool.core.date.DateUnit; +import cn.hutool.core.util.ObjectUtil; +import com.alibaba.fastjson2.JSONObject; +import com.tongran.agent.server.core.session.SessionManager; +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.netty.model.UpMsgResponse; +import com.tongran.agent.server.utils.AssertLog; +import io.netty.buffer.ByteBuf; +import io.netty.channel.ChannelHandler; +import io.netty.channel.ChannelHandlerContext; +import io.netty.channel.ChannelInboundHandlerAdapter; +import io.netty.util.AttributeKey; +import io.netty.util.CharsetUtil; +import org.springframework.stereotype.Component; + +import javax.annotation.Resource; +import java.util.Objects; + +@Component +@ChannelHandler.Sharable +public class AgentDecoderHandler extends ChannelInboundHandlerAdapter { + + // 连接标识属性键 + private static final AttributeKey CONNECTION_KEY = AttributeKey.valueOf("connectionKey"); + + protected final SessionManager sessionManager; + + public AgentDecoderHandler() { + this.sessionManager = SessionManager.getInstance(); + } + + @Resource + AgentDispatcherManager agentDispatcherManager; + + // 用来临时保留没有处理过的请求报文 + private final Cache lruCache = CacheUtil.newLRUCache(5000); + + /** + * 消息解码器 + */ + @Override + public void channelRead(ChannelHandlerContext ctx, Object msg) { + // 检查是否为 ByteBuf(未解码的原始数据) + if (msg instanceof ByteBuf) { + ByteBuf byteBuf = (ByteBuf) msg; + String messages = byteBuf.toString(CharsetUtil.UTF_8); // 指定字符集解码 + AssertLog.info("<<[up1]:[up-content]==>{}", messages); + + String content = ""; + StringBuilder sb = new StringBuilder(); + boolean isClear = false; + String tempMsg = lruCache.get(sessionManager.client(ctx)); + tempMsg = tempMsg == null ? "" : tempMsg; + int tmpMsgSize = tempMsg.length(); + + boolean startsWith = messages.startsWith("agent-client:"); + boolean endsWith = messages.endsWith("@tong-ran"); + String ms = messages.replaceAll("agent-client:",""); + String[] arr_msg = ms.split("@tong-ran"); + if(startsWith){ + //判定是否整包 + if(arr_msg.length == 1){ + //判定是否拆包 + if(endsWith){ + lruCache.remove(sessionManager.client(ctx)); + tmpMsgSize = 0; +// sb.append(arr_msg[0]+"@tong-ran"); + content = arr_msg[0]+"@tong-ran"; + }else{ + lruCache.remove(sessionManager.client(ctx)); + tmpMsgSize = 0; +// sb.append(arr_msg[0]); + content = arr_msg[0]; + lruCache.put(sessionManager.client(ctx), content, DateUnit.SECOND.getMillis() * 5000); + } + } + //判定是否粘包 + if(arr_msg.length > 1){ + if(!endsWith){ + lruCache.remove(sessionManager.client(ctx)); + tmpMsgSize = 0; + for (int i = 0; i < arr_msg.length; i++) { + if(i == arr_msg.length -1){ + lruCache.put(sessionManager.client(ctx), arr_msg[i], DateUnit.SECOND.getMillis() * 5000); + }else{ + content += arr_msg[i].replaceAll("agent-client:","")+"@tong-ran"; +// sb.append(arr_msg[i].replaceAll("agent-client:","")+"@tong-ran"); + } + } + }else{ + lruCache.remove(sessionManager.client(ctx)); + tmpMsgSize = 0; +// sb.append(String.join("@tong-ran", arr_msg)); + content = ms; + } + } + } +// content = sb.toString(); + if (tmpMsgSize > 0) { + content = tempMsg + messages; + endsWith = content.endsWith("@tong-ran"); + if(endsWith){ + isClear = true; + }else{ + lruCache.remove(sessionManager.client(ctx)); + tmpMsgSize = 0; + lruCache.put(sessionManager.client(ctx), content, DateUnit.SECOND.getMillis() * 5000); + } + + } + AssertLog.info("<<[up2]:IP:{},[up-content]==>{}", sessionManager.client(ctx),content); + endsWith = content.endsWith("@tong-ran"); + //判断是否是整包 + if(endsWith){ + content = content.replaceAll("agent-client:",""); + String[] arr = content.split("@tong-ran"); + try { + for (String message : arr) { + JSONObject jsonObject = JSONObject.parseObject(message); + String clientId = jsonObject.getString("clientId"); + String dataType = jsonObject.getString("dataType"); + String data = jsonObject.getString("data"); + if(Objects.nonNull(agentDispatcherManager)){ + AgentHandler msgHandler = agentDispatcherManager.getHandler(dataType + "&" + AgentDispatcher.VersionEnum.V1.value); + if (ObjectUtil.isNotEmpty(msgHandler)) { + AssertLog.info("<<[up-after-handle]:clientId:{},type={},[handle-content]={}", clientId, dataType, data); + UpMsgResponse response = msgHandler.upHandle(data, clientId ); + if(Objects.nonNull(response)){ + Message agentMessage = Message.builder().build(); + agentMessage.setClientId(clientId); + agentMessage.setDataType(response.getDataType()); + agentMessage.setData(response.getContent()); + ctx.fireChannelRead(agentMessage);//传递到下一个handler + } + } + } + + } + }catch (Exception e){ + AssertLog.error("=====channelRead:{}=====" + e.getMessage()); + } + } + if (isClear) { + lruCache.remove(sessionManager.client(ctx)); + } + byteBuf.release(); // 释放 ByteBuf 资源(重要!) + } else { + System.out.println("Unexpected message type: " + msg.getClass()); + } + + } + + @Override + public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { + String connectionKey = ctx.channel().attr(CONNECTION_KEY).get(); + System.err.println("连接 " + connectionKey + " 发生异常: " + cause.getMessage()); + ctx.close(); + } + + public byte[] subByte(byte[] b, int off, int length) { + byte[] bytes = new byte[length]; + System.arraycopy(b, off, bytes, 0, length); + return bytes; + } + + +} diff --git a/src/main/java/com/tongran/agent/server/netty/handler/AgentDispatcherHandler.java b/src/main/java/com/tongran/agent/server/netty/handler/AgentDispatcherHandler.java new file mode 100644 index 0000000..4e6282f --- /dev/null +++ b/src/main/java/com/tongran/agent/server/netty/handler/AgentDispatcherHandler.java @@ -0,0 +1,54 @@ +package com.tongran.agent.server.netty.handler; + +import com.tongran.agent.server.core.session.Session; +import com.tongran.agent.server.core.session.SessionManager; +import com.tongran.agent.server.netty.model.Message; +import com.tongran.agent.server.utils.AssertLog; +import io.netty.channel.ChannelHandler; +import io.netty.channel.ChannelHandlerContext; +import io.netty.channel.SimpleChannelInboundHandler; +import org.apache.commons.lang3.StringUtils; +import org.springframework.stereotype.Component; + +@Component +@ChannelHandler.Sharable +public class AgentDispatcherHandler extends SimpleChannelInboundHandler { + + protected final SessionManager sessionManager; + + public AgentDispatcherHandler() { + this.sessionManager = SessionManager.getInstance(); + } + + @Override + protected void channelRead0(ChannelHandlerContext ctx, Message msg) { +// AssertLog.info(">>>>>>[要处理的终端数据] {}", msg.toString()); + String clientId = msg.getClientId(); + if (StringUtils.isBlank(clientId)) { + AssertLog.error("<<<<<<错误的信息from:{}", ctx.channel().remoteAddress()); + return; + } + Session session = Session.buildSession(ctx, clientId); + sessionManager.put(clientId, session); + if (StringUtils.isNotBlank(msg.getData())) { +// AssertLog.info(">>>>>>[getSessionById的终端数据] {}", sessionManager.getSessionById(clientId)); + sessionManager.writeAndFlush(sessionManager.getSessionById(clientId).getChannel(), msg); + } + } + +// @Override +// public void channelActive(ChannelHandlerContext ctx) { +// AssertLog.info("与服务器建立连接: {}", ctx.channel().remoteAddress()); +// // 连接建立后可以发送登录认证消息 +// long timestamp = System.currentTimeMillis(); +// String clientId = "clientId001"; +// JSONObject object = new JSONObject(); +// object.put("uuid",clientId); +// object.put("timestamp",timestamp); +// Message message = Message.builder().clientId(clientId).dataType("LOGIN").data(object.toString()).build(); +// // 将对象转为 JSON 字符串 +// String json = JSON.toJSONString(message); +// System.out.println("连接建立成功,发送登录认证="+json); +// ctx.writeAndFlush(message); +// } +} diff --git a/src/main/java/com/tongran/agent/server/netty/handler/AgentEncoderHandler.java b/src/main/java/com/tongran/agent/server/netty/handler/AgentEncoderHandler.java new file mode 100644 index 0000000..cdc5426 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/netty/handler/AgentEncoderHandler.java @@ -0,0 +1,49 @@ +package com.tongran.agent.server.netty.handler; + +import com.alibaba.fastjson2.JSON; +import com.tongran.agent.server.netty.model.Message; +import com.tongran.agent.server.core.session.SessionManager; +import com.tongran.agent.server.utils.AssertLog; +import io.netty.buffer.ByteBuf; +import io.netty.buffer.Unpooled; +import io.netty.channel.ChannelHandler; +import io.netty.channel.ChannelHandlerContext; +import io.netty.channel.ChannelOutboundHandlerAdapter; +import io.netty.channel.ChannelPromise; +import org.apache.commons.lang3.StringUtils; +import org.springframework.stereotype.Component; + +import java.nio.charset.StandardCharsets; + +@Component +@ChannelHandler.Sharable +public class AgentEncoderHandler extends ChannelOutboundHandlerAdapter { + + protected final SessionManager sessionManager; + + public AgentEncoderHandler() { + this.sessionManager = SessionManager.getInstance(); + } + + /** + * 消息解码器 + */ + @Override + public void write(ChannelHandlerContext ctx, Object msg, ChannelPromise promise) throws Exception { + if (msg instanceof Message) { + Message entity = (Message) msg; + if (StringUtils.isBlank(entity.getData())) { + AssertLog.error(">>[down]:IP:{},errorContent:{}", sessionManager.client(ctx), msg); + return; + } + AssertLog.info(">>[down]:IP:{},[content]==>{}", sessionManager.client(ctx), entity.getData()); +// String json = JSON.toJSONString(entity); + String json = "agent-server:"+JSON.toJSONString(entity)+"@tong-ran"; + byte[] bytes = json.getBytes(StandardCharsets.UTF_8); // 显式指定 UTF-8 +// byte[] bytes = EscapeUtil.hexStringToByteArray(entity.getContent()); + ByteBuf buf = Unpooled.wrappedBuffer(bytes); + ctx.write(buf, promise); + } + } + +} 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/netty/model/Message.java b/src/main/java/com/tongran/agent/server/netty/model/Message.java new file mode 100644 index 0000000..5cc96bd --- /dev/null +++ b/src/main/java/com/tongran/agent/server/netty/model/Message.java @@ -0,0 +1,37 @@ +package com.tongran.agent.server.netty.model; + + +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; + +import java.io.Serializable; + +/** + * @author egrias + */ +@Data +@Builder +@AllArgsConstructor +@NoArgsConstructor +public class Message implements Serializable { + + private static final long serialVersionUID = -1267013167162440610L; + + /** + * clientId + */ + private String clientId; + + /** + * 数据类型:LOGIN、HEARTBEAT、CPU、MEMORY、SYSTEM、POINT、NET、DISK、DOCKER、SWITCHBOARD + */ + private String dataType; + + /** + * 发送内容 + */ + private String data; + +} diff --git a/src/main/java/com/tongran/agent/server/netty/model/UpMsgResponse.java b/src/main/java/com/tongran/agent/server/netty/model/UpMsgResponse.java new file mode 100644 index 0000000..8933661 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/netty/model/UpMsgResponse.java @@ -0,0 +1,22 @@ +package com.tongran.agent.server.netty.model; + +import lombok.AllArgsConstructor; +import lombok.Data; +import lombok.NoArgsConstructor; +import lombok.experimental.SuperBuilder; + +@Data +@SuperBuilder +@AllArgsConstructor +@NoArgsConstructor +public class UpMsgResponse { + + private String clientId; + /** + * type + */ + private String dataType; + + private String content; + +} diff --git a/src/main/java/com/tongran/agent/server/rockermq/AgentProducer.java b/src/main/java/com/tongran/agent/server/rockermq/AgentProducer.java new file mode 100644 index 0000000..4531f53 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/rockermq/AgentProducer.java @@ -0,0 +1,54 @@ +package com.tongran.agent.server.rockermq; + + +import com.tongran.agent.server.rockermq.model.MqMsg; +import com.tongran.agent.server.utils.AssertLog; +import org.apache.rocketmq.client.producer.SendCallback; +import org.apache.rocketmq.client.producer.SendResult; +import org.apache.rocketmq.spring.core.RocketMQTemplate; +import org.springframework.messaging.support.MessageBuilder; +import org.springframework.stereotype.Component; + +import javax.annotation.Resource; + +@Component +public class AgentProducer implements RocketMqService { + + @Resource + private RocketMQTemplate rocketMQTemplate; + + @Override + public void send(MqMsg msg) { + AssertLog.info("send发送消息:==>{}", msg); + rocketMQTemplate.send(msg.getTopic(), MessageBuilder.withPayload(msg.getContent()).build()); + } + + @Override + public void asyncSend(MqMsg msg) { + AssertLog.info("asyncSend发送消息:==>{}", msg); + rocketMQTemplate.asyncSend(msg.getTopic(), msg.getContent(), new SendCallback() { + @Override + public void onSuccess(SendResult sendResult) { +// AssertLog.info("asyncSend发送消息成功:==>{}", sendResult.getSendStatus()); + } + + @Override + public void onException(Throwable throwable) { + AssertLog.error("asyncSend发送消息失败:==>{}", throwable.getMessage()); + } + }); + } + + @Override + public void syncSendOrderly(MqMsg msg) { + AssertLog.info("syncSendOrderly发送消息:==>{}", msg); + rocketMQTemplate.sendOneWay(msg.getTopic(), msg.getContent()); + } + + @Override + public void delayedSendOrderly(MqMsg msg) { + AssertLog.info("delayedSendOrderly发送消息:==>{}", msg); + rocketMQTemplate.syncSend(msg.getTopic(), MessageBuilder.withPayload(msg.getContent()).build(), 30000, msg.getLevel()); + } + +} diff --git a/src/main/java/com/tongran/agent/server/rockermq/RocketMqService.java b/src/main/java/com/tongran/agent/server/rockermq/RocketMqService.java new file mode 100644 index 0000000..6cb5894 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/rockermq/RocketMqService.java @@ -0,0 +1,45 @@ +package com.tongran.agent.server.rockermq; + + +import com.tongran.agent.server.rockermq.model.MqMsg; + +/** + * Rocket MQ 对应服务封装 + * @author BAO + * + */ +public interface RocketMqService { + + /** + * 同步发送消息
+ *

+ * 当发送的消息很重要是,且对响应时间不敏感的时候采用sync方式; + * + * @param msg 发送消息实体类 + */ + void send(MqMsg msg); + + /** + * 异步发送消息,异步返回消息结果 当发送的消息很重要,且对响应时间非常敏感的时候采用async方式; + * + * @param msg 发送消息实体类 + */ + void asyncSend(MqMsg msg); + + /** + * 直接发送发送消息,不关心返回结果,容易消息丢失,适合日志收集、不精确统计等消息发送; 发送的消息不重要时,采用one-way方式,以提高吞吐量; + * + * @param msg 发送消息实体类 + */ + void syncSendOrderly(MqMsg msg); + + /** + * 发送延时消息,不关心返回结果
+ *

+ * 当发送的消息不重要时,采用one-way方式,以提高吞吐量; + * + * @param msg 发送消息实体类 + */ + default void delayedSendOrderly(MqMsg msg){}; + +} \ No newline at end of file diff --git a/src/main/java/com/tongran/agent/server/rockermq/config/AgentMqConfig.java b/src/main/java/com/tongran/agent/server/rockermq/config/AgentMqConfig.java new file mode 100644 index 0000000..474771c --- /dev/null +++ b/src/main/java/com/tongran/agent/server/rockermq/config/AgentMqConfig.java @@ -0,0 +1,14 @@ +package com.tongran.agent.server.rockermq.config; + +import lombok.Data; +import org.springframework.boot.context.properties.ConfigurationProperties; +import org.springframework.context.annotation.Configuration; + +@Data +@Configuration +@ConfigurationProperties(prefix = "rocketmq.producer") +public class AgentMqConfig { + + private String agentTopic; + +} 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 new file mode 100644 index 0000000..3674d32 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/rockermq/listener/AgentConsumerDownListener.java @@ -0,0 +1,51 @@ +package com.tongran.agent.server.rockermq.listener; + +import cn.hutool.core.util.ObjectUtil; +import com.alibaba.fastjson2.JSON; +import com.tongran.agent.server.core.session.SessionManager; +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.utils.AssertLog; +import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; +import org.apache.rocketmq.spring.core.RocketMQListener; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.stereotype.Component; + +import javax.annotation.Resource; +import java.util.Objects; + +@Component +@ConditionalOnProperty(name = "rocketmq.enabled", havingValue = "true") +@RocketMQMessageListener(topic = "tr_agent_down", consumerGroup = "tr_agent_down_group") +public class AgentConsumerDownListener implements RocketMQListener { + + protected final SessionManager sessionManager; + + @Resource + private AgentDispatcherManager agentDispatcherManager; + + public AgentConsumerDownListener() { + this.sessionManager = SessionManager.getInstance(); + } + + @Override + public void onMessage(String message) { + AssertLog.info("consumer==> received down message: {}", message); + try { + Message dto = JSON.parseObject(message, Message.class); + 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/rockermq/model/MqMsg.java b/src/main/java/com/tongran/agent/server/rockermq/model/MqMsg.java new file mode 100644 index 0000000..398c894 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/rockermq/model/MqMsg.java @@ -0,0 +1,38 @@ +package com.tongran.agent.server.rockermq.model; + +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; + +import java.io.Serializable; + +/** + * RocketMQ消息发送对象 + * + * @author BAO + */ +@Data +@Builder +@AllArgsConstructor +@NoArgsConstructor +public class MqMsg implements Serializable { + + private static final long serialVersionUID = 4164379745748817325L; + + /** + * 消息topic + */ + private String topic; + + /** + * 消息content + */ + private Object content; + + /** + * 消息延时等级 + */ + private Integer level; + +} \ No newline at end of file diff --git a/src/main/java/com/tongran/agent/server/service/AgentService.java b/src/main/java/com/tongran/agent/server/service/AgentService.java new file mode 100644 index 0000000..b7d962a --- /dev/null +++ b/src/main/java/com/tongran/agent/server/service/AgentService.java @@ -0,0 +1,10 @@ +package com.tongran.agent.server.service; + +import com.tongran.agent.server.core.vo.RegisterVO; + +public interface AgentService { + + boolean createConn(String clientId, String host, int port, int timeout); + + void register(RegisterVO registerVO); +} 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 new file mode 100644 index 0000000..fd289fc --- /dev/null +++ b/src/main/java/com/tongran/agent/server/service/impl/AgentServiceImpl.java @@ -0,0 +1,55 @@ +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.model.Message; +import com.tongran.agent.server.service.AgentService; +import org.springframework.stereotype.Service; + +@Service +public class AgentServiceImpl implements AgentService { + +// @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; + return false; + } + + @Override + public void register(RegisterVO registerVO) { + // 发送注册 + JSONObject object = new JSONObject(); + object.put("clientId", registerVO.getClientId()); + object.put("timestamp", registerVO.getTimestamp()); + JSONObject switchBoard = new JSONObject(); + switchBoard.put("community", registerVO.getSwitchBoard().getCommunity()); + switchBoard.put("ip", registerVO.getSwitchBoard().getIp()); + switchBoard.put("port", registerVO.getSwitchBoard().getPort()); + switchBoard.put("oid", String.join(",", registerVO.getSwitchBoard().getFanOID())); + object.put("switchBoard", switchBoard.toString()); + Message message = Message.builder().clientId(registerVO.getClientId()).dataType("REGISTER").data(object.toString()).build(); + // 将对象转为 JSON 字符串 + String json = JSON.toJSONString(message); + System.out.println("发送注册消息="+json); +// client.sendMessages(registerVO.getClientId(), message); + } +} diff --git a/src/main/java/com/tongran/agent/server/utils/AssertLog.java b/src/main/java/com/tongran/agent/server/utils/AssertLog.java new file mode 100644 index 0000000..2423570 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/utils/AssertLog.java @@ -0,0 +1,73 @@ +package com.tongran.agent.server.utils; + +import lombok.extern.slf4j.Slf4j; + +/** + * 日志断言类 + * + * @author BAO + * + */ +@Slf4j +public class AssertLog { + + /** + * 打印Info 日志 + * + * @param format + * @param arguments + */ + public static void info(String format, Object... arguments) { + if (log.isInfoEnabled()) { + log.info(format, arguments); + } + } + + /** + * 打印Debug 日志 + * + * @param format + * @param arguments + */ + public static void debug(String format, Object... arguments) { + if (log.isDebugEnabled()) { + log.debug(format, arguments); + } + } + + /** + * 打印Error 日志 + * + * @param format + * @param arguments + */ + public static void error(String format, Object... arguments) { + if (log.isErrorEnabled()) { + log.error(format, arguments); + } + } + + /** + * 打印Trace 日志 + * + * @param format + * @param arguments + */ + public static void trace(String format, Object... arguments) { + if (log.isTraceEnabled()) { + log.trace(format, arguments); + } + } + + /** + * 打印Warn 日志 + * + * @param format + * @param arguments + */ + public static void warn(String format, Object... arguments) { + if (log.isWarnEnabled()) { + log.warn(format, arguments); + } + } +} diff --git a/src/main/java/com/tongran/agent/server/utils/ClientExample.java b/src/main/java/com/tongran/agent/server/utils/ClientExample.java new file mode 100644 index 0000000..84eab55 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/utils/ClientExample.java @@ -0,0 +1,46 @@ +package com.tongran.agent.server.utils; + +import com.tongran.agent.server.netty.config.ConnectionConfig; + +import java.util.Arrays; +import java.util.List; +import java.util.Map; + +public class ClientExample { + public static void main(String[] args) { + EnhancedConnectionManager manager = new EnhancedConnectionManager(); + + try { + // 配置多个连接 + List configs = Arrays.asList( +// new ConnectionConfig("db_server", "192.168.1.101", 3306, 5), +// new ConnectionConfig("redis_server", "192.168.1.102", 6379, 3), +// new ConnectionConfig("api_server", "api.example.com", 8080, 10), + new ConnectionConfig("server1", "127.0.0.1", 6610, 5) + ); + + // 批量创建连接 + Map results = manager.createConnections(configs); + System.out.println("连接创建结果: " + results); + + // 等待连接建立 + Thread.sleep(3000); + + // 根据消息类型路由 + manager.routeMessageByType("TYPE_A", "Database query"); +// manager.routeMessageByType("TYPE_B", "Cache operation"); +// manager.routeMessageByType("UNKNOWN", "Broadcast message"); + + // 检查连接状态 + Map status = manager.getConnectionStatus(); + System.out.println("连接状态: " + status); + + Thread.sleep(5000); + + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } finally { +// 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 new file mode 100644 index 0000000..d7060d8 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/utils/EnhancedConnectionManager.java @@ -0,0 +1,81 @@ +package com.tongran.agent.server.utils; + +import com.tongran.agent.server.netty.config.ConnectionConfig; + +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; + +/** + * 增强版连接管理器 + */ +public class EnhancedConnectionManager { +// private final MultiTargetNettyClient client; + private final Map connectionConfigs; + + public EnhancedConnectionManager() { +// this.client = new MultiTargetNettyClient(); + this.connectionConfigs = new ConcurrentHashMap<>(); + } + + /** + * 批量创建连接 + */ + 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); +// }); + + return results; + } + + /** + * 根据消息类型路由到不同连接 + */ + 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)); +// } + } + + /** + * 获取所有连接状态 + */ + public Map getConnectionStatus() { + Map status = new HashMap<>(); +// connectionConfigs.forEach((key, config) -> { +// // 这里可以添加更详细的状态检查 +// status.put(key, client.sendMessage(key, "PING")); +// }); + return status; + } + +// public void shutdown() { +// client.shutdown(); +// } +} \ No newline at end of file diff --git a/src/main/java/com/tongran/agent/server/utils/R.java b/src/main/java/com/tongran/agent/server/utils/R.java new file mode 100644 index 0000000..e566d72 --- /dev/null +++ b/src/main/java/com/tongran/agent/server/utils/R.java @@ -0,0 +1,76 @@ +package com.tongran.agent.server.utils; + +import com.tongran.agent.server.exception.code.ErrorCode; +import com.tongran.agent.server.exception.code.GlobalErrorCode; +import io.swagger.v3.oas.annotations.media.Schema; +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; + +import java.io.Serializable; + +/** + * @author BAO + */ +@Data +@Builder +@NoArgsConstructor +@AllArgsConstructor +@Schema(description = "接口交互统一数据返回标准") +public class R implements Serializable { + + private static final long serialVersionUID = 1L; + + @Schema(description = "返回代码") + private Integer code; + + @Schema(description = "消息描述") + private String msg; + + @Schema(description = "结果对象") + private T data; + + public static R success(String msg, T t) { + R r = new R<>(); + r.setData(t); + r.setMsg(msg); + r.setCode(GlobalErrorCode.SUCCESS.getCode()); + return r; + } + + public static R success(T t) { + return R.success(GlobalErrorCode.SUCCESS.getMsg(), t); + } + + public static R success() { + return R.success(null); + } + + public static R error(String msg, Integer code) { + R r = new R<>(); + r.setMsg(msg); + r.setCode(code); + return r; + } + + public static R error(Integer code, String msg) { + R r = new R<>(); + r.setMsg(msg); + r.setCode(code); + return r; + } + + public static R error(ErrorCode err) { + return R.error(err.getMsg(), err.getCode()); + } + + public static R error() { + return R.error(GlobalErrorCode.ERROR.getMsg(), GlobalErrorCode.ERROR.getCode()); + } + + public static R error(String msg) { + return R.error(msg, GlobalErrorCode.ERROR.getCode()); + } + +} \ No newline at end of file diff --git a/src/main/resources/application-dev.yml b/src/main/resources/application-dev.yml new file mode 100644 index 0000000..4f8ca1f --- /dev/null +++ b/src/main/resources/application-dev.yml @@ -0,0 +1,29 @@ +server: + port: -1 + servlet: + context-path: /tr-agent-server + +# 接口文档配置 +knife4j: + enable: true + production: false # 开启屏蔽文档资源 + +# 日志配置 +logging: + file: +# path: /usr/local/tongran/logs + path: /usr/local/tongran_server/logs +# path: D:/job/agent-logs/logs + +rocketmq: + name-server: 172.16.15.103:9876 + enabled: true + +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 new file mode 100644 index 0000000..b504121 --- /dev/null +++ b/src/main/resources/application.yml @@ -0,0 +1,43 @@ +spring: + profiles: + active: dev + mvc: + pathmatch: + matching-strategy: ant_path_matcher + application: + name: tr-agent-server + version: 1.0 + web: + resources: + static-locations: classpath*:/META-INF/resources/ +logging: + config: classpath:logback/logback-${spring.profiles.active}.xml + +# springdoc-openapi项目配置 +knife4j: + setting: + enable-footer-custom: true + footer-custom-content: Apache License 2.0 +springdoc: + config: + title: AGENT服务 + description: AGENT服务接口文档 + contact: SERVER + email: + version: ${spring.application.version} + group-configs: + - group: 'AGENT服务' + paths-to-match: '/**' + packages-to-scan: com.tongran.agent.server + +rocketmq: + producer: + group: tongran-group + topic: default_topic + agent-group: tr_agent_up_group + agent-topic: tr_agent_up + consumer: + group: tongran-group + topic: default_topic + agent-group: tr_agent_down_group + agent-topic: tr_agent_down \ No newline at end of file diff --git a/src/main/resources/logback/logback-dev.xml b/src/main/resources/logback/logback-dev.xml new file mode 100644 index 0000000..cf1e7db --- /dev/null +++ b/src/main/resources/logback/logback-dev.xml @@ -0,0 +1,116 @@ + + + + + ${appName} + + + + + + + [requestId:%X{requestId}] [time:%d{yyyy-MM-dd HH:mm:ss.SSS}] [thread:%thread] [level:%-5level] [appName:${appName}] msg- %msg%n + + + + + + + ${LOG_HOME}/debug.log + + + + ${LOG_HOME}/debug.%d{yyyy-MM-dd}.%i.AssertLog.zip + + 30 + + 10GB + + + + 10MB + + + + + + [requestId:%X{requestId}] [time:%d{yyyy-MM-dd HH:mm:ss.SSS}] [thread:%thread] [level:%-5level] [appName:${appName}] msg- %msg%n + + + + + + + + INFO + ACCEPT + DENY + + + ${LOG_HOME}/info.log + + + + ${LOG_HOME}/info.%d{yyyy-MM-dd}.%i.AssertLog.zip + + 30 + + 10GB + + + + 10MB + + + + + + [requestId:%X{requestId}] [time:%d{yyyy-MM-dd HH:mm:ss.SSS}] [thread:%thread] [level:%-5level] [appName:${appName}] msg- %msg%n + + + + + + + + ERROR + ACCEPT + DENY + + ${LOG_HOME}/error.log + + + + ${LOG_HOME}/error.%d{yyyy-MM-dd}.%i.AssertLog.zip + + 30 + + 10GB + + + + 10MB + + + + + + [requestId:%X{requestId}] [time:%d{yyyy-MM-dd HH:mm:ss.SSS}] [thread:%thread] [level:%-5level] [appName:${appName}] msg- %msg%n + + + + + + + + + + + + + + + + + + \ No newline at end of file diff --git a/src/main/resources/logback/logback-prod.xml b/src/main/resources/logback/logback-prod.xml new file mode 100644 index 0000000..859494f --- /dev/null +++ b/src/main/resources/logback/logback-prod.xml @@ -0,0 +1,115 @@ + + + + + ${appName} + + + + + + + [requestId:%X{requestId}] [time:%d{yyyy-MM-dd HH:mm:ss.SSS}] [thread:%thread] [level:%-5level] [appName:${appName}] msg- %msg%n + + + + + + + ${LOG_HOME}/debug.log + + + + ${LOG_HOME}/debug.%d{yyyy-MM-dd}.%i.AssertLog.zip + + 30 + + 10GB + + + + 10MB + + + + + + [requestId:%X{requestId}] [time:%d{yyyy-MM-dd HH:mm:ss.SSS}] [thread:%thread] [level:%-5level] [appName:${appName}] msg- %msg%n + + + + + + + + INFO + ACCEPT + DENY + + + ${LOG_HOME}/info.log + + + + ${LOG_HOME}/info.%d{yyyy-MM-dd}.%i.AssertLog.zip + + 30 + + 10GB + + + + 100MB + + + + + + [requestId:%X{requestId}] [time:%d{yyyy-MM-dd HH:mm:ss.SSS}] [thread:%thread] [level:%-5level] [appName:${appName}] msg- %msg%n + + + + + + + + ERROR + ACCEPT + DENY + + ${LOG_HOME}/error.log + + + + ${LOG_HOME}/error.%d{yyyy-MM-dd}.%i.AssertLog.zip + + 30 + + 10GB + + + + 10MB + + + + + + [requestId:%X{requestId}] [time:%d{yyyy-MM-dd HH:mm:ss.SSS}] [thread:%thread] [level:%-5level] [appName:${appName}] msg- %msg%n + + + + + + + + + + + + + + + + + \ No newline at end of file