mirror of
https://gitee.com/san-bing/JChargePointProtocol
synced 2026-05-06 02:49:57 +08:00
@@ -15,6 +15,8 @@ import sanbing.jcpp.infrastructure.cache.CacheValueWrapper;
|
||||
import sanbing.jcpp.infrastructure.cache.TransactionalCache;
|
||||
import sanbing.jcpp.infrastructure.queue.discovery.ServiceInfoProvider;
|
||||
import sanbing.jcpp.proto.gen.ProtocolProto.DownlinkRequestMessage;
|
||||
import sanbing.jcpp.proto.gen.ProtocolProto.LoginRequest;
|
||||
import sanbing.jcpp.proto.gen.ProtocolProto.UplinkQueueMessage;
|
||||
import sanbing.jcpp.protocol.adapter.DownlinkController;
|
||||
|
||||
import java.util.UUID;
|
||||
@@ -55,22 +57,50 @@ public abstract class DownlinkCallService {
|
||||
if (downlinkMessageBuilder.getSessionIdLSB() == 0) {
|
||||
downlinkMessageBuilder.setSessionIdLSB(protocolSessionId.getLeastSignificantBits());
|
||||
}
|
||||
if(downlinkMessageBuilder.getProtocolName() == null){
|
||||
if (downlinkMessageBuilder.getProtocolName() == null) {
|
||||
downlinkMessageBuilder.setProtocolName(pileSession.getProtocolName());
|
||||
}
|
||||
|
||||
String nodeId = pileSession.getNodeId();
|
||||
String nodeIp = pileSession.getNodeIp();
|
||||
int nodeRestPort = pileSession.getNodeRestPort();
|
||||
int nodeGrpcPort = pileSession.getNodeGrpcPort();
|
||||
|
||||
sendDownlinkMessage(downlinkMessageBuilder, nodeId, nodeIp, nodeRestPort, nodeGrpcPort);
|
||||
}
|
||||
|
||||
public void sendDownlinkMessage(DownlinkRequestMessage.Builder downlinkMessageBuilder, UplinkQueueMessage uplinkQueueMessage, LoginRequest loginRequest) {
|
||||
|
||||
if (downlinkMessageBuilder.getSessionIdMSB() == 0) {
|
||||
downlinkMessageBuilder.setSessionIdMSB(uplinkQueueMessage.getSessionIdMSB());
|
||||
}
|
||||
if (downlinkMessageBuilder.getSessionIdLSB() == 0) {
|
||||
downlinkMessageBuilder.setSessionIdLSB(uplinkQueueMessage.getSessionIdLSB());
|
||||
}
|
||||
if (downlinkMessageBuilder.getProtocolName() == null) {
|
||||
downlinkMessageBuilder.setProtocolName(uplinkQueueMessage.getProtocolName());
|
||||
}
|
||||
|
||||
String nodeId = loginRequest.getNodeId();
|
||||
String nodeIp = loginRequest.getNodeHostAddress();
|
||||
int nodeRestPort = loginRequest.getNodeRestPort();
|
||||
int nodeGrpcPort = loginRequest.getNodeGrpcPort();
|
||||
|
||||
sendDownlinkMessage(downlinkMessageBuilder, nodeId, nodeIp, nodeRestPort, nodeGrpcPort);
|
||||
}
|
||||
|
||||
private void sendDownlinkMessage(DownlinkRequestMessage.Builder downlinkMessageBuilder, String nodeId, String nodeIp, int nodeRestPort, int nodeGrpcPort) {
|
||||
if (serviceInfoProvider.isMonolith() &&
|
||||
("caffeine".equalsIgnoreCase(cacheType) || serviceInfoProvider.getServiceId().equalsIgnoreCase(pileSession.getNodeId()))) {
|
||||
("caffeine".equalsIgnoreCase(cacheType) || serviceInfoProvider.getServiceId().equalsIgnoreCase(nodeId))) {
|
||||
|
||||
downlinkController.onDownlink(downlinkMessageBuilder.build())
|
||||
.setResultHandler(result -> log.debug("下行消息发送完成"));
|
||||
|
||||
} else {
|
||||
|
||||
_sendDownlinkMessage(downlinkMessageBuilder.build(), pileSession);
|
||||
_sendDownlinkMessage(downlinkMessageBuilder.build(), nodeIp, nodeRestPort, nodeGrpcPort);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
protected abstract void _sendDownlinkMessage(DownlinkRequestMessage downlinkMessage, PileSession pileSession);
|
||||
protected abstract void _sendDownlinkMessage(DownlinkRequestMessage downlinkMessage, String nodeIp, int nodeRestPort, int nodeGrpcPort);
|
||||
}
|
||||
Reference in New Issue
Block a user