- else:
- log.debug("vim_type is misconfigured or unsupported; %s",
- vim_type)
-
- elif message.topic == "vim_account":
- if message.key == "create" or message.key == "edit":
- auth_manager.store_auth_credentials(message)
- if message.key == "delete":
- auth_manager.delete_auth_credentials(message)
-
- # TODO: Remove in the near future. Modify tests accordingly.
- elif message.topic == "access_credentials":
- # Check the vim desired by the message
- vim_type = get_vim_type(message)
- if vim_type == "openstack":
- log.info("This message is for the OpenStack plugin.")
- auth_token = openstack_auth._authenticate(message=message)
-
- elif vim_type == "aws":
- log.info("This message is for the CloudWatch plugin.")
- aws_access_credentials.access_credential_calls(message)
-
- elif vim_type == "vmware":
- log.info("This access_credentials message is for the vROPs plugin.")
- vrops_rcvr.consume(message)
+kafka_logger = logging.getLogger('kafka')
+kafka_logger.setLevel(logging.getLevelName(cfg.OSMMON_KAFKA_LOG_LEVEL))
+kafka_formatter = logging.Formatter('%(asctime)s - %(name)s - %(levelname)s - %(message)s')
+kafka_handler = logging.StreamHandler(sys.stdout)
+kafka_handler.setFormatter(kafka_formatter)
+kafka_logger.addHandler(kafka_handler)
+
+
+class CommonConsumer:
+
+ def __init__(self):
+ cfg = Config.instance()
+
+ self.auth_manager = AuthManager()
+ self.database_manager = DatabaseManager()
+ self.database_manager.create_tables()
+
+ # Create OpenStack alarming and metric instances
+ self.openstack_metrics = metrics.Metrics()
+ self.openstack_alarms = alarming.Alarming()
+
+ # Create CloudWatch alarm and metric instances
+ self.cloudwatch_alarms = plugin_alarms()
+ self.cloudwatch_metrics = plugin_metrics()
+ self.aws_connection = Connection()
+ self.aws_access_credentials = AccessCredentials()
+
+ # Create vROps plugin_receiver class instance
+ self.vrops_rcvr = plugin_receiver.PluginReceiver()
+
+ log.info("Connecting to MongoDB...")
+ self.common_db = dbmongo.DbMongo()
+ common_db_uri = cfg.MONGO_URI.split(':')
+ self.common_db.db_connect({'host': common_db_uri[0], 'port': int(common_db_uri[1]), 'name': 'osm'})
+ log.info("Connection successful.")
+
+ def get_vim_type(self, vim_uuid):
+ """Get the vim type that is required by the message."""
+ credentials = self.database_manager.get_credentials(vim_uuid)
+ return credentials.type
+
+ def get_vdur(self, nsr_id, member_index, vdu_name):
+ vnfr = self.get_vnfr(nsr_id, member_index)
+ for vdur in vnfr['vdur']:
+ if vdur['name'] == vdu_name:
+ return vdur
+ raise ValueError('vdur not found for nsr-id %s, member_index %s and vdu_name %s', nsr_id, member_index,
+ vdu_name)
+
+ def get_vnfr(self, nsr_id, member_index):
+ vnfr = self.common_db.get_one("vnfrs",
+ {"nsr-id-ref": nsr_id, "member-vnf-index-ref": str(member_index)})
+ return vnfr
+
+ def run(self):
+ common_consumer = KafkaConsumer(bootstrap_servers=cfg.BROKER_URI,
+ key_deserializer=bytes.decode,
+ value_deserializer=bytes.decode,
+ group_id="mon-consumer",
+ session_timeout_ms=60000,
+ heartbeat_interval_ms=20000)
+
+ topics = ['metric_request', 'alarm_request', 'vim_account']
+ common_consumer.subscribe(topics)
+
+ log.info("Listening for messages...")
+ for message in common_consumer:
+ t = threading.Thread(target=self.consume_message, args=(message,))
+ t.start()
+
+ def consume_message(self, message):
+ log.info("Message arrived: %s", message)
+ try:
+ try:
+ values = json.loads(message.value)
+ except ValueError:
+ values = yaml.safe_load(message.value)
+
+ if message.topic == "vim_account":
+ if message.key == "create" or message.key == "edit":
+ self.auth_manager.store_auth_credentials(values)
+ if message.key == "delete":
+ self.auth_manager.delete_auth_credentials(values)