增加指令更新采集间隔

This commit is contained in:
baoqm
2025-08-23 23:54:23 +08:00
parent 6c80850516
commit 94856a31e8
6 changed files with 263 additions and 37 deletions
@@ -38,7 +38,7 @@ public class AgentConsumerDownListener implements RocketMQListener<String> {
AgentMsgHandler msgHandler = agentDispatcherManager.getHandler(dto.getDataType() + "&" + AgentDispatcher.VersionEnum.V1.value);
if (ObjectUtil.isNotEmpty(msgHandler)) {
String msg = msgHandler.downHandle(dto.getData());
Message chargeMessage = Message.builder().clientId(dto.getClientId()).data(msg).build();
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);
}
@@ -11,7 +11,7 @@ public interface AgentMsgHandler {
* @return 需要发送给终端的消息
*/
default UpMsgResponse upHandle(String data) {
default UpMsgResponse upHandle(String data, String clientId) {
return null;
}
@@ -24,10 +24,10 @@ public class AgentEndpoint {
@AgentDispatcher(msgId = MsgEnum.登录认证)
public class LoginHandler implements AgentMsgHandler {
@Override
public UpMsgResponse upHandle(String data) {
public UpMsgResponse upHandle(String data, String clientId) {
JSONObject jsonObject = JSONObject.parseObject(data);
String clientId = jsonObject.getString("clientId");
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(data).build();
jsonObject.put("dataType",MsgEnum.登录认证.getValue());
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(jsonObject.toString()).build();
agentProducer.asyncSend(mqMsg);
JSONObject json = new JSONObject();
json.put("resCode",1);
@@ -59,10 +59,10 @@ public class AgentEndpoint {
@AgentDispatcher(msgId = MsgEnum.心跳包)
public class HeartBeatHandler implements AgentMsgHandler {
@Override
public UpMsgResponse upHandle(String data) {
public UpMsgResponse upHandle(String data, String clientId) {
JSONObject jsonObject = JSONObject.parseObject(data);
String clientId = jsonObject.getString("clientId");
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(data).build();
jsonObject.put("dataType",MsgEnum.心跳包.getValue());
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(jsonObject.toString()).build();
agentProducer.asyncSend(mqMsg);
// T0x04 obj = new T0x04();
@@ -111,10 +111,10 @@ public class AgentEndpoint {
@AgentDispatcher(msgId = MsgEnum.CPU包)
public class CpuHandler implements AgentMsgHandler {
@Override
public UpMsgResponse upHandle(String data) {
public UpMsgResponse upHandle(String data, String clientId) {
JSONObject jsonObject = JSONObject.parseObject(data);
String clientId = jsonObject.getString("clientId");
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(data).build();
jsonObject.put("dataType",MsgEnum.CPU包.getValue());
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(jsonObject.toString()).build();
agentProducer.asyncSend(mqMsg);
JSONObject json = new JSONObject();
json.put("resCode",1);
@@ -135,10 +135,10 @@ public class AgentEndpoint {
@AgentDispatcher(msgId = MsgEnum.磁盘包)
public class DiskHandler implements AgentMsgHandler {
@Override
public UpMsgResponse upHandle(String data) {
public UpMsgResponse upHandle(String data, String clientId) {
JSONObject jsonObject = JSONObject.parseObject(data);
String clientId = jsonObject.getString("clientId");
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(data).build();
jsonObject.put("dataType",MsgEnum.CPU包.getValue());
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(jsonObject.toString()).build();
agentProducer.asyncSend(mqMsg);
JSONObject json = new JSONObject();
json.put("resCode",1);
@@ -159,10 +159,10 @@ public class AgentEndpoint {
@AgentDispatcher(msgId = MsgEnum.内存包)
public class MemoryHandler implements AgentMsgHandler {
@Override
public UpMsgResponse upHandle(String data) {
public UpMsgResponse upHandle(String data, String clientId) {
JSONObject jsonObject = JSONObject.parseObject(data);
String clientId = jsonObject.getString("clientId");
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(data).build();
jsonObject.put("dataType",MsgEnum.内存包.getValue());
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(jsonObject.toString()).build();
agentProducer.asyncSend(mqMsg);
JSONObject json = new JSONObject();
json.put("resCode",1);
@@ -183,10 +183,10 @@ public class AgentEndpoint {
@AgentDispatcher(msgId = MsgEnum.网卡包)
public class NetHandler implements AgentMsgHandler {
@Override
public UpMsgResponse upHandle(String data) {
public UpMsgResponse upHandle(String data, String clientId) {
JSONObject jsonObject = JSONObject.parseObject(data);
String clientId = jsonObject.getString("clientId");
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(data).build();
jsonObject.put("dataType",MsgEnum.网卡包.getValue());
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(jsonObject.toString()).build();
agentProducer.asyncSend(mqMsg);
JSONObject json = new JSONObject();
json.put("resCode",1);
@@ -207,10 +207,10 @@ public class AgentEndpoint {
@AgentDispatcher(msgId = MsgEnum.挂载点包)
public class PointHandler implements AgentMsgHandler {
@Override
public UpMsgResponse upHandle(String data) {
public UpMsgResponse upHandle(String data, String clientId) {
JSONObject jsonObject = JSONObject.parseObject(data);
String clientId = jsonObject.getString("clientId");
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(data).build();
jsonObject.put("dataType",MsgEnum.挂载点包.getValue());
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(jsonObject.toString()).build();
agentProducer.asyncSend(mqMsg);
JSONObject json = new JSONObject();
json.put("resCode",1);
@@ -231,10 +231,10 @@ public class AgentEndpoint {
@AgentDispatcher(msgId = MsgEnum.系统包)
public class SystemHandler implements AgentMsgHandler {
@Override
public UpMsgResponse upHandle(String data) {
public UpMsgResponse upHandle(String data, String clientId) {
JSONObject jsonObject = JSONObject.parseObject(data);
String clientId = jsonObject.getString("clientId");
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(data).build();
jsonObject.put("dataType",MsgEnum.系统包.getValue());
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(jsonObject.toString()).build();
agentProducer.asyncSend(mqMsg);
JSONObject json = new JSONObject();
json.put("resCode",1);
@@ -255,10 +255,10 @@ public class AgentEndpoint {
@AgentDispatcher(msgId = MsgEnum.交换机包)
public class SwitchBoardHandler implements AgentMsgHandler {
@Override
public UpMsgResponse upHandle(String data) {
public UpMsgResponse upHandle(String data, String clientId) {
JSONObject jsonObject = JSONObject.parseObject(data);
String clientId = jsonObject.getString("clientId");
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(data).build();
jsonObject.put("dataType",MsgEnum.交换机包.getValue());
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(jsonObject.toString()).build();
agentProducer.asyncSend(mqMsg);
JSONObject json = new JSONObject();
json.put("resCode",1);
@@ -276,7 +276,197 @@ public class AgentEndpoint {
}
}
@AgentDispatcher(msgId = MsgEnum.更新CPU采集间隔)
public class TimeCpuHandler implements AgentMsgHandler {
@Override
public String downHandle(String data) {
JSONObject jsonObject = JSONObject.parseObject(data);
long intervalMillis = jsonObject.getLong("intervalMillis");
JSONObject json = new JSONObject();
json.put("intervalMillis",intervalMillis);
return json.toString();
}
}
@AgentDispatcher(msgId = MsgEnum.更新CPU采集间隔应答)
public class TimeCpuRspHandler implements AgentMsgHandler {
@Override
public UpMsgResponse upHandle(String data, String clientId) {
JSONObject jsonObject = JSONObject.parseObject(data);
jsonObject.put("dataType",MsgEnum.更新CPU采集间隔应答.getValue());
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(jsonObject.toString()).build();
agentProducer.asyncSend(mqMsg);
return null;
}
}
@AgentDispatcher(msgId = MsgEnum.更新容器采集间隔)
public class TimeDockerHandler implements AgentMsgHandler {
@Override
public String downHandle(String data) {
JSONObject jsonObject = JSONObject.parseObject(data);
long intervalMillis = jsonObject.getLong("intervalMillis");
JSONObject json = new JSONObject();
json.put("intervalMillis",intervalMillis);
return json.toString();
}
}
@AgentDispatcher(msgId = MsgEnum.更新容器采集间隔应答)
public class TimeDockerRspHandler implements AgentMsgHandler {
@Override
public UpMsgResponse upHandle(String data, String clientId) {
JSONObject jsonObject = JSONObject.parseObject(data);
jsonObject.put("dataType",MsgEnum.更新容器采集间隔应答.getValue());
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(jsonObject.toString()).build();
agentProducer.asyncSend(mqMsg);
return null;
}
}
@AgentDispatcher(msgId = MsgEnum.更新网卡采集间隔)
public class TimeNetHandler implements AgentMsgHandler {
@Override
public String downHandle(String data) {
JSONObject jsonObject = JSONObject.parseObject(data);
long intervalMillis = jsonObject.getLong("intervalMillis");
JSONObject json = new JSONObject();
json.put("intervalMillis",intervalMillis);
return json.toString();
}
}
@AgentDispatcher(msgId = MsgEnum.更新网卡采集间隔应答)
public class TimeNetRspHandler implements AgentMsgHandler {
@Override
public UpMsgResponse upHandle(String data, String clientId) {
JSONObject jsonObject = JSONObject.parseObject(data);
jsonObject.put("dataType",MsgEnum.更新网卡采集间隔应答.getValue());
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(jsonObject.toString()).build();
agentProducer.asyncSend(mqMsg);
return null;
}
}
@AgentDispatcher(msgId = MsgEnum.更新交换机采集间隔)
public class TimeSwitchHandler implements AgentMsgHandler {
@Override
public String downHandle(String data) {
JSONObject jsonObject = JSONObject.parseObject(data);
long intervalMillis = jsonObject.getLong("intervalMillis");
JSONObject json = new JSONObject();
json.put("intervalMillis",intervalMillis);
return json.toString();
}
}
@AgentDispatcher(msgId = MsgEnum.更新交换机采集间隔应答)
public class TimeSwitchRspHandler implements AgentMsgHandler {
@Override
public UpMsgResponse upHandle(String data, String clientId) {
JSONObject jsonObject = JSONObject.parseObject(data);
jsonObject.put("dataType",MsgEnum.更新交换机采集间隔应答.getValue());
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(jsonObject.toString()).build();
agentProducer.asyncSend(mqMsg);
return null;
}
}
@AgentDispatcher(msgId = MsgEnum.更新磁盘采集间隔)
public class TimeDiskHandler implements AgentMsgHandler {
@Override
public String downHandle(String data) {
JSONObject jsonObject = JSONObject.parseObject(data);
long intervalMillis = jsonObject.getLong("intervalMillis");
JSONObject json = new JSONObject();
json.put("intervalMillis",intervalMillis);
return json.toString();
}
}
@AgentDispatcher(msgId = MsgEnum.更新磁盘采集间隔应答)
public class TimeDiskRspHandler implements AgentMsgHandler {
@Override
public UpMsgResponse upHandle(String data, String clientId) {
JSONObject jsonObject = JSONObject.parseObject(data);
jsonObject.put("dataType",MsgEnum.更新磁盘采集间隔应答.getValue());
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(jsonObject.toString()).build();
agentProducer.asyncSend(mqMsg);
return null;
}
}
@AgentDispatcher(msgId = MsgEnum.更新挂载采集间隔)
public class TimePointHandler implements AgentMsgHandler {
@Override
public String downHandle(String data) {
JSONObject jsonObject = JSONObject.parseObject(data);
long intervalMillis = jsonObject.getLong("intervalMillis");
JSONObject json = new JSONObject();
json.put("intervalMillis",intervalMillis);
return json.toString();
}
}
@AgentDispatcher(msgId = MsgEnum.更新挂载采集间隔应答)
public class TimePointRspHandler implements AgentMsgHandler {
@Override
public UpMsgResponse upHandle(String data, String clientId) {
JSONObject jsonObject = JSONObject.parseObject(data);
jsonObject.put("dataType",MsgEnum.更新挂载采集间隔应答.getValue());
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(jsonObject.toString()).build();
agentProducer.asyncSend(mqMsg);
return null;
}
}
@AgentDispatcher(msgId = MsgEnum.更新内存采集间隔)
public class TimeMemoryHandler implements AgentMsgHandler {
@Override
public String downHandle(String data) {
JSONObject jsonObject = JSONObject.parseObject(data);
long intervalMillis = jsonObject.getLong("intervalMillis");
JSONObject json = new JSONObject();
json.put("intervalMillis",intervalMillis);
return json.toString();
}
}
@AgentDispatcher(msgId = MsgEnum.更新内存采集间隔应答)
public class TimeMemoryRspHandler implements AgentMsgHandler {
@Override
public UpMsgResponse upHandle(String data, String clientId) {
JSONObject jsonObject = JSONObject.parseObject(data);
jsonObject.put("dataType",MsgEnum.更新内存采集间隔应答.getValue());
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(jsonObject.toString()).build();
agentProducer.asyncSend(mqMsg);
return null;
}
}
@AgentDispatcher(msgId = MsgEnum.更新系统采集间隔)
public class TimeSystemHandler implements AgentMsgHandler {
@Override
public String downHandle(String data) {
JSONObject jsonObject = JSONObject.parseObject(data);
long intervalMillis = jsonObject.getLong("intervalMillis");
JSONObject json = new JSONObject();
json.put("intervalMillis",intervalMillis);
return json.toString();
}
}
@AgentDispatcher(msgId = MsgEnum.更新系统采集间隔应答)
public class TimeSystemRspHandler implements AgentMsgHandler {
@Override
public UpMsgResponse upHandle(String data, String clientId) {
JSONObject jsonObject = JSONObject.parseObject(data);
jsonObject.put("dataType",MsgEnum.更新系统采集间隔应答.getValue());
MqMsg mqMsg = MqMsg.builder().topic(agentMqConfig.getAgentTopic()).content(jsonObject.toString()).build();
agentProducer.asyncSend(mqMsg);
return null;
}
}
@@ -19,6 +19,7 @@ import io.netty.util.CharsetUtil;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
import java.util.Objects;
@Component
@ChannelHandler.Sharable
@@ -97,15 +98,18 @@ public class AgentDecoderHandler extends ChannelInboundHandlerAdapter {
String[] arr = content.split("@tong-ran");
for (String message : arr) {
JSONObject jsonObject = JSONObject.parseObject(message);
String clientId = jsonObject.getString("clientId");
String dataType = jsonObject.getString("dataType");
AgentMsgHandler msgHandler = agentDispatcherManager.getHandler(dataType + "&" + AgentDispatcher.VersionEnum.V1.value);
UpMsgResponse response = msgHandler.upHandle(message);
UpMsgResponse response = msgHandler.upHandle(message, clientId );
AssertLog.info("<<[up-after-handle]:clientId:{},type={},[handle-content]={}", response.getClientId(), dataType, message);
Message agentMessage = Message.builder().build();
agentMessage.setClientId(response.getClientId());
agentMessage.setDataType(dataType);
agentMessage.setData(response.getContent());
ctx.fireChannelRead(agentMessage);//传递到下一个handler
if(Objects.nonNull(response)){
Message agentMessage = Message.builder().build();
agentMessage.setClientId(response.getClientId());
agentMessage.setDataType(dataType);
agentMessage.setData(response.getContent());
ctx.fireChannelRead(agentMessage);//传递到下一个handler
}
}
}
if (isClear) {
@@ -37,7 +37,7 @@ public class AgentEncoderHandler extends ChannelOutboundHandlerAdapter {
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);
@@ -50,7 +50,39 @@ public enum MsgEnum {
交换机包("SWITCHBOARD"),
交换机包应答("SWITCHBOARD_RSP");
交换机包应答("SWITCHBOARD_RSP"),
更新CPU采集间隔("TIME_CPU"),
更新CPU采集间隔应答("TIME_CPU_RSP"),
更新容器采集间隔("TIME_DOCKER"),
更新容器采集间隔应答("TIME_DOCKER_RSP"),
更新网卡采集间隔("TIME_NET"),
更新网卡采集间隔应答("TIME_NET_RSP"),
更新交换机采集间隔("TIME_SWITCH"),
更新交换机采集间隔应答("TIME_SWITCH_RSP"),
更新磁盘采集间隔("TIME_DISK"),
更新磁盘采集间隔应答("TIME_DISK_RSP"),
更新挂载采集间隔("TIME_POINT"),
更新挂载采集间隔应答("TIME_POINT_RSP"),
更新内存采集间隔("TIME_MEMORY"),
更新内存采集间隔应答("TIME_MEMORY_RSP"),
更新系统采集间隔("TIME_SYSTEM"),
更新系统采集间隔应答("TIME_SYSTEM_RSP");
private String value;