"""Kafka message wrapper.""" import sentry_sdk import simplejson as json from neo4j_sync_dlq import config from neo4j_sync_dlq.const import neo4j_streams class KafkaMessage: """Kafka message wrapper.""" def __init__(self, record): """Create KafkaMessage instance. Args: record (ConsumerRecord): ConsumerRecord instance from KafkaConsumer """ try: self.message = json.loads(record.value.decode()) self.headers = self._load_headers(record) except json.JSONDecodeError: with sentry_sdk.push_scope() as scope: scope.set_extra('raw_message', record) sentry_sdk.capture_exception() def _load_headers(self, message): """Load and parse Kafka message headers. Args: message (dict): parsed kafka message. Returns: dict: parsed message headers. """ headers = {} if hasattr(message, 'headers'): headers_dict = dict(message.headers) for name, value in headers_dict.items(): headers[name] = value.decode('utf8').strip("'") return headers def _header(self, name): """Get a kafka message header by name. Args: name (str): header name. Returns: str: header value. """ return self.headers.get(f'{config.KAFKA_DLQ_HEADER_PREFIX}{name}') @property def table(self): """Database table name.""" return self.message.get('table') @property def database(self): """Database name.""" return self.message.get('database') @property def data(self): """Sync event data.""" return self.message.get('data') @property def offset(self): """Kafka topic offset.""" return self._header(neo4j_streams.ERROR_OFFSET) @property def partition(self): """Kafka partition.""" return self._header(neo4j_streams.ERROR_PARTITION) @property def error_class_name(self): """Neo4J streams error class name.""" return self._header(neo4j_streams.ERROR_CLASS_NAME) @property def error_exception_name(self): """Neo4J streams exception class name.""" return self._header(neo4j_streams.ERROR_EXCEPTION_NAME) @property def error_exception_message(self): """Neo4J streams exception message.""" return self._header(neo4j_streams.ERROR_EXCEPTION_MESSAGE) @property def error_exception_stacktrace(self): """Neo4J streams exception stacktrace.""" return self._header(neo4j_streams.ERROR_EXCEPTION_STACKTRACE) def to_sentry_message(self, include_stacktrace=False): """Prepare a dict with the necessary fields for logging/sentry. Args: include_stacktrace (bool): include the stacktrace or not. Returns: dict: filtered data for logging. """ message = { 'table': self.table, 'exception_name': self.error_exception_name, 'error_message': self.error_exception_message, 'data': self.data } if include_stacktrace: message['stack_trace'] = self.error_exception_stacktrace return message