2 * ============LICENSE_START=======================================================
4 * ================================================================================
5 * Copyright (C) 2017-2019 AT&T Intellectual Property. All rights reserved.
6 * ================================================================================
7 * Copyright (C) 2017 Amdocs
8 * ================================================================================
9 * Modifications Copyright (C) 2019 Ericsson
10 * =============================================================================
11 * Licensed under the Apache License, Version 2.0 (the "License");
12 * you may not use this file except in compliance with the License.
13 * You may obtain a copy of the License at
15 * http://www.apache.org/licenses/LICENSE-2.0
17 * Unless required by applicable law or agreed to in writing, software
18 * distributed under the License is distributed on an "AS IS" BASIS,
19 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
20 * See the License for the specific language governing permissions and
21 * limitations under the License.
23 * ============LICENSE_END=========================================================
26 package org.onap.appc.listener.impl;
28 import com.att.eelf.configuration.EELFLogger;
29 import com.att.eelf.configuration.EELFManager;
31 import org.onap.appc.adapter.factory.MessageService;
32 import org.onap.appc.adapter.message.Consumer;
33 import org.onap.appc.adapter.message.Producer;
34 import org.onap.appc.adapter.messaging.dmaap.http.HttpDmaapConsumerImpl;
35 import org.onap.appc.adapter.messaging.dmaap.http.HttpDmaapProducerImpl;
36 import org.onap.appc.listener.EventHandler;
37 import org.onap.appc.listener.ListenerProperties;
38 import org.onap.appc.listener.util.Mapper;
39 import org.onap.appc.logging.LoggingConstants;
42 import java.util.ArrayList;
43 import java.util.Collection;
44 import java.util.HashSet;
45 import java.util.List;
48 * This class is a wrapper for the DMaaP client provided in appc-dmaap-adapter. Its aim is to ensure
49 * that only well formed messages are sent and received on DMaaP.
51 public class EventHandlerImpl implements EventHandler {
53 private final EELFLogger LOG = EELFManager.getInstance().getLogger(EventHandlerImpl.class);
56 * The amount of time in seconds to keep a connection to a topic open while waiting for data
58 private int READ_TIMEOUT = 60;
61 * The pool of hosts to query against
63 private Collection<String> pool;
66 * The topic to read messages from
68 private String readTopic;
71 * The topic to write messages to
73 private String writeTopic;
76 * The client (group) name to use for reading messages
78 private String clientName;
81 * The id of the client (group) that is reading messages
83 private String clientId;
86 * The api public key to use for authentication
88 private String apiKey;
91 * The api secret key to use for authentication
93 private String apiSecret;
96 * A json object containing filter arguments.
98 private String filter_json;
102 * Blacklist time for a server with response problem in seconds
104 private String responseProblemBlacklistTime;
107 * Blacklist time for a server with server problem in seconds
109 private String serverProblemBlacklistTime;
112 * Blacklist time for a server with DNS problem in seconds
114 private String dnsIssueBlacklistTime;
117 * Blacklist time for a server with IO Exception problem in seconds
119 private String ioExceptionBlacklistTime;
121 private MessageService messageService;
123 private Consumer reader = null;
124 private Producer producer = null;
126 public EventHandlerImpl(ListenerProperties props) {
127 pool = new HashSet<>();
130 readTopic = props.getProperty(ListenerProperties.KEYS.TOPIC_READ);
131 clientName = props.getProperty(ListenerProperties.KEYS.CLIENT_NAME, "APP-C");
132 clientId = props.getProperty(ListenerProperties.KEYS.CLIENT_ID, "0");
133 apiKey = props.getProperty(ListenerProperties.KEYS.AUTH_USER_KEY);
134 apiSecret = props.getProperty(ListenerProperties.KEYS.AUTH_SECRET_KEY);
135 responseProblemBlacklistTime = props.getProperty(ListenerProperties.KEYS.PROBLEM_WITH_RESPONSE_BLACKLIST_TIME);
136 serverProblemBlacklistTime = props.getProperty(ListenerProperties.KEYS.PROBLEM_SERVERSIDE_ERROR_BLACKLIST_TIME);
137 dnsIssueBlacklistTime = props.getProperty(ListenerProperties.KEYS.PROBLEM_DNS_BLACKLIST_TIME);
138 ioExceptionBlacklistTime = props.getProperty(ListenerProperties.KEYS.PROBLEM_IO_EXCEPTION_BLACKLIST_TIME);
140 filter_json = props.getProperty(ListenerProperties.KEYS.TOPIC_READ_FILTER);
142 READ_TIMEOUT = Integer
143 .valueOf(props.getProperty(ListenerProperties.KEYS.TOPIC_READ_TIMEOUT, String.valueOf(READ_TIMEOUT)));
145 String hostnames = props.getProperty(ListenerProperties.KEYS.HOSTS);
146 if (hostnames != null && !hostnames.isEmpty()) {
147 for (String name : hostnames.split(",")) {
152 String writeTopicStr = props.getProperty(ListenerProperties.KEYS.TOPIC_WRITE);
153 if (writeTopicStr != null) {
154 for (String topic : writeTopicStr.split(",")) {
159 messageService = MessageService.parse(props.getProperty(ListenerProperties.KEYS.MESSAGE_SERVICE));
161 LOG.info(String.format(
162 "Configured to use %s client on host pool [%s]. Reading from [%s] filtered by %s. Writing to [%s]. Authenticated using %s",
163 messageService, hostnames, readTopic, filter_json, writeTopic, apiKey));
168 public List<String> getIncomingEvents() {
169 return getIncomingEvents(1000);
173 public List<String> getIncomingEvents(int limit) {
174 List<String> out = new ArrayList<>();
175 LOG.info(String.format("Getting up to %d incoming events", limit));
176 // reuse the consumer object instead of creating a new one every time
177 if (reader == null) {
178 LOG.info("Getting Consumer...");
179 reader = getConsumer();
181 if (reader != null) {
182 List<String> items = reader.fetch(READ_TIMEOUT * 1000, limit);
183 for (String item : items) {
187 LOG.info(String.format("Read %d messages from %s as %s/%s.", out.size(), readTopic, clientName, clientId));
192 public <T> List<T> getIncomingEvents(Class<T> cls) {
193 return getIncomingEvents(cls, 1000);
197 public <T> List<T> getIncomingEvents(Class<T> cls, int limit) {
198 List<String> incomingStrings = getIncomingEvents(limit);
199 return Mapper.mapList(incomingStrings, cls);
203 public void postStatus(String event) {
204 postStatus(null, event);
208 public void postStatus(String partition, String event) {
209 LOG.debug(String.format("Posting Message [%s]", event));
210 if (producer == null) {
211 LOG.info("Getting Producer...");
212 producer = getProducer();
214 producer.post(partition, event);
218 * Returns a consumer object for direct access to our Cambria consumer interface
220 * @return An instance of the consumer interface
222 protected Consumer getConsumer() {
223 LOG.debug(String.format("Getting Consumer: %s %s/%s/%s", pool, readTopic, clientName, clientId));
224 if (filter_json == null && writeTopic.equals(readTopic)) {
226 "*****We will be writing and reading to the same topic without a filter. This will cause an infinite loop.*****");
230 out = new HttpDmaapConsumerImpl(pool, readTopic, clientName, clientId, filter_json, apiKey, apiSecret);
232 if (out != null && responseProblemBlacklistTime != null && responseProblemBlacklistTime.length() > 0)
234 out.setResponseProblemBlacklistTime(responseProblemBlacklistTime);
237 if (out != null && serverProblemBlacklistTime != null && serverProblemBlacklistTime.length() > 0)
239 out.setServerProblemBlacklistTime(serverProblemBlacklistTime);
242 if (out != null && dnsIssueBlacklistTime != null && dnsIssueBlacklistTime.length() > 0)
244 out.setDnsIssueBlacklistTime(dnsIssueBlacklistTime);
247 if (out != null && ioExceptionBlacklistTime != null && ioExceptionBlacklistTime.length() > 0)
249 out.setIOExceptionBlacklistTime(ioExceptionBlacklistTime);
252 for (String url : pool) {
253 if (url.contains("3905") || url.contains("https")) {
263 * Returns a consumer object for direct access to our Cambria producer interface
265 * @return An instance of the producer interface
267 protected Producer getProducer() {
268 LOG.debug(String.format("Getting Producer: %s %s", pool, readTopic));
271 out = new HttpDmaapProducerImpl(pool, writeTopic, apiKey, apiSecret);
272 if (out != null && responseProblemBlacklistTime != null && responseProblemBlacklistTime.length() > 0)
274 out.setResponseProblemBlacklistTime(responseProblemBlacklistTime);
277 if (out != null && serverProblemBlacklistTime != null && serverProblemBlacklistTime.length() > 0)
279 out.setServerProblemBlacklistTime(serverProblemBlacklistTime);
282 if (out != null && dnsIssueBlacklistTime != null && dnsIssueBlacklistTime.length() > 0)
284 out.setDnsIssueBlacklistTime(dnsIssueBlacklistTime);
287 if (out != null && ioExceptionBlacklistTime != null && ioExceptionBlacklistTime.length() > 0)
289 out.setIOExceptionBlacklistTime(ioExceptionBlacklistTime);
292 for (String url : pool) {
293 if (url.contains("3905") || url.contains("https")) {
303 public void closeClients() {
304 LOG.debug("Closing Consumer and Producer DMaaP clients");
305 if (reader != null) {
308 if (producer != null) {
314 public String getClientId() {
319 public void setClientId(String clientId) {
320 this.clientId = clientId;
324 public String getClientName() {
329 public void setClientName(String clientName) {
330 this.clientName = clientName;
331 MDC.put(LoggingConstants.MDCKeys.PARTNER_NAME, clientName);
335 public void addToPool(String hostname) {
340 public Collection<String> getPool() {
345 public void removeFromPool(String hostname) {
346 pool.remove(hostname);
350 public String getReadTopic() {
355 public void setReadTopic(String readTopic) {
356 this.readTopic = readTopic;
360 public String getWriteTopic() {
365 public void setWriteTopic(String writeTopic) {
366 this.writeTopic = writeTopic;
370 public void setResponseProblemBlacklistTime(String duration){
371 this.responseProblemBlacklistTime = duration;
375 public void setServerProblemBlacklistTime(String duration){
376 this.serverProblemBlacklistTime = duration;
380 public void setDnsIssueBlacklistTime(String duration){
381 this.dnsIssueBlacklistTime = duration;
385 public void setIOExceptionBlacklistTime(String duration){
386 this.ioExceptionBlacklistTime = duration;
390 public void clearCredentials() {
396 public void setCredentials(String key, String secret) {