2 * ============LICENSE_START=======================================================
4 * ================================================================================
5 * Copyright (C) 2017 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 * ECOMP is a trademark and service mark of AT&T Intellectual Property.
22 * ============LICENSE_END=========================================================
25 package org.openecomp.appc.listener.impl;
26 import com.att.eelf.configuration.EELFLogger;
27 import com.att.eelf.configuration.EELFManager;
29 import org.openecomp.appc.adapter.factory.DmaapMessageAdapterFactoryImpl;
30 import org.openecomp.appc.adapter.factory.MessageService;
31 import org.openecomp.appc.adapter.message.Consumer;
32 import org.openecomp.appc.adapter.message.MessageAdapterFactory;
33 import org.openecomp.appc.adapter.message.Producer;
34 import org.openecomp.appc.listener.EventHandler;
35 import org.openecomp.appc.listener.ListenerProperties;
36 import org.openecomp.appc.listener.util.Mapper;
37 import org.openecomp.appc.logging.LoggingConstants;
38 import org.osgi.framework.BundleContext;
39 import org.osgi.framework.FrameworkUtil;
40 import org.osgi.framework.ServiceReference;
46 * This class is a wrapper for the DMaaP client provided in appc-dmaap-adapter. Its aim is to ensure that only well formed
47 * messages are sent and received on DMaaP.
50 public class EventHandlerImpl implements EventHandler {
52 private final EELFLogger LOG = EELFManager.getInstance().getLogger(EventHandlerImpl.class);
55 * The amount of time in seconds to keep a connection to a topic open while waiting for data
57 private int READ_TIMEOUT = 60;
60 * The pool of hosts to query against
62 private Collection<String> pool;
65 * The topic to read messages from
67 private String readTopic;
70 * The topic to write messages to
72 private Set<String> writeTopics;
75 * The client (group) name to use for reading messages
77 private String clientName;
80 * The id of the client (group) that is reading messages
82 private String clientId;
85 * The api public key to use for authentication
87 private String apiKey;
90 * The api secret key to use for authentication
92 private String apiSecret;
95 * A json object containing filter arguments.
97 private String filter_json;
99 private MessageService messageService;
101 private Consumer reader = null;
102 private Producer producer = null;
104 public EventHandlerImpl(ListenerProperties props) {
105 pool = new HashSet<String>();
106 writeTopics = new HashSet<String>();
109 readTopic = props.getProperty(ListenerProperties.KEYS.TOPIC_READ);
110 clientName = props.getProperty(ListenerProperties.KEYS.CLIENT_NAME, "APP-C");
111 clientId = props.getProperty(ListenerProperties.KEYS.CLIENT_ID, "0");
112 apiKey = props.getProperty(ListenerProperties.KEYS.AUTH_USER_KEY);
113 apiSecret = props.getProperty(ListenerProperties.KEYS.AUTH_SECRET_KEY);
115 filter_json = props.getProperty(ListenerProperties.KEYS.TOPIC_READ_FILTER);
117 READ_TIMEOUT = Integer
118 .valueOf(props.getProperty(ListenerProperties.KEYS.TOPIC_READ_TIMEOUT, String.valueOf(READ_TIMEOUT)));
120 String hostnames = props.getProperty(ListenerProperties.KEYS.HOSTS);
121 if (hostnames != null && !hostnames.isEmpty()) {
122 for (String name : hostnames.split(",")) {
127 String writeTopicStr = props.getProperty(ListenerProperties.KEYS.TOPIC_WRITE);
128 if (writeTopicStr != null) {
129 for (String topic : writeTopicStr.split(",")) {
130 writeTopics.add(topic);
134 messageService = MessageService.parse(props.getProperty(ListenerProperties.KEYS.MESSAGE_SERVICE));
136 LOG.info(String.format(
137 "Configured to use %s client on host pool [%s]. Reading from [%s] filtered by %s. Wriring to [%s]. Authenticated using %s",
138 messageService, hostnames, readTopic, filter_json, writeTopics, apiKey));
143 public List<String> getIncomingEvents() {
144 return getIncomingEvents(1000);
148 public List<String> getIncomingEvents(int limit) {
149 List<String> out = new ArrayList<String>();
150 LOG.info(String.format("Getting up to %d incoming events", limit));
151 // reuse the consumer object instead of creating a new one every time
152 if (reader == null) {
153 LOG.info("Getting Consumer...");
154 reader = getConsumer();
157 List<String> items = null;
159 items = reader.fetch(READ_TIMEOUT * 1000, limit);
161 LOG.error("EvenHandlerImpl.getIncomingEvents",r);
165 for (String item : items) {
168 LOG.info(String.format("Read %d messages from %s as %s/%s.", out.size(), readTopic, clientName, clientId));
173 public <T> List<T> getIncomingEvents(Class<T> cls) {
174 return getIncomingEvents(cls, 1000);
178 public <T> List<T> getIncomingEvents(Class<T> cls, int limit) {
179 List<String> incomingStrings = getIncomingEvents(limit);
180 return Mapper.mapList(incomingStrings, cls);
184 public void postStatus(String event) {
185 postStatus(null, event);
189 public void postStatus(String partition, String event) {
190 LOG.debug(String.format("Posting Message [%s]", event));
191 if (producer == null) {
192 LOG.info("Getting Producer...");
193 producer = getProducer();
195 producer.post(partition, event);
199 * Returns a consumer object for direct access to our Cambria consumer interface
201 * @return An instance of the consumer interface
203 protected Consumer getConsumer() {
204 LOG.debug(String.format("Getting Consumer: %s %s/%s/%s", pool, readTopic, clientName, clientId));
205 if (filter_json == null && writeTopics.contains(readTopic)) {
207 "*****We will be writing and reading to the same topic without a filter. This will cause an infinite loop.*****");
211 BundleContext ctx = FrameworkUtil.getBundle(EventHandlerImpl.class).getBundleContext();
213 ServiceReference svcRef = ctx.getServiceReference(MessageAdapterFactory.class.getName());
214 if (svcRef != null) {
216 out = ((MessageAdapterFactory) ctx.getService(svcRef)).createConsumer(pool, readTopic, clientName, clientId, filter_json, apiKey, apiSecret);
218 //TODO:create eelf message
219 LOG.error("EvenHandlerImp.getConsumer calling MessageAdapterFactor.createConsumer",e);
221 for (String url : pool) {
222 if (url.contains("3905") || url.contains("https")) {
233 * Returns a consumer object for direct access to our Cambria producer interface
235 * @return An instance of the producer interface
237 protected Producer getProducer() {
238 LOG.debug(String.format("Getting Producer: %s %s", pool, readTopic));
241 BundleContext ctx = FrameworkUtil.getBundle(EventHandlerImpl.class).getBundleContext();
243 ServiceReference svcRef = ctx.getServiceReference(MessageAdapterFactory.class.getName());
244 if (svcRef != null) {
245 out = ((MessageAdapterFactory) ctx.getService(svcRef)).createProducer(pool, writeTopics,apiKey, apiSecret);
246 for (String url : pool) {
247 if (url.contains("3905") || url.contains("https")) {
258 public void closeClients() {
259 LOG.debug("Closing Consumer and Producer DMaaP clients");
260 if (reader != null) {
263 if (producer != null) {
269 public String getClientId() {
274 public void setClientId(String clientId) {
275 this.clientId = clientId;
279 public String getClientName() {
284 public void setClientName(String clientName) {
285 this.clientName = clientName;
286 MDC.put(LoggingConstants.MDCKeys.PARTNER_NAME, clientName);
290 public void addToPool(String hostname) {
295 public Collection<String> getPool() {
300 public void removeFromPool(String hostname) {
301 pool.remove(hostname);
305 public String getReadTopic() {
310 public void setReadTopic(String readTopic) {
311 this.readTopic = readTopic;
315 public Set<String> getWriteTopics() {
320 public void setWriteTopics(Set<String> writeTopics) {
321 this.writeTopics = writeTopics;
325 public void clearCredentials() {
331 public void setCredentials(String key, String secret) {