Updating README.md
[ccsdk/cds.git] / ms / blueprintsprocessor / functions / message-prioritizaion / src / main / kotlin / org / onap / ccsdk / cds / blueprintsprocessor / functions / message / prioritization / nats / NatsMessagePrioritizationConsumer.kt
1 /*
2  * Copyright © 2018-2019 AT&T Intellectual Property.
3  *
4  * Licensed under the Apache License, Version 2.0 (the "License");
5  * you may not use this file except in compliance with the License.
6  * You may obtain a copy of the License at
7  *
8  *     http://www.apache.org/licenses/LICENSE-2.0
9  *
10  * Unless required by applicable law or agreed to in writing, software
11  * distributed under the License is distributed on an "AS IS" BASIS,
12  * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13  * See the License for the specific language governing permissions and
14  * limitations under the License.
15  */
16
17 package org.onap.ccsdk.cds.blueprintsprocessor.functions.message.prioritization.nats
18
19 import io.nats.streaming.MessageHandler
20 import io.nats.streaming.Subscription
21 import kotlinx.coroutines.runBlocking
22 import org.onap.ccsdk.cds.blueprintsprocessor.functions.message.prioritization.MessagePrioritizationService
23 import org.onap.ccsdk.cds.blueprintsprocessor.functions.message.prioritization.db.MessagePrioritization
24 import org.onap.ccsdk.cds.blueprintsprocessor.nats.asJsonType
25 import org.onap.ccsdk.cds.blueprintsprocessor.nats.service.BluePrintNatsLibPropertyService
26 import org.onap.ccsdk.cds.blueprintsprocessor.nats.service.BluePrintNatsService
27 import org.onap.ccsdk.cds.blueprintsprocessor.nats.utils.NatsClusterUtils
28 import org.onap.ccsdk.cds.blueprintsprocessor.nats.utils.SubscriptionOptionsUtils
29 import org.onap.ccsdk.cds.controllerblueprints.core.BluePrintProcessorException
30 import org.onap.ccsdk.cds.controllerblueprints.core.asType
31 import org.onap.ccsdk.cds.controllerblueprints.core.logger
32 import org.onap.ccsdk.cds.controllerblueprints.core.utils.ClusterUtils
33
34 open class NatsMessagePrioritizationConsumer(
35     private val bluePrintNatsLibPropertyService: BluePrintNatsLibPropertyService,
36     private val natsMessagePrioritizationService: MessagePrioritizationService
37 ) {
38     private val log = logger(NatsMessagePrioritizationConsumer::class)
39
40     lateinit var bluePrintNatsService: BluePrintNatsService
41     private lateinit var subscription: Subscription
42
43     suspend fun startConsuming() {
44         val prioritizationConfiguration = natsMessagePrioritizationService.getConfiguration()
45         val natsConfiguration = prioritizationConfiguration.natsConfiguration
46             ?: throw BluePrintProcessorException("couldn't get NATS consumer configuration")
47
48         check((natsMessagePrioritizationService is AbstractNatsMessagePrioritizationService)) {
49             "messagePrioritizationService is not of type AbstractNatsMessagePrioritizationService."
50         }
51         bluePrintNatsService = consumerService(natsConfiguration.connectionSelector)
52         natsMessagePrioritizationService.bluePrintNatsService = bluePrintNatsService
53         val inputSubject = NatsClusterUtils.currentApplicationSubject(natsConfiguration.inputSubject)
54         val loadBalanceGroup = ClusterUtils.applicationName()
55         val messageHandler = createMessageHandler()
56         val subscriptionOptions = SubscriptionOptionsUtils.durable(NatsClusterUtils.currentNodeDurable(inputSubject))
57         subscription = bluePrintNatsService.loadBalanceSubscribe(
58             inputSubject,
59             loadBalanceGroup,
60             messageHandler,
61             subscriptionOptions
62         )
63         log.info(
64             "Nats prioritization consumer listening on subject($inputSubject) on loadBalance group($loadBalanceGroup)."
65         )
66     }
67
68     suspend fun shutDown() {
69         if (::subscription.isInitialized) {
70             subscription.unsubscribe()
71         }
72         log.info("Nats prioritization consumer listener shutdown complete")
73     }
74
75     private fun consumerService(selector: String): BluePrintNatsService {
76         return bluePrintNatsLibPropertyService.bluePrintNatsService(selector)
77     }
78
79     private fun createMessageHandler(): MessageHandler {
80         return MessageHandler { message ->
81             try {
82                 val messagePrioritization = message.asJsonType().asType(MessagePrioritization::class.java)
83                 runBlocking {
84                     natsMessagePrioritizationService.prioritize(messagePrioritization)
85                 }
86             } catch (e: Exception) {
87                 log.error("failed to process prioritize message", e)
88             }
89         }
90     }
91 }