This is an automated email from the ASF dual-hosted git repository.

mikexue pushed a commit to branch develop
in repository https://gitbox.apache.org/repos/asf/incubator-eventmesh.git


The following commit(s) were added to refs/heads/develop by this push:
     new 0acafcc  resolve_exception->Consumer subscription topic is invalid 
#590 (#598)
0acafcc is described below

commit 0acafccc080e0d31d7d24e6384ca4527702354a0
Author: hagsyn <[email protected]>
AuthorDate: Fri Nov 19 17:22:34 2021 +0800

    resolve_exception->Consumer subscription topic is invalid #590 (#598)
    
    close #590
---
 .../http/processor/SubscribeProcessor.java         | 115 +++++++++++----------
 1 file changed, 58 insertions(+), 57 deletions(-)

diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SubscribeProcessor.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SubscribeProcessor.java
index d91b462..ab411f9 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SubscribeProcessor.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SubscribeProcessor.java
@@ -70,49 +70,49 @@ public class SubscribeProcessor implements 
HttpRequestProcessor {
 
     @Override
     public void processRequest(ChannelHandlerContext ctx, 
AsyncContext<HttpCommand> asyncContext)
-        throws Exception {
+            throws Exception {
         HttpCommand responseEventMeshCommand;
         httpLogger.info("cmd={}|{}|client2eventMesh|from={}|to={}",
-            
RequestCode.get(Integer.valueOf(asyncContext.getRequest().getRequestCode())),
-            EventMeshConstants.PROTOCOL_HTTP,
-            RemotingHelper.parseChannelRemoteAddr(ctx.channel()), 
IPUtil.getLocalAddress());
+                
RequestCode.get(Integer.valueOf(asyncContext.getRequest().getRequestCode())),
+                EventMeshConstants.PROTOCOL_HTTP,
+                RemotingHelper.parseChannelRemoteAddr(ctx.channel()), 
IPUtil.getLocalAddress());
         SubscribeRequestHeader subscribeRequestHeader =
-            (SubscribeRequestHeader) asyncContext.getRequest().getHeader();
+                (SubscribeRequestHeader) asyncContext.getRequest().getHeader();
         SubscribeRequestBody subscribeRequestBody =
-            (SubscribeRequestBody) asyncContext.getRequest().getBody();
+                (SubscribeRequestBody) asyncContext.getRequest().getBody();
 
         SubscribeResponseHeader subscribeResponseHeader =
-            SubscribeResponseHeader
-                
.buildHeader(Integer.valueOf(asyncContext.getRequest().getRequestCode()),
-                    
eventMeshHTTPServer.getEventMeshHttpConfiguration().eventMeshCluster,
-                    IPUtil.getLocalAddress(),
-                    
eventMeshHTTPServer.getEventMeshHttpConfiguration().eventMeshEnv,
-                    
eventMeshHTTPServer.getEventMeshHttpConfiguration().eventMeshIDC);
+                SubscribeResponseHeader
+                        
.buildHeader(Integer.valueOf(asyncContext.getRequest().getRequestCode()),
+                                
eventMeshHTTPServer.getEventMeshHttpConfiguration().eventMeshCluster,
+                                IPUtil.getLocalAddress(),
+                                
eventMeshHTTPServer.getEventMeshHttpConfiguration().eventMeshEnv,
+                                
eventMeshHTTPServer.getEventMeshHttpConfiguration().eventMeshIDC);
 
         //validate header
         if (StringUtils.isBlank(subscribeRequestHeader.getIdc())
-            || StringUtils.isBlank(subscribeRequestHeader.getPid())
-            || !StringUtils.isNumeric(subscribeRequestHeader.getPid())
-            || StringUtils.isBlank(subscribeRequestHeader.getSys())) {
+                || StringUtils.isBlank(subscribeRequestHeader.getPid())
+                || !StringUtils.isNumeric(subscribeRequestHeader.getPid())
+                || StringUtils.isBlank(subscribeRequestHeader.getSys())) {
             responseEventMeshCommand = 
asyncContext.getRequest().createHttpCommandResponse(
-                subscribeResponseHeader,
-                SubscribeResponseBody
-                    
.buildBody(EventMeshRetCode.EVENTMESH_PROTOCOL_HEADER_ERR.getRetCode(),
-                        
EventMeshRetCode.EVENTMESH_PROTOCOL_HEADER_ERR.getErrMsg()));
+                    subscribeResponseHeader,
+                    SubscribeResponseBody
+                            
.buildBody(EventMeshRetCode.EVENTMESH_PROTOCOL_HEADER_ERR.getRetCode(),
+                                    
EventMeshRetCode.EVENTMESH_PROTOCOL_HEADER_ERR.getErrMsg()));
             asyncContext.onComplete(responseEventMeshCommand);
             return;
         }
 
         //validate body
         if (StringUtils.isBlank(subscribeRequestBody.getUrl())
-            || CollectionUtils.isEmpty(subscribeRequestBody.getTopics())
-            || StringUtils.isBlank(subscribeRequestBody.getConsumerGroup())) {
+                || CollectionUtils.isEmpty(subscribeRequestBody.getTopics())
+                || 
StringUtils.isBlank(subscribeRequestBody.getConsumerGroup())) {
 
             responseEventMeshCommand = 
asyncContext.getRequest().createHttpCommandResponse(
-                subscribeResponseHeader,
-                SubscribeResponseBody
-                    
.buildBody(EventMeshRetCode.EVENTMESH_PROTOCOL_BODY_ERR.getRetCode(),
-                        
EventMeshRetCode.EVENTMESH_PROTOCOL_BODY_ERR.getErrMsg()));
+                    subscribeResponseHeader,
+                    SubscribeResponseBody
+                            
.buildBody(EventMeshRetCode.EVENTMESH_PROTOCOL_BODY_ERR.getRetCode(),
+                                    
EventMeshRetCode.EVENTMESH_PROTOCOL_BODY_ERR.getErrMsg()));
             asyncContext.onComplete(responseEventMeshCommand);
             return;
         }
@@ -128,17 +128,17 @@ public class SubscribeProcessor implements 
HttpRequestProcessor {
             for (SubscriptionItem item : subTopicList) {
                 try {
                     Acl.doAclCheckInHttpReceive(remoteAddr, user, pass, 
subsystem, item.getTopic(),
-                        requestCode);
+                            requestCode);
                 } catch (Exception e) {
 
                     responseEventMeshCommand = 
asyncContext.getRequest().createHttpCommandResponse(
-                        subscribeResponseHeader,
-                        SendMessageResponseBody
-                            
.buildBody(EventMeshRetCode.EVENTMESH_ACL_ERR.getRetCode(),
-                                e.getMessage()));
+                            subscribeResponseHeader,
+                            SendMessageResponseBody
+                                    
.buildBody(EventMeshRetCode.EVENTMESH_ACL_ERR.getRetCode(),
+                                            e.getMessage()));
                     asyncContext.onComplete(responseEventMeshCommand);
                     aclLogger
-                        .warn("CLIENT HAS NO PERMISSION,SubscribeProcessor 
subscribe failed", e);
+                            .warn("CLIENT HAS NO PERMISSION,SubscribeProcessor 
subscribe failed", e);
                     return;
                 }
             }
@@ -153,7 +153,7 @@ public class SubscribeProcessor implements 
HttpRequestProcessor {
 
             for (SubscriptionItem subTopic : subTopicList) {
                 List<Client> groupTopicClients = 
eventMeshHTTPServer.localClientInfoMapping
-                    .get(consumerGroup + "@" + subTopic.getTopic());
+                        .get(consumerGroup + "@" + subTopic.getTopic());
 
                 if (CollectionUtils.isEmpty(groupTopicClients)) {
                     httpLogger.error("group {} topic {} clients is empty", 
consumerGroup, subTopic);
@@ -170,7 +170,7 @@ public class SubscribeProcessor implements 
HttpRequestProcessor {
                     }
                 }
                 ConsumerGroupConf consumerGroupConf =
-                    
eventMeshHTTPServer.localConsumerGroupMapping.get(consumerGroup);
+                        
eventMeshHTTPServer.localConsumerGroupMapping.get(consumerGroup);
                 if (consumerGroupConf == null) {
                     // new subscription
                     consumerGroupConf = new ConsumerGroupConf(consumerGroup);
@@ -188,7 +188,17 @@ public class SubscribeProcessor implements 
HttpRequestProcessor {
                 } else {
                     // already subscribed
                     Map<String, ConsumerGroupTopicConf> map =
-                        consumerGroupConf.getConsumerGroupTopicConf();
+                            consumerGroupConf.getConsumerGroupTopicConf();
+                    if (!map.containsKey(subTopic.getTopic())) {
+                        //If there are multiple topics, append it
+                        ConsumerGroupTopicConf newTopicConf = new 
ConsumerGroupTopicConf();
+                        newTopicConf.setConsumerGroup(consumerGroup);
+                        newTopicConf.setTopic(subTopic.getTopic());
+                        newTopicConf.setSubscriptionItem(subTopic);
+                        newTopicConf.setUrls(new 
HashSet<>(Arrays.asList(url)));
+                        newTopicConf.setIdcUrls(idcUrls);
+                        map.put(subTopic.getTopic(), newTopicConf);
+                    }
                     for (String key : map.keySet()) {
                         if (StringUtils.equals(subTopic.getTopic(), key)) {
                             ConsumerGroupTopicConf latestTopicConf = new 
ConsumerGroupTopicConf();
@@ -202,15 +212,6 @@ public class SubscribeProcessor implements 
HttpRequestProcessor {
                             latestTopicConf.setIdcUrls(idcUrls);
 
                             map.put(key, latestTopicConf);
-                        } else {
-                            //If there are multiple topics, append it
-                            ConsumerGroupTopicConf newTopicConf = new 
ConsumerGroupTopicConf();
-                            newTopicConf.setConsumerGroup(consumerGroup);
-                            newTopicConf.setTopic(subTopic.getTopic());
-                            newTopicConf.setSubscriptionItem(subTopic);
-                            newTopicConf.setUrls(new 
HashSet<>(Arrays.asList(url)));
-                            newTopicConf.setIdcUrls(idcUrls);
-                            map.put(subTopic.getTopic(), newTopicConf);
                         }
                     }
                 }
@@ -221,7 +222,7 @@ public class SubscribeProcessor implements 
HttpRequestProcessor {
             try {
                 // subscription relationship change notification
                 
eventMeshHTTPServer.getConsumerManager().notifyConsumerManager(consumerGroup,
-                    
eventMeshHTTPServer.localConsumerGroupMapping.get(consumerGroup));
+                        
eventMeshHTTPServer.localConsumerGroupMapping.get(consumerGroup));
 
                 final CompleteHandler<HttpCommand> handler = new 
CompleteHandler<HttpCommand>() {
                     @Override
@@ -232,8 +233,8 @@ public class SubscribeProcessor implements 
HttpRequestProcessor {
                             }
                             eventMeshHTTPServer.sendResponse(ctx, 
httpCommand.httpResponse());
                             
eventMeshHTTPServer.metrics.summaryMetrics.recordHTTPReqResTimeCost(
-                                System.currentTimeMillis()
-                                    - asyncContext.getRequest().getReqTime());
+                                    System.currentTimeMillis()
+                                            - 
asyncContext.getRequest().getReqTime());
                         } catch (Exception ex) {
                             // ignore
                         }
@@ -241,22 +242,22 @@ public class SubscribeProcessor implements 
HttpRequestProcessor {
                 };
 
                 responseEventMeshCommand = 
asyncContext.getRequest().createHttpCommandResponse(
-                    EventMeshRetCode.SUCCESS.getRetCode(), 
EventMeshRetCode.SUCCESS.getErrMsg());
+                        EventMeshRetCode.SUCCESS.getRetCode(), 
EventMeshRetCode.SUCCESS.getErrMsg());
                 asyncContext.onComplete(responseEventMeshCommand, handler);
             } catch (Exception e) {
                 HttpCommand err = 
asyncContext.getRequest().createHttpCommandResponse(
-                    subscribeResponseHeader,
-                    SubscribeResponseBody
-                        
.buildBody(EventMeshRetCode.EVENTMESH_SUBSCRIBE_ERR.getRetCode(),
-                            
EventMeshRetCode.EVENTMESH_SUBSCRIBE_ERR.getErrMsg()
-                                + EventMeshUtil.stackTrace(e, 2)));
+                        subscribeResponseHeader,
+                        SubscribeResponseBody
+                                
.buildBody(EventMeshRetCode.EVENTMESH_SUBSCRIBE_ERR.getRetCode(),
+                                        
EventMeshRetCode.EVENTMESH_SUBSCRIBE_ERR.getErrMsg()
+                                                + EventMeshUtil.stackTrace(e, 
2)));
                 asyncContext.onComplete(err);
                 long endTime = System.currentTimeMillis();
                 httpLogger.error(
-                    "message|eventMesh2mq|REQ|ASYNC|send2MQCost={}ms|topic={}"
-                        + "|bizSeqNo={}|uniqueId={}", endTime - startTime,
-                    JsonUtils.serialize(subscribeRequestBody.getTopics()),
-                    subscribeRequestBody.getUrl(), e);
+                        
"message|eventMesh2mq|REQ|ASYNC|send2MQCost={}ms|topic={}"
+                                + "|bizSeqNo={}|uniqueId={}", endTime - 
startTime,
+                        JsonUtils.serialize(subscribeRequestBody.getTopics()),
+                        subscribeRequestBody.getUrl(), e);
                 
eventMeshHTTPServer.metrics.summaryMetrics.recordSendMsgFailed();
                 
eventMeshHTTPServer.metrics.summaryMetrics.recordSendMsgCost(endTime - 
startTime);
             }
@@ -286,7 +287,7 @@ public class SubscribeProcessor implements 
HttpRequestProcessor {
 
             if 
(eventMeshHTTPServer.localClientInfoMapping.containsKey(groupTopicKey)) {
                 List<Client> localClients =
-                    
eventMeshHTTPServer.localClientInfoMapping.get(groupTopicKey);
+                        
eventMeshHTTPServer.localClientInfoMapping.get(groupTopicKey);
                 boolean isContains = false;
                 for (Client localClient : localClients) {
                     if (StringUtils.equals(localClient.url, client.url)) {

---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to