import asyncio
import json
import logging
+import time
from osm_mon.core.config import Config
from osm_mon.core.message_bus_client import MessageBusClient
def run(self):
self.loop.run_until_complete(self.start())
- async def start(self):
+ async def start(self, wait_time=5):
topics = [
"alarm_request"
]
- await self.msg_bus.aioread(topics, self._process_msg)
+ while True:
+ try:
+ await self.msg_bus.aioread(topics, self._process_msg)
+ break
+ except Exception as e:
+ # Failed to subscribe to kafka topic
+ log.exception("Error when subscribing to topics %s", str(topics))
+ log.exception("Exception %s", str(e))
+ # Wait for some time for kaka to stabilize and then reattempt to subscribe again
+ time.sleep(wait_time)
+ log.info("Retrying to subscribe the kafka topic(s) %s", str(topics))
async def _process_msg(self, topic, key, values):
log.info("Message arrived: %s", values)
alarm_details['severity'].lower(),
alarm_details['statistic'].lower(),
alarm_details['metric_name'],
- alarm_details['vdu_name'],
- alarm_details['vnf_member_index'],
- alarm_details['ns_id']
+ alarm_details['tags']
)
response = response_builder.generate_response('create_alarm_response',
cor_id=cor_id,