Fix issue: requestID issue on TCAgen2 CL output
[dcaegen2/analytics/tca-gen2.git] / dcae-analytics / dcae-analytics-web / src / main / java / org / onap / dcae / analytics / web / dmaap / MrMessageSplitter.java
1 /*
2  * ================================================================================
3  * Copyright (c) 2018 AT&T Intellectual Property. All rights reserved.
4  * ================================================================================
5  * Licensed under the Apache License, Version 2.0 (the "License");
6  * you may not use this file except in compliance with the License.
7  * You may obtain a copy of the License at
8  *
9  *      http://www.apache.org/licenses/LICENSE-2.0
10  *
11  * Unless required by applicable law or agreed to in writing, software
12  * distributed under the License is distributed on an "AS IS" BASIS,
13  * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14  * See the License for the specific language governing permissions and
15  * limitations under the License.
16  * ============LICENSE_END=========================================================
17  *
18  */
19
20 package org.onap.dcae.analytics.web.dmaap;
21
22 import static org.apache.commons.text.StringEscapeUtils.unescapeJava;
23 import static org.apache.commons.text.StringEscapeUtils.unescapeJson;
24 import static org.onap.dcae.analytics.model.AnalyticsHttpConstants.REQUEST_ID_HEADER_KEY;
25
26 import java.io.IOException;
27 import java.util.Collections;
28 import java.util.Date;
29 import java.util.LinkedList;
30 import java.util.List;
31 import java.util.stream.IntStream;
32
33 import javax.annotation.Nonnull;
34 import javax.annotation.Nullable;
35
36 import org.onap.dcae.analytics.model.AnalyticsHttpConstants;
37 import org.onap.dcae.analytics.model.DmaapMrConstants;
38 import org.onap.dcae.analytics.tca.core.exception.AnalyticsParsingException;
39 import org.onap.dcae.analytics.tca.core.util.LogSpec;
40 import org.onap.dcae.analytics.web.util.AnalyticsHttpUtils;
41 import org.onap.dcae.utils.eelf.logger.api.log.EELFLogFactory;
42 import org.onap.dcae.utils.eelf.logger.api.log.EELFLogger;
43 import org.onap.dcae.utils.eelf.logger.api.spec.AuditLogSpec;
44 import org.onap.dcae.utils.eelf.logger.api.spec.ErrorLogSpec;
45 import org.springframework.integration.splitter.AbstractMessageSplitter;
46 import org.springframework.integration.support.MessageBuilder;
47 import org.springframework.messaging.Message;
48
49 import com.fasterxml.jackson.databind.JsonNode;
50 import com.fasterxml.jackson.databind.ObjectMapper;
51
52 /**
53  * DMaaP MR message splitter split the incoming messages into batch of given batch size
54  *
55  * @author Rajiv Singla
56  */
57 public class MrMessageSplitter extends AbstractMessageSplitter {
58
59     private static final EELFLogger eelfLogger = EELFLogFactory.getLogger(MrMessageSplitter.class);
60
61     private final ObjectMapper objectMapper;
62     private final Integer batchSize;
63
64     public MrMessageSplitter(@Nonnull final ObjectMapper objectMapper,
65                              @Nonnull final Integer batchSize) {
66         this.objectMapper = objectMapper;
67         this.batchSize = batchSize;
68     }
69
70     @Override
71     protected Object splitMessage(final Message<?> message) {
72
73         final String requestId = AnalyticsHttpUtils.getRequestId(message.getHeaders());
74         final List<String> dmaapMessages = convertJsonToStringMessages(requestId, String.class.cast(message.getPayload()).trim());
75
76         final String transactionId = AnalyticsHttpUtils.getTransactionId(message.getHeaders());
77
78         final Date requestBeginTimestamp = AnalyticsHttpUtils.getTimestampFromHeaders(message.getHeaders(),
79                 AnalyticsHttpConstants.REQUEST_BEGIN_TS_HEADER_KEY);
80         final AuditLogSpec auditLogSpec = LogSpec.createAuditLogSpec(requestId, requestBeginTimestamp);
81
82         eelfLogger.auditLog().info("Request Id: {}, Transaction Id: {}, dmaapMessages: {},"
83                 + " Received new messages from DMaaP MR. Count: {}",
84                 auditLogSpec, requestId, transactionId, dmaapMessages.toString(), String.valueOf(dmaapMessages.size()));
85
86         final List<List<String>> messagePartitions = partition(dmaapMessages, batchSize);
87         eelfLogger.auditLog().info("Request Id: {}, Transaction Id: {},  Max allowed messages per batch: {}. " +
88                 "No of batches created: {}",
89                 auditLogSpec, requestId, transactionId, String.valueOf(batchSize), String.valueOf(messagePartitions.size()));
90
91         // append batch id to request id header
92         return messagePartitions.isEmpty() ? null : IntStream.range(0, messagePartitions.size())
93                 .mapToObj(batchIndex ->
94                         MessageBuilder
95                                 .withPayload(messagePartitions.get(0))
96                                 .copyHeaders(message.getHeaders())
97                                 .setHeader(REQUEST_ID_HEADER_KEY, requestId)
98                                 .build()
99
100                 );
101     }
102
103     /**
104      * Converts DMaaP MR subscriber messages json string to List of messages. If message Json String is empty
105      * or null
106      *
107      * @param messagesJsonString json messages String
108      *
109      * @return List containing DMaaP MR Messages
110      */
111     private List<String> convertJsonToStringMessages(String requestId, @Nullable final String messagesJsonString) {
112
113         final LinkedList<String> messages = new LinkedList<>();
114
115         // If message string is not null or not empty parse json message array to List of string messages
116         if (messagesJsonString != null && !messagesJsonString.trim().isEmpty()
117                 && !DmaapMrConstants.SUBSCRIBER_EMPTY_MESSAGE_RESPONSE_STRING.equals(messagesJsonString.trim())) {
118
119             try {
120                 // get root node
121                 final JsonNode rootNode = objectMapper.readTree(messagesJsonString);
122                 // iterate over root node and parse arrays messages
123                 for (JsonNode jsonNode : rootNode) {
124                     // if array parse it is array of messages
125                     final String incomingMessageString = jsonNode.toString();
126                     if (jsonNode.isArray()) {
127                         final List messageList = objectMapper.readValue(incomingMessageString, List.class);
128                         for (Object message : messageList) {
129                             final String jsonMessageString = objectMapper.writeValueAsString(message);
130                             addUnescapedJsonToMessage(messages, jsonMessageString);
131                         }
132                     } else {
133                         // parse it as object
134                         addUnescapedJsonToMessage(messages, incomingMessageString);
135                     }
136                 }
137
138             } catch (IOException e) {
139                 ErrorLogSpec errorLogSpec = LogSpec.createErrorLogSpec(requestId);
140                 eelfLogger.errorLog().error("Unable to convert subscriber Json String to Messages. " +
141                         "Subscriber Response String: {}, Json Error: {}", errorLogSpec, messagesJsonString, e.toString());
142                 String errorMessage = String.format("Unable to convert subscriber Json String to Messages. " +
143                         "Subscriber Response String: %s, Json Error: %s", messagesJsonString, e);
144                 throw new AnalyticsParsingException(errorMessage, e);
145             }
146
147         }
148         return messages;
149     }
150
151     /**
152      * Adds unescaped Json messages to given messages list
153      *
154      * @param messages message list in which unescaped messages will be added
155      * @param incomingMessageString incoming message string that may need to be escaped
156      */
157     private static void addUnescapedJsonToMessage(List<String> messages, String incomingMessageString) {
158         if (incomingMessageString.startsWith("\"") && incomingMessageString.endsWith("\"")) {
159             messages.add(unescapeJava(unescapeJson(
160                     incomingMessageString.substring(1, incomingMessageString.length() - 1))));
161         } else {
162             messages.add(unescapeJava(unescapeJson(incomingMessageString)));
163         }
164     }
165
166     /**
167      * Partition list into multiple lists
168      *
169      * @param list input list that needs to be broken into chunks
170      * @param batchSize batch size for each list
171      * @param <E> element type of the list
172      *
173      * @return List containing list of entries of specified batch size
174      */
175     private static <E> List<List<E>> partition(List<E> list, final Integer batchSize) {
176
177         if (list == null || batchSize == null || batchSize <= 0 || list.size() < batchSize) {
178             return Collections.singletonList(list);
179         }
180
181         final List<List<E>> result = new LinkedList<>();
182
183         for (int i = 0; i < list.size(); i++) {
184
185             if (i == 0 || i % batchSize == 0) {
186                 List<E> sublist = new LinkedList<>();
187                 result.add(sublist);
188             }
189
190             final List<E> lastSubList = result.get(result.size() - 1);
191             lastSubList.add(list.get(i));
192
193         }
194         return result;
195     }
196 }