class Server:
- def __init__(self, config: Config, loop=None):
+ def __init__(self, config: Config):
self.conf = config
- if not loop:
- loop = asyncio.get_event_loop()
- self.loop = loop
self.msg_bus = MessageBusClient(config)
self.service = ServerService(config)
self.service.populate_prometheus()
def run(self):
- self.loop.run_until_complete(self.start())
+ asyncio.run(self.start())
async def start(self, wait_time=5):
topics = ["alarm_request"]
while True:
try:
await self.msg_bus.aioread(topics, self._process_msg)
- log.info("Sucessfully subscribed to kafka topic(s) %s", str(topics))
break
except Exception as e:
# Failed to subscribe to kafka topic
async def _process_msg(self, topic, key, values):
log.info("Message arrived: %s", values)
try:
-
if topic == "alarm_request":
if key == "create_alarm_request":
alarm_details = values["alarm_create_request"]