"""Kafka connectors.""" import ssl import jks from neo4j_sync_dlq import config def create_ssl_context(cacerts_file, cacerts_pass='changeit'): """Create SSL context with JVM CA certificates. Args: cacerts_file (str): absolute path to JVM cacerts truststore file cacerts_pass (str): truststore password ("changeit" is the JVM default) Returns: ssl.SSLContext: SSLContext instance to use with KafkaClient. """ truststore = jks.KeyStore.load(cacerts_file, cacerts_pass) ctx = ssl.SSLContext(protocol=ssl.PROTOCOL_TLSv1_2) for alias, cert in truststore.certs.items(): ctx.load_verify_locations(cadata=cert.cert) return ctx def get_kafka_config(): """Get kafka connection configuration. Returns: dict: kafka connection params. """ conf = { 'bootstrap_servers': config.KAFKA_BOOTSTRAP_SERVERS, 'security_protocol': 'SSL' } if config.ENVIRONMENT == config.ENVIRONMENT_DEV: conf['ssl_context'] = create_ssl_context( config.CACERTS_FILE, config.CACERTS_PASSWORD) conf['auto_offset_reset'] = 'earliest' conf['client_id'] = 'PythonNeo4J' return conf