2 * ============LICENSE_START=======================================================
4 * ================================================================================
5 * Copyright (C) 2017-2019 AT&T Intellectual Property. All rights reserved.
6 * ================================================================================
7 * Copyright (C) 2017 Amdocs
8 * =============================================================================
9 * Licensed under the Apache License, Version 2.0 (the "License");
10 * you may not use this file except in compliance with the License.
11 * You may obtain a copy of the License at
13 * http://www.apache.org/licenses/LICENSE-2.0
15 * Unless required by applicable law or agreed to in writing, software
16 * distributed under the License is distributed on an "AS IS" BASIS,
17 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
18 * See the License for the specific language governing permissions and
19 * limitations under the License.
21 * ============LICENSE_END=========================================================
24 package org.onap.appc.listener.impl;
26 import com.att.eelf.configuration.EELFLogger;
27 import com.att.eelf.configuration.EELFManager;
29 import org.onap.appc.adapter.factory.DmaapMessageAdapterFactoryImpl;
30 import org.onap.appc.adapter.factory.MessageService;
31 import org.onap.appc.adapter.message.Consumer;
32 import org.onap.appc.adapter.message.MessageAdapterFactory;
33 import org.onap.appc.adapter.message.Producer;
34 import org.onap.appc.listener.EventHandler;
35 import org.onap.appc.listener.ListenerProperties;
36 import org.onap.appc.listener.util.Mapper;
37 import org.onap.appc.logging.LoggingConstants;
38 import org.osgi.framework.Bundle;
39 import org.osgi.framework.BundleContext;
40 import org.osgi.framework.FrameworkUtil;
41 import org.osgi.framework.ServiceReference;
44 import java.util.ArrayList;
45 import java.util.Collection;
46 import java.util.HashSet;
47 import java.util.List;
51 * This class is a wrapper for the DMaaP client provided in appc-dmaap-adapter. Its aim is to ensure
52 * that only well formed messages are sent and received on DMaaP.
54 public class EventHandlerImpl implements EventHandler {
56 private final EELFLogger LOG = EELFManager.getInstance().getLogger(EventHandlerImpl.class);
59 * The amount of time in seconds to keep a connection to a topic open while waiting for data
61 private int READ_TIMEOUT = 60;
64 * The pool of hosts to query against
66 private Collection<String> pool;
69 * The topic to read messages from
71 private String readTopic;
74 * The topic to write messages to
76 private Set<String> writeTopics;
79 * The client (group) name to use for reading messages
81 private String clientName;
84 * The id of the client (group) that is reading messages
86 private String clientId;
89 * The api public key to use for authentication
91 private String apiKey;
94 * The api secret key to use for authentication
96 private String apiSecret;
99 * A json object containing filter arguments.
101 private String filter_json;
105 * Blacklist time for a server with response problem in seconds
107 private String responseProblemBlacklistTime;
110 * Blacklist time for a server with server problem in seconds
112 private String serverProblemBlacklistTime;
115 * Blacklist time for a server with DNS problem in seconds
117 private String dnsIssueBlacklistTime;
120 * Blacklist time for a server with IO Exception problem in seconds
122 private String ioExceptionBlacklistTime;
124 private MessageService messageService;
126 private Consumer reader = null;
127 private Producer producer = null;
129 public EventHandlerImpl(ListenerProperties props) {
130 pool = new HashSet<>();
131 writeTopics = new HashSet<>();
134 readTopic = props.getProperty(ListenerProperties.KEYS.TOPIC_READ);
135 clientName = props.getProperty(ListenerProperties.KEYS.CLIENT_NAME, "APP-C");
136 clientId = props.getProperty(ListenerProperties.KEYS.CLIENT_ID, "0");
137 apiKey = props.getProperty(ListenerProperties.KEYS.AUTH_USER_KEY);
138 apiSecret = props.getProperty(ListenerProperties.KEYS.AUTH_SECRET_KEY);
139 responseProblemBlacklistTime = props.getProperty(ListenerProperties.KEYS.PROBLEM_WITH_RESPONSE_BLACKLIST_TIME);
140 serverProblemBlacklistTime = props.getProperty(ListenerProperties.KEYS.PROBLEM_SERVERSIDE_ERROR_BLACKLIST_TIME);
141 dnsIssueBlacklistTime = props.getProperty(ListenerProperties.KEYS.PROBLEM_DNS_BLACKLIST_TIME);
142 ioExceptionBlacklistTime = props.getProperty(ListenerProperties.KEYS.PROBLEM_IO_EXCEPTION_BLACKLIST_TIME);
144 filter_json = props.getProperty(ListenerProperties.KEYS.TOPIC_READ_FILTER);
146 READ_TIMEOUT = Integer
147 .valueOf(props.getProperty(ListenerProperties.KEYS.TOPIC_READ_TIMEOUT, String.valueOf(READ_TIMEOUT)));
149 String hostnames = props.getProperty(ListenerProperties.KEYS.HOSTS);
150 if (hostnames != null && !hostnames.isEmpty()) {
151 for (String name : hostnames.split(",")) {
156 String writeTopicStr = props.getProperty(ListenerProperties.KEYS.TOPIC_WRITE);
157 if (writeTopicStr != null) {
158 for (String topic : writeTopicStr.split(",")) {
159 writeTopics.add(topic);
163 messageService = MessageService.parse(props.getProperty(ListenerProperties.KEYS.MESSAGE_SERVICE));
165 LOG.info(String.format(
166 "Configured to use %s client on host pool [%s]. Reading from [%s] filtered by %s. Writing to [%s]. Authenticated using %s",
167 messageService, hostnames, readTopic, filter_json, writeTopics, apiKey));
172 public List<String> getIncomingEvents() {
173 return getIncomingEvents(1000);
177 public List<String> getIncomingEvents(int limit) {
178 List<String> out = new ArrayList<>();
179 LOG.info(String.format("Getting up to %d incoming events", limit));
180 // reuse the consumer object instead of creating a new one every time
181 if (reader == null) {
182 LOG.info("Getting Consumer...");
183 reader = getConsumer();
185 if (reader != null) {
186 List<String> items = reader.fetch(READ_TIMEOUT * 1000, limit);
187 for (String item : items) {
191 LOG.info(String.format("Read %d messages from %s as %s/%s.", out.size(), readTopic, clientName, clientId));
196 public <T> List<T> getIncomingEvents(Class<T> cls) {
197 return getIncomingEvents(cls, 1000);
201 public <T> List<T> getIncomingEvents(Class<T> cls, int limit) {
202 List<String> incomingStrings = getIncomingEvents(limit);
203 return Mapper.mapList(incomingStrings, cls);
207 public void postStatus(String event) {
208 postStatus(null, event);
212 public void postStatus(String partition, String event) {
213 LOG.debug(String.format("Posting Message [%s]", event));
214 if (producer == null) {
215 LOG.info("Getting Producer...");
216 producer = getProducer();
218 producer.post(partition, event);
222 * Returns a consumer object for direct access to our Cambria consumer interface
224 * @return An instance of the consumer interface
226 protected Consumer getConsumer() {
227 LOG.debug(String.format("Getting Consumer: %s %s/%s/%s", pool, readTopic, clientName, clientId));
228 if (filter_json == null && writeTopics.contains(readTopic)) {
230 "*****We will be writing and reading to the same topic without a filter. This will cause an infinite loop.*****");
234 BundleContext ctx = null;
235 Bundle bundle = FrameworkUtil.getBundle(EventHandlerImpl.class);
237 ctx = bundle.getBundleContext();
241 ServiceReference svcRef = ctx.getServiceReference(MessageAdapterFactory.class.getName());
242 if (svcRef != null) {
244 out = ((MessageAdapterFactory) ctx.getService(svcRef))
245 .createConsumer(pool, readTopic, clientName, clientId, filter_json, apiKey, apiSecret);
247 if (out != null && responseProblemBlacklistTime != null && responseProblemBlacklistTime.length() > 0)
249 out.setResponseProblemBlacklistTime(responseProblemBlacklistTime);
252 if (out != null && serverProblemBlacklistTime != null && serverProblemBlacklistTime.length() > 0)
254 out.setServerProblemBlacklistTime(serverProblemBlacklistTime);
257 if (out != null && dnsIssueBlacklistTime != null && dnsIssueBlacklistTime.length() > 0)
259 out.setDnsIssueBlacklistTime(dnsIssueBlacklistTime);
262 if (out != null && ioExceptionBlacklistTime != null && ioExceptionBlacklistTime.length() > 0)
264 out.setIOExceptionBlacklistTime(ioExceptionBlacklistTime);
266 } catch (Exception e) {
267 //TODO:create eelf message
268 LOG.error("EvenHandlerImp.getConsumer calling MessageAdapterFactor.createConsumer", e);
271 for (String url : pool) {
272 if (url.contains("3905") || url.contains("https")) {
284 * Returns a consumer object for direct access to our Cambria producer interface
286 * @return An instance of the producer interface
288 protected Producer getProducer() {
289 LOG.debug(String.format("Getting Producer: %s %s", pool, readTopic));
292 BundleContext ctx = FrameworkUtil.getBundle(EventHandlerImpl.class).getBundleContext();
294 ServiceReference svcRef = ctx.getServiceReference(MessageAdapterFactory.class.getName());
295 if (svcRef != null) {
296 out = ((MessageAdapterFactory) ctx.getService(svcRef))
297 .createProducer(pool, writeTopics, apiKey, apiSecret);
298 if (out != null && responseProblemBlacklistTime != null && responseProblemBlacklistTime.length() > 0)
300 out.setResponseProblemBlacklistTime(responseProblemBlacklistTime);
303 if (out != null && serverProblemBlacklistTime != null && serverProblemBlacklistTime.length() > 0)
305 out.setServerProblemBlacklistTime(serverProblemBlacklistTime);
308 if (out != null && dnsIssueBlacklistTime != null && dnsIssueBlacklistTime.length() > 0)
310 out.setDnsIssueBlacklistTime(dnsIssueBlacklistTime);
313 if (out != null && ioExceptionBlacklistTime != null && ioExceptionBlacklistTime.length() > 0)
315 out.setIOExceptionBlacklistTime(ioExceptionBlacklistTime);
318 for (String url : pool) {
319 if (url.contains("3905") || url.contains("https")) {
331 public void closeClients() {
332 LOG.debug("Closing Consumer and Producer DMaaP clients");
333 if (reader != null) {
336 if (producer != null) {
342 public String getClientId() {
347 public void setClientId(String clientId) {
348 this.clientId = clientId;
352 public String getClientName() {
357 public void setClientName(String clientName) {
358 this.clientName = clientName;
359 MDC.put(LoggingConstants.MDCKeys.PARTNER_NAME, clientName);
363 public void addToPool(String hostname) {
368 public Collection<String> getPool() {
373 public void removeFromPool(String hostname) {
374 pool.remove(hostname);
378 public String getReadTopic() {
383 public void setReadTopic(String readTopic) {
384 this.readTopic = readTopic;
388 public Set<String> getWriteTopics() {
393 public void setWriteTopics(Set<String> writeTopics) {
394 this.writeTopics = writeTopics;
398 public void setResponseProblemBlacklistTime(String duration){
399 this.responseProblemBlacklistTime = duration;
403 public void setServerProblemBlacklistTime(String duration){
404 this.serverProblemBlacklistTime = duration;
408 public void setDnsIssueBlacklistTime(String duration){
409 this.dnsIssueBlacklistTime = duration;
413 public void setIOExceptionBlacklistTime(String duration){
414 this.ioExceptionBlacklistTime = duration;
418 public void clearCredentials() {
424 public void setCredentials(String key, String secret) {