diff --git a/site/app.py b/site/app.py index 48077fe..b77e06e 100644 --- a/site/app.py +++ b/site/app.py @@ -7,9 +7,6 @@ Env vars (set by Terraform): KAFKA_TOPIC — topic name KAFKA_USERNAME — SASL username KAFKA_PASSWORD — SASL password - KAFKA_CA_CRT — CA certificate (PEM) - KAFKA_USER_CRT — user certificate (PEM) - KAFKA_USER_KEY — user private key (PEM) KAFKA_GROUP_ID — consumer group id REDIS_HOST — Redis host REDIS_PORT — Redis port (default 6379) @@ -20,7 +17,7 @@ Env vars (set by Terraform): import json import os -import tempfile +import ssl import threading from flask import Flask, jsonify @@ -35,9 +32,6 @@ KAFKA_BROKERS = os.environ["KAFKA_BROKERS"] KAFKA_TOPIC = os.environ["KAFKA_TOPIC"] KAFKA_USERNAME = os.environ["KAFKA_USERNAME"] KAFKA_PASSWORD = os.environ["KAFKA_PASSWORD"] -KAFKA_CA_CRT = os.environ.get("KAFKA_CA_CRT", "") -KAFKA_USER_CRT = os.environ.get("KAFKA_USER_CRT", "") -KAFKA_USER_KEY = os.environ.get("KAFKA_USER_KEY", "") KAFKA_GROUP_ID = os.environ.get("KAFKA_GROUP_ID", "iot-consumer-group") REDIS_HOST = os.environ["REDIS_HOST"] @@ -47,21 +41,6 @@ REDIS_PASSWORD = os.environ.get("REDIS_PASSWORD", "") MONGO_URI = os.environ["MONGO_URI"] MONGO_DB = os.environ.get("MONGO_DB", "iot") -# --- TLS certs → temp files --- -_cert_dir = tempfile.mkdtemp(prefix="kafka-certs-") - -def _write_cert(name, content): - if not content: - return None - path = os.path.join(_cert_dir, name) - with open(path, "w") as f: - f.write(content) - return path - -_ca_file = _write_cert("ca.crt", KAFKA_CA_CRT) -_cert_file = _write_cert("user.crt", KAFKA_USER_CRT) -_key_file = _write_cert("user.key", KAFKA_USER_KEY) - # --- Clients --- events_consumed = 0 errors_count = 0 @@ -129,26 +108,23 @@ def store_event(event): def consume_loop(): global events_consumed, errors_count - kwargs = { - "bootstrap_servers": KAFKA_BROKERS, - "security_protocol": "SASL_SSL", - "sasl_mechanism": "PLAIN", - "sasl_plain_username": KAFKA_USERNAME, - "sasl_plain_password": KAFKA_PASSWORD, - "group_id": KAFKA_GROUP_ID, - "auto_offset_reset": "earliest", - "value_deserializer": lambda m: json.loads(m.decode("utf-8")), - "max_poll_records": 20, - } - if _ca_file: - kwargs["ssl_cafile"] = _ca_file - if _cert_file: - kwargs["ssl_certfile"] = _cert_file - if _key_file: - kwargs["ssl_keyfile"] = _key_file - kwargs["ssl_check_hostname"] = False + ctx = ssl.create_default_context() + ctx.check_hostname = False + ctx.verify_mode = ssl.CERT_NONE - consumer = KafkaConsumer(KAFKA_TOPIC, **kwargs) + consumer = KafkaConsumer( + KAFKA_TOPIC, + bootstrap_servers=KAFKA_BROKERS, + security_protocol="SASL_SSL", + sasl_mechanism="PLAIN", + sasl_plain_username=KAFKA_USERNAME, + sasl_plain_password=KAFKA_PASSWORD, + ssl_context=ctx, + group_id=KAFKA_GROUP_ID, + auto_offset_reset="earliest", + value_deserializer=lambda m: json.loads(m.decode("utf-8")), + max_poll_records=20, + ) print(f"[CONSUMER] Listening on topic '{KAFKA_TOPIC}'...") for message in consumer: