Commit 096d8137 authored by Lluis Gifre Renom's avatar Lluis Gifre Renom
Browse files

NBI component - VNTM recommendations:

- Hardened parsing of notifications from e2e orchestrator
parent bdca3f5d
Loading
Loading
Loading
Loading
+31 −21
Original line number Diff line number Diff line
@@ -45,32 +45,42 @@ class VntRecommServerNamespace(Namespace):
        LOGGER.debug(MSG.format(str(request.sid), str(reason)))
        leave_room(SIO_ROOM, namespace=SIO_NAMESPACE)

    def on_vlink_created(self, data):
        MSG = '[on_vlink_created] begin: sid={:s}, data={:s}'
        LOGGER.debug(MSG.format(str(request.sid), str(data)))
    @staticmethod
    def _parse_payload(data):
        if isinstance(data, str):
            return json.loads(data)
        if isinstance(data, dict):
            return dict(data)
        raise TypeError('Unsupported recommendation callback payload type: {:s}'.format(type(data).__name__))

    def _publish_reply(self, event_name: str, data) -> None:
        sid = getattr(request, 'sid', '<unknown>')
        LOGGER.info('[%s] begin: sid=%s payload=%s', event_name, sid, str(data))

        data = json.loads(data)
        request_key = str(data.pop('_request_key')).encode('utf-8')
        vntm_reply = json.dumps({'event': 'vlink_created', 'data': data}).encode('utf-8')
        LOGGER.debug('[on_vlink_created] request_key={:s}/{:s}'.format(str(type(request_key)), str(request_key)))
        LOGGER.debug('[on_vlink_created] vntm_reply={:s}/{:s}'.format(str(type(vntm_reply)), str(vntm_reply)))
        json_data = self._parse_payload(data)
        request_key = str(json_data.pop('_request_key')).encode('utf-8')
        vntm_reply = json.dumps({'event': event_name, 'data': json_data}).encode('utf-8')

        LOGGER.info(
            '[%s] Publishing Kafka reply: request_key=%s payload=%s',
            event_name, request_key.decode('utf-8'), vntm_reply.decode('utf-8')
        )
        self.kafka_producer.send(
            KafkaTopic.VNTMANAGER_RESPONSE.value, key=request_key, value=vntm_reply
        )
        self.kafka_producer.flush()
        LOGGER.info('[%s] Kafka reply published', event_name)

    def on_vlink_removed(self, data):
        MSG = '[on_vlink_removed] begin: sid={:s}, data={:s}'
        LOGGER.debug(MSG.format(str(request.sid), str(data)))

        data = json.loads(data)
        request_key = str(data.pop('_request_key')).encode('utf-8')
        vntm_reply = json.dumps({'event': 'vlink_removed', 'data': data}).encode('utf-8')
        LOGGER.debug('[on_vlink_removed] request_key={:s}/{:s}'.format(str(type(request_key)), str(request_key)))
        LOGGER.debug('[on_vlink_removed] vntm_reply={:s}/{:s}'.format(str(type(vntm_reply)), str(vntm_reply)))
    def on_vlink_created(self, data):
        try:
            self._publish_reply('vlink_created', data)
        except Exception:
            LOGGER.exception('[on_vlink_created] Failed to process callback')
            raise

        self.kafka_producer.send(
            KafkaTopic.VNTMANAGER_RESPONSE.value, key=request_key, value=vntm_reply
        )
        self.kafka_producer.flush()
    def on_vlink_removed(self, data):
        try:
            self._publish_reply('vlink_removed', data)
        except Exception:
            LOGGER.exception('[on_vlink_removed] Failed to process callback')
            raise