Add Message tracing logger service.
[ccsdk/cds.git] / ms / blueprintsprocessor / modules / commons / message-lib / src / main / kotlin / org / onap / ccsdk / cds / blueprintsprocessor / message / service / KafkaBasicAuthMessageProducerService.kt
1 /*
2  *  Copyright © 2019 IBM.
3  *  Modifications Copyright © 2018-2019 AT&T Intellectual Property.
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  */
17
18 package org.onap.ccsdk.cds.blueprintsprocessor.message.service
19
20 import org.apache.commons.lang.builder.ToStringBuilder
21 import org.apache.kafka.clients.producer.Callback
22 import org.apache.kafka.clients.producer.KafkaProducer
23 import org.apache.kafka.clients.producer.ProducerConfig.*
24 import org.apache.kafka.clients.producer.ProducerRecord
25 import org.apache.kafka.common.header.internals.RecordHeader
26 import org.apache.kafka.common.serialization.ByteArraySerializer
27 import org.apache.kafka.common.serialization.StringSerializer
28 import org.onap.ccsdk.cds.blueprintsprocessor.message.KafkaBasicAuthMessageProducerProperties
29 import org.onap.ccsdk.cds.controllerblueprints.core.asJsonString
30 import org.onap.ccsdk.cds.controllerblueprints.core.defaultToUUID
31 import org.slf4j.LoggerFactory
32 import java.nio.charset.Charset
33
34 class KafkaBasicAuthMessageProducerService(
35         private val messageProducerProperties: KafkaBasicAuthMessageProducerProperties)
36     : BlueprintMessageProducerService {
37
38     private val log = LoggerFactory.getLogger(KafkaBasicAuthMessageProducerService::class.java)!!
39
40     private var kafkaProducer: KafkaProducer<String, ByteArray>? = null
41
42     private val messageLoggerService = MessageLoggerService()
43
44     override suspend fun sendMessageNB(message: Any): Boolean {
45         checkNotNull(messageProducerProperties.topic) { "default topic is not configured" }
46         return sendMessageNB(messageProducerProperties.topic!!, message)
47     }
48
49     override suspend fun sendMessageNB(message: Any, headers: MutableMap<String, String>?): Boolean {
50         checkNotNull(messageProducerProperties.topic) { "default topic is not configured" }
51         return sendMessageNB(messageProducerProperties.topic!!, message, headers)
52     }
53
54     override suspend fun sendMessageNB(topic: String, message: Any,
55                                        headers: MutableMap<String, String>?): Boolean {
56         val byteArrayMessage = when (message) {
57             is String -> message.toByteArray(Charset.defaultCharset())
58             else -> message.asJsonString().toByteArray(Charset.defaultCharset())
59         }
60
61         val record = ProducerRecord<String, ByteArray>(topic, defaultToUUID(), byteArrayMessage)
62         val recordHeaders = record.headers()
63         messageLoggerService.messageProducing(recordHeaders)
64         headers?.let {
65             headers.forEach { (key, value) -> recordHeaders.add(RecordHeader(key, value.toByteArray())) }
66         }
67         val callback = Callback { metadata, exception ->
68             log.info("message published offset(${metadata.offset()}, headers :$headers )")
69         }
70         messageTemplate().send(record, callback).get()
71         return true
72     }
73
74     fun messageTemplate(additionalConfig: Map<String, ByteArray>? = null): KafkaProducer<String, ByteArray> {
75         log.trace("Client Properties : ${ToStringBuilder.reflectionToString(messageProducerProperties)}")
76         val configProps = hashMapOf<String, Any>()
77         configProps[BOOTSTRAP_SERVERS_CONFIG] = messageProducerProperties.bootstrapServers
78         configProps[KEY_SERIALIZER_CLASS_CONFIG] = StringSerializer::class.java
79         configProps[VALUE_SERIALIZER_CLASS_CONFIG] = ByteArraySerializer::class.java
80         if (messageProducerProperties.clientId != null) {
81             configProps[CLIENT_ID_CONFIG] = messageProducerProperties.clientId!!
82         }
83         // TODO("Security Implementation based on type")
84
85         // Add additional Properties
86         if (additionalConfig != null) {
87             configProps.putAll(additionalConfig)
88         }
89
90         if (kafkaProducer == null) {
91             kafkaProducer = KafkaProducer(configProps)
92         }
93         return kafkaProducer!!
94     }
95 }
96