Add Message tracing logger service.
[ccsdk/cds.git] / ms / blueprintsprocessor / modules / commons / message-lib / src / main / kotlin / org / onap / ccsdk / cds / blueprintsprocessor / message / service / KafkaBasicAuthMessageConsumerService.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 kotlinx.coroutines.channels.Channel
21 import kotlinx.coroutines.delay
22 import kotlinx.coroutines.launch
23 import kotlinx.coroutines.runBlocking
24 import org.apache.kafka.clients.CommonClientConfigs
25 import org.apache.kafka.clients.consumer.Consumer
26 import org.apache.kafka.clients.consumer.ConsumerConfig
27 import org.apache.kafka.clients.consumer.KafkaConsumer
28 import org.apache.kafka.common.serialization.ByteArrayDeserializer
29 import org.apache.kafka.common.serialization.StringDeserializer
30 import org.onap.ccsdk.cds.blueprintsprocessor.message.KafkaBasicAuthMessageConsumerProperties
31 import org.onap.ccsdk.cds.controllerblueprints.core.logger
32 import java.nio.charset.Charset
33 import java.time.Duration
34 import kotlin.concurrent.thread
35
36 open class KafkaBasicAuthMessageConsumerService(
37         private val messageConsumerProperties: KafkaBasicAuthMessageConsumerProperties)
38     : BlueprintMessageConsumerService {
39
40     val log = logger(KafkaBasicAuthMessageConsumerService::class)
41
42     val channel = Channel<String>()
43     var kafkaConsumer: Consumer<String, ByteArray>? = null
44
45     @Volatile
46     var keepGoing = true
47
48     fun kafkaConsumer(additionalConfig: Map<String, Any>? = null): Consumer<String, ByteArray> {
49         val configProperties = hashMapOf<String, Any>()
50         configProperties[CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG] = messageConsumerProperties.bootstrapServers
51         configProperties[ConsumerConfig.GROUP_ID_CONFIG] = messageConsumerProperties.groupId
52         configProperties[ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG] = messageConsumerProperties.autoCommit
53         /**
54          * earliest: automatically reset the offset to the earliest offset
55          * latest: automatically reset the offset to the latest offset
56          */
57         configProperties[ConsumerConfig.AUTO_OFFSET_RESET_CONFIG] = messageConsumerProperties.autoOffsetReset
58         configProperties[ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG] = StringDeserializer::class.java
59         configProperties[ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG] = ByteArrayDeserializer::class.java
60         if (messageConsumerProperties.clientId != null) {
61             configProperties[ConsumerConfig.CLIENT_ID_CONFIG] = messageConsumerProperties.clientId!!
62         }
63         /** To handle Back pressure, Get only configured record for processing */
64         if (messageConsumerProperties.pollRecords > 0) {
65             configProperties[ConsumerConfig.MAX_POLL_RECORDS_CONFIG] = messageConsumerProperties.pollRecords
66         }
67         // TODO("Security Implementation based on type")
68         /** add or override already set properties */
69         additionalConfig?.let { configProperties.putAll(it) }
70         /** Create Kafka consumer */
71         return KafkaConsumer(configProperties)
72     }
73
74     override suspend fun subscribe(additionalConfig: Map<String, Any>?): Channel<String> {
75         /** get to topic names */
76         val consumerTopic = messageConsumerProperties.topic?.split(",")?.map { it.trim() }
77         check(!consumerTopic.isNullOrEmpty()) { "couldn't get topic information" }
78         return subscribe(consumerTopic, additionalConfig)
79     }
80
81
82     override suspend fun subscribe(consumerTopic: List<String>, additionalConfig: Map<String, Any>?): Channel<String> {
83         /** Create Kafka consumer */
84         kafkaConsumer = kafkaConsumer(additionalConfig)
85
86         checkNotNull(kafkaConsumer) {
87             "failed to create kafka consumer for " +
88                     "server(${messageConsumerProperties.bootstrapServers})'s " +
89                     "topics(${messageConsumerProperties.bootstrapServers})"
90         }
91
92         kafkaConsumer!!.subscribe(consumerTopic)
93         log.info("Successfully consumed topic($consumerTopic)")
94
95         thread(start = true, name = "KafkaConsumer") {
96             keepGoing = true
97             kafkaConsumer!!.use { kc ->
98                 while (keepGoing) {
99                     val consumerRecords = kc.poll(Duration.ofMillis(messageConsumerProperties.pollMillSec))
100                     log.info("Consumed Records : ${consumerRecords.count()}")
101                     runBlocking {
102                         consumerRecords?.forEach { consumerRecord ->
103                             /** execute the command block */
104                             consumerRecord.value()?.let {
105                                 launch {
106                                     if (!channel.isClosedForSend) {
107                                         channel.send(String(it, Charset.defaultCharset()))
108                                     } else {
109                                         log.error("Channel is closed to receive message")
110                                     }
111                                 }
112                             }
113                         }
114                     }
115                 }
116                 log.info("message listener shutting down.....")
117             }
118         }
119         return channel
120     }
121
122     override suspend fun shutDown() {
123         /** stop the polling loop */
124         keepGoing = false
125         /** Close the Channel */
126         channel.cancel()
127         /** TO shutdown gracefully, need to wait for the maximum poll time */
128         delay(messageConsumerProperties.pollMillSec)
129     }
130 }