From 154ebc88efb77b56783e741c88be359b50ea0f6a Mon Sep 17 00:00:00 2001 From: Valerij Maljulin Date: Wed, 23 Sep 2026 12:25:51 +0200 Subject: [PATCH] Add Kafka listener JIRA: RHELWF-13776 Assisted-by: Opus 4.6 --- docker-compose.yml | 40 +- docker/broker.xml | 10 - docker/dev.env | 2 + docker/greenwave-settings.py | 24 ++ docker/resultsdb-settings.py | 17 +- docker/waiverdb-settings.py | 16 +- docker/waiverdb.env | 2 + greenwave/config.py | 26 ++ greenwave/listeners/base.py | 141 ++++++- greenwave/listeners/kafka.py | 112 ++++++ greenwave/listeners/resultsdb.py | 5 +- greenwave/listeners/waiverdb.py | 5 +- greenwave/tests/test_listeners.py | 2 + greenwave/tests/test_listeners_kafka.py | 471 ++++++++++++++++++++++++ pyproject.toml | 1 + uv.lock | 25 ++ 16 files changed, 850 insertions(+), 49 deletions(-) delete mode 100644 docker/broker.xml create mode 100644 greenwave/listeners/kafka.py create mode 100644 greenwave/tests/test_listeners_kafka.py diff --git a/docker-compose.yml b/docker-compose.yml index d2036162..c8647ec7 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -43,9 +43,11 @@ services: test: "pg_isready -U postgres || exit 1" resultsdb: - image: quay.io/factory2/resultsdb:latest@sha256:d75d6e737c677104fb7e96469139141239afaef72945b43c6ad1f613f5de21fd + image: quay.io/redhat-user-workloads/exd-sp-rhel-wf-tenant/resultsdb:latest@sha256:c3862d7008c788b43e084e8b272c383b7f654913bbf9ea9bd2edc3f3b332f338 environment: - GREENWAVE_LISTENERS=${GREENWAVE_LISTENERS:-1} + - RESULTSDB_KAFKA_SASL_USERNAME=resultsdb + - RESULTSDB_KAFKA_SASL_PASSWORD=resultsdb command: ["bash", "-c", "/start.sh"] volumes: - ./docker/home:/home/dev:rw,z @@ -67,7 +69,7 @@ services: env_file: ["docker/waiverdb-db.env"] waiverdb: - image: quay.io/factory2/waiverdb:latest@sha256:fa167759506262ff21003e9f773d50c3e1bf0054580e5fec835c3a0d028187a9 + image: quay.io/redhat-user-workloads/exd-sp-rhel-wf-tenant/waiverdb:latest@sha256:84f032c6473510e2e374323d3f641cead1d66a0e03ccb16f1797a9829b80f784 env_file: ["docker/waiverdb.env"] environment: - GREENWAVE_LISTENERS=${GREENWAVE_LISTENERS:-1} @@ -126,14 +128,36 @@ services: replicas: ${GREENWAVE_LISTENERS:-1} message-broker: - image: docker.io/apache/activemq-artemis:latest-alpine@sha256:1bce124d2324faeb1253e2db4bb68d64decf46412ecfd8d9ef46c08b1b4af5b4 + image: docker.io/apache/kafka:3.9.0@sha256:fbc7d7c428e3755cf36518d4976596002477e4c052d1f80b5b9eafd06d0fff2f + hostname: message-broker restart: unless-stopped - volumes: - - ./docker/broker.xml:/var/lib/artemis-instance/etc-override/broker.xml:ro,z ports: - - 127.0.0.1:5671:5671 # amqp - - 127.0.0.1:61612:61612 # stomp - - 127.0.0.1:8162:8161 # http + - 127.0.0.1:9092:9092 + environment: + KAFKA_NODE_ID: 1 + KAFKA_PROCESS_ROLES: broker,controller + KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093 + KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://message-broker:9092 + KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT + KAFKA_CONTROLLER_QUORUM_VOTERS: 1@message-broker:9093 + KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER + KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT + KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 + KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 + KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 + KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0 + KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true" + CLUSTER_ID: MkU3OEVBNTcwNTJENDM2Qk + healthcheck: + test: + [ + "CMD-SHELL", + "/opt/kafka/bin/kafka-broker-api-versions.sh --bootstrap-server localhost:9092", + ] + interval: 5s + timeout: 10s + retries: 12 + start_period: 20s deploy: replicas: ${GREENWAVE_LISTENERS:-1} diff --git a/docker/broker.xml b/docker/broker.xml deleted file mode 100644 index b3990079..00000000 --- a/docker/broker.xml +++ /dev/null @@ -1,10 +0,0 @@ - - - false - - tcp://0.0.0.0:61616?tcpSendBufferSize=1048576;tcpReceiveBufferSize=1048576;protocols=CORE,AMQP,STOMP - tcp://0.0.0.0:5671?protocols=AMQP - tcp://0.0.0.0:61612?protocols=STOMP - - - diff --git a/docker/dev.env b/docker/dev.env index dfcc7323..400b66fb 100644 --- a/docker/dev.env +++ b/docker/dev.env @@ -6,3 +6,5 @@ RESULTSDB_TEST_URL=http://resultsdb:5001/ PYTEST_ADDOPTS=-o cache_dir=/home/dev/.pytest_cache PYTHONPATH=/src GREENWAVE_STATSD_HOST=statsd:9125 +GREENWAVE_KAFKA_SASL_USERNAME=greenwave +GREENWAVE_KAFKA_SASL_PASSWORD=greenwave diff --git a/docker/greenwave-settings.py b/docker/greenwave-settings.py index 4e7727b7..b61cadd4 100644 --- a/docker/greenwave-settings.py +++ b/docker/greenwave-settings.py @@ -15,6 +15,27 @@ "reconnect_sleep_max": 10.0, "reconnect_attempts_max": 5, } +LISTENER_BACKEND = "kafka" +KAFKA = { + "resultsdb_topic": "eng.resultsdb.result.new", + "waiverdb_topic": "eng.waiverdb.waiver.new", + "decision_topic": "eng.greenwave.decision.update", + "consumer": { + "bootstrap.servers": "message-broker:9092", + "client.id": "greenwave", + "enable.auto.commit": False, + "auto.offset.reset": "earliest", + "security.protocol": "PLAINTEXT", + }, + "producer": { + "bootstrap.servers": "message-broker:9092", + "client.id": "greenwave", + "retries": 3, + "retry.backoff.ms": 100, + "security.protocol": "PLAINTEXT", + }, + "flush_timeout_seconds": 20.0, +} CACHE = { # 'backend': 'dogpile.cache.null', "backend": "dogpile.cache.pymemcache", @@ -35,6 +56,9 @@ "stomp.py": { "level": "DEBUG", }, + "confluent_kafka": { + "level": "DEBUG", + }, }, "handlers": { "console": { diff --git a/docker/resultsdb-settings.py b/docker/resultsdb-settings.py index 3e0a9948..1d0e563c 100644 --- a/docker/resultsdb-settings.py +++ b/docker/resultsdb-settings.py @@ -11,12 +11,15 @@ ADDITIONAL_RESULT_OUTCOMES = ("RUNNING", "QUEUED", "ERROR") MESSAGE_BUS_PUBLISH = os.environ.get("GREENWAVE_LISTENERS", "") not in ("", "0") -MESSAGE_BUS_PLUGIN = "stomp" -MESSAGE_BUS_KWARGS = { - "modname": "resultsdb", - "destination": "/topic/VirtualTopic.eng.resultsdb.result.new", - "connection": { - "host_and_ports": [("message-broker", 61612)], - "use_ssl": False, +MESSAGE_BUS_PLUGIN = "kafka" +KAFKA = { + "topic": "eng.resultsdb.result.new", + "producer": { + "bootstrap.servers": "message-broker:9092", + "client.id": "resultsdb", + "retries": 3, + "retry.backoff.ms": 100, + "security.protocol": "PLAINTEXT", }, + "flush_timeout_seconds": 20.0, } diff --git a/docker/waiverdb-settings.py b/docker/waiverdb-settings.py index d26ce1dc..b6efee7a 100644 --- a/docker/waiverdb-settings.py +++ b/docker/waiverdb-settings.py @@ -14,11 +14,15 @@ RESULTSDB_API_URL = "http://resultsdb:5001/api/v2.0" # NOSONAR MESSAGE_BUS_PUBLISH = os.environ.get("GREENWAVE_LISTENERS", "") not in ("", "0") -MESSAGE_PUBLISHER = "stomp" -STOMP_CONFIGS = { - "destination": "/topic/VirtualTopic.eng.waiverdb.waiver.new", - "connection": { - "host_and_ports": [("message-broker", 61612)], - "use_ssl": False, +MESSAGE_PUBLISHER = "kafka" +KAFKA = { + "topic": "eng.waiverdb.waiver.new", + "producer": { + "bootstrap.servers": "message-broker:9092", + "client.id": "waiverdb", + "retries": 3, + "retry.backoff.ms": 100, + "security.protocol": "PLAINTEXT", }, + "flush_timeout_seconds": 20.0, } diff --git a/docker/waiverdb.env b/docker/waiverdb.env index bd3f546a..95b37920 100644 --- a/docker/waiverdb.env +++ b/docker/waiverdb.env @@ -3,3 +3,5 @@ SECRET_KEY=waiverdb WAIVERDB_CONFIG=/etc/waiverdb/settings.py WAIVERDB_TEST_URL=http://waiverdb:5004/ PYTEST_ADDOPTS=-o cache_dir=/tmp/.pytest_cache +WAIVERDB_KAFKA_SASL_USERNAME=waiverdb +WAIVERDB_KAFKA_SASL_PASSWORD=waiverdb diff --git a/greenwave/config.py b/greenwave/config.py index 90a918f4..8e45ccc0 100644 --- a/greenwave/config.py +++ b/greenwave/config.py @@ -97,6 +97,32 @@ class Config: "ca_certs": "/etc/pki/umb/umb-ca", } + # Listener message bus: "stomp" (UMB) or "kafka" (MSK). + LISTENER_BACKEND = "stomp" + # Used when LISTENER_BACKEND is "kafka". + # "consumer" and "producer" keys are passed to confluent-kafka (librdkafka). + # Full reference: https://github.com/confluentinc/librdkafka/blob/master/CONFIGURATION.md + # SASL credentials are read from GREENWAVE_KAFKA_SASL_USERNAME and + # GREENWAVE_KAFKA_SASL_PASSWORD environment variables. + KAFKA = { + "resultsdb_topic": "eng.resultsdb.result.new", + "waiverdb_topic": "eng.waiverdb.waiver.new", + "decision_topic": "eng.greenwave.decision.update", + "consumer": { + "bootstrap.servers": "localhost:9092", + "client.id": "greenwave", + "enable.auto.commit": False, + "auto.offset.reset": "latest", + }, + "producer": { + "bootstrap.servers": "localhost:9092", + "client.id": "greenwave", + "retries": 3, + "retry.backoff.ms": 100, + }, + "flush_timeout_seconds": 20.0, + } + OTEL_EXPORTER_OTLP_TRACES_ENDPOINT = None OTEL_EXPORTER_SERVICE_NAME = "greenwave" diff --git a/greenwave/listeners/base.py b/greenwave/listeners/base.py index 547de72d..a75d7e2f 100644 --- a/greenwave/listeners/base.py +++ b/greenwave/listeners/base.py @@ -7,6 +7,7 @@ import uuid import stomp +from confluent_kafka import KafkaError from opentelemetry.context import Context from opentelemetry.trace.propagation.tracecontext import ( TraceContextTextMapPropagator, @@ -15,6 +16,7 @@ from werkzeug.exceptions import HTTPException import greenwave.app_factory +from greenwave.listeners.kafka import KafkaBus from greenwave.logger import init_logging, log_to_stdout from greenwave.monitor import ( decision_changed_counter, @@ -83,15 +85,22 @@ def __init__(self, uid_suffix, config_obj=None): self.connecting = False self.stop = False + self.uid_suffix = uid_suffix self.uid = f"{GREENWAVE_LISTENER_PREFIX}-{uid_suffix}-{uuid.uuid1().hex}" init_logging() log_to_stdout(logging.DEBUG) self.app = greenwave.app_factory.create_app(config_obj) - self.destination = self.app.config["LISTENER_DECISION_UPDATE_DESTINATION"] + self._backend = self.app.config.get("LISTENER_BACKEND", "stomp") + if self._backend == "kafka": + self.destination = self.app.config["KAFKA"]["decision_topic"] + else: + self.destination = self.app.config["LISTENER_DECISION_UPDATE_DESTINATION"] self.context = None + self._kafka_bus = None + self._kafka_thread = None def on_error(self, frame): self.app.logger.warning("Received an error: %s", frame.body) @@ -160,10 +169,6 @@ def on_receiver_loop_completed(self, frame): self._terminate() def listen(self): - if self.connection is not None: - self.app.logger.warning("Already connected") - return - def handler(signum, frame): self.app.logger.warning("Stopping listener on signal %s", signum) self.disconnect() @@ -171,6 +176,16 @@ def handler(signum, frame): if threading.current_thread() is threading.main_thread(): signal.signal(signal.SIGTERM, handler) + if self._backend == "kafka": + self._listen_kafka() + else: + self._listen_stomp() + + def _listen_stomp(self): + if self.connection is not None: + self.app.logger.warning("Already connected") + return + hosts = self.app.config["LISTENER_HOSTS"] hosts_and_ports = [tuple(url.split(":")) for url in hosts.split(",")] connection_args = self.app.config["LISTENER_CONNECTION"] @@ -191,14 +206,97 @@ def handler(signum, frame): self.app.logger.info("Listening on %s", self.topic) + def _listen_kafka(self): + if self._kafka_bus is not None: + self.app.logger.warning("Already connected") + return + + self._kafka_bus = KafkaBus( + self.app.config, group_id=f"greenwave-{self.uid_suffix}" + ) + self._kafka_bus.subscribe(self.topic) + self._kafka_thread = threading.Thread( + target=self._kafka_loop, + name=f"greenwave-kafka-{self.uid_suffix}", + daemon=True, + ) + self._kafka_thread.start() + self.app.logger.info("Listening on %s", self.topic) + + def _kafka_loop(self): + try: + while True: + with self.connection_condition: + if self.stop: + return + msg = self._kafka_bus.poll(1.0) + if msg is None: + continue + error = msg.error() + if error: + if error.code() == KafkaError._PARTITION_EOF: + continue + self.app.logger.error("Kafka consumer error: %s", error) + if error.fatal(): + self._terminate() + return + continue + self._handle_kafka_message(msg) + except Exception: + self.app.logger.exception("Kafka consumer loop failed") + self._terminate() + + def _handle_kafka_message(self, msg): + with self.connection_condition: + if self.stop: + return + + self.app.logger.debug("Received a Kafka message: %s", msg.offset()) + self._inc(messaging_rx_counter) + + try: + data = json.loads(msg.value()) + except (json.JSONDecodeError, TypeError) as e: + self.app.logger.debug("Failed to decode JSON message: %s", e) + self._inc(messaging_rx_ignored_counter) + with self.connection_condition: + if not self.stop: + self._kafka_bus.commit(msg) + return + + try: + with self.app.app_context(): + processed = self._consume_message(data) + except BaseException: # NOSONAR + self._inc(messaging_rx_failed_counter) + raise + + if processed: + self._inc(messaging_rx_processed_ok_counter) + else: + self._inc(messaging_rx_ignored_counter) + with self.connection_condition: + if not self.stop: + self._kafka_bus.commit(msg) + + def _stop_connections(self): + self.stop = True + if self._kafka_bus is not None: + self._kafka_bus.close() + if self.connection is not None: + self.connection.disconnect() + def disconnect(self): self.app.logger.debug("Disconnecting listener") with self.connection_condition: - self.stop = True - self.connection.disconnect() + self._stop_connections() + + if self._kafka_thread is not None: + self._kafka_thread.join(timeout=5) def _terminate(self): - self.disconnect() + with self.connection_condition: + self._stop_connections() os.kill(os.getpid(), signal.SIGQUIT) # NOSONAR def _consume_message(self, message): @@ -217,16 +315,27 @@ def _publish_decision_update(self, decision): TraceContextTextMapPropagator().inject(decision, self.context) message = {"msg": decision, "topic": self.destination} body = json.dumps(message) + headers = { + "subject_type": decision["subject_type"], + "subject_identifier": decision["subject_identifier"], + "product_version": decision["product_version"], + "decision_context": decision["decision_context"], + "policies_satisfied": str(decision["policies_satisfied"]).lower(), + "summary": decision["summary"], + } + if self._backend == "kafka": + try: + with self.connection_condition: + self._kafka_bus.publish(self.destination, body, headers) + except Exception: + self.app.logger.exception("Error sending decision update message") + self._inc(messaging_tx_failed_counter) + raise + self._inc(messaging_tx_sent_ok_counter) + return + while True: try: - headers = { - "subject_type": decision["subject_type"], - "subject_identifier": decision["subject_identifier"], - "product_version": decision["product_version"], - "decision_context": decision["decision_context"], - "policies_satisfied": str(decision["policies_satisfied"]).lower(), - "summary": decision["summary"], - } self.connection.send( body=body, headers=headers, destination=self.destination ) diff --git a/greenwave/listeners/kafka.py b/greenwave/listeners/kafka.py new file mode 100644 index 00000000..bd473fac --- /dev/null +++ b/greenwave/listeners/kafka.py @@ -0,0 +1,112 @@ +# SPDX-License-Identifier: GPL-2.0+ +import logging +import os +from typing import Any + +from confluent_kafka import Consumer, KafkaError, KafkaException, Producer +from confluent_kafka import Message as KafkaMessage + +_log = logging.getLogger(__name__) + +REQUIRED_KAFKA_KEYS = ( + "resultsdb_topic", + "waiverdb_topic", + "decision_topic", + "consumer", + "producer", +) + + +def parse_kafka_config( + app_config, +) -> tuple[dict[str, Any], dict[str, Any], dict[str, Any]]: + config = app_config.get("KAFKA") + if not isinstance(config, dict): + raise RuntimeError( + f"KAFKA configuration is invalid, expected a dict, got: {config!r}" + ) + + missing = [key for key in REQUIRED_KAFKA_KEYS if key not in config] + if missing: + raise RuntimeError(f"Invalid KAFKA configuration: missing {', '.join(missing)}") + if not isinstance(config["consumer"], dict) or not isinstance( + config["producer"], dict + ): + raise RuntimeError( + "Invalid KAFKA configuration: consumer and producer must be dicts" + ) + + username = os.environ.get("GREENWAVE_KAFKA_SASL_USERNAME") + password = os.environ.get("GREENWAVE_KAFKA_SASL_PASSWORD") + if not username or not password: + raise RuntimeError( + "GREENWAVE_KAFKA_SASL_USERNAME and GREENWAVE_KAFKA_SASL_PASSWORD " + "environment variables are required" + ) + + sasl = {"sasl.username": username, "sasl.password": password} + consumer_config = {**config["consumer"], **sasl} + producer_config = {**config["producer"], **sasl} + return config, consumer_config, producer_config + + +class KafkaBus: + """Kafka consumer and producer used by Greenwave listeners.""" + + def __init__(self, app_config, group_id: str) -> None: + kafka_config, consumer_config, producer_config = parse_kafka_config(app_config) + consumer_config = {**consumer_config, "group.id": group_id} + consumer_config.setdefault("enable.auto.commit", False) + self.config = kafka_config + self.flush_timeout_seconds = float( + kafka_config.get("flush_timeout_seconds", 20.0) + ) + self.consumer = Consumer(consumer_config) + self.producer = Producer(producer_config) + + def subscribe(self, topic: str) -> None: + self.consumer.subscribe([topic]) + + def poll(self, timeout: float = 1.0) -> KafkaMessage | None: + return self.consumer.poll(timeout) + + def commit(self, msg: KafkaMessage) -> None: + self.consumer.commit(message=msg) + + def publish(self, topic: str, body: str, headers: dict[str, str]) -> None: + delivery_error = None + + def _delivery_callback(err: KafkaError | None, _msg: KafkaMessage) -> None: + nonlocal delivery_error + if err is not None: + delivery_error = delivery_error or KafkaException(err) + + kafka_headers: list[tuple[str, str | bytes | None]] = [ + (key, value.encode("utf-8")) for key, value in headers.items() + ] + self.producer.produce( + topic, + value=body.encode("utf-8"), + headers=kafka_headers, + on_delivery=_delivery_callback, + ) + remaining = self.producer.flush(timeout=self.flush_timeout_seconds) + if remaining > 0: + raise KafkaException( + KafkaError( + KafkaError._MSG_TIMED_OUT, + f"{remaining} message(s) were not delivered within timeout", + ) + ) + if delivery_error is not None: + raise delivery_error + + def close(self) -> None: + try: + self.consumer.close() + except Exception: + _log.debug("Error closing Kafka consumer", exc_info=True) + try: + self.producer.flush() + except Exception: + _log.debug("Error flushing Kafka producer", exc_info=True) diff --git a/greenwave/listeners/resultsdb.py b/greenwave/listeners/resultsdb.py index 13d487ea..f67c0f53 100644 --- a/greenwave/listeners/resultsdb.py +++ b/greenwave/listeners/resultsdb.py @@ -34,7 +34,10 @@ class ResultsDBListener(BaseListener): def __init__(self, config_obj=None): super().__init__(uid_suffix="resultsdb", config_obj=config_obj) - self.topic = self.app.config["LISTENER_RESULTSDB_QUEUE"] + if self._backend == "kafka": + self.topic = self.app.config["KAFKA"]["resultsdb_topic"] + else: + self.topic = self.app.config["LISTENER_RESULTSDB_QUEUE"] self.koji_base_url = self.app.config["KOJI_BASE_URL"] @staticmethod diff --git a/greenwave/listeners/waiverdb.py b/greenwave/listeners/waiverdb.py index af35cdd7..02512f3f 100644 --- a/greenwave/listeners/waiverdb.py +++ b/greenwave/listeners/waiverdb.py @@ -8,7 +8,10 @@ class WaiverDBListener(BaseListener): def __init__(self, config_obj=None): super().__init__(uid_suffix="waiverdb", config_obj=config_obj) - self.topic = self.app.config["LISTENER_WAIVERDB_QUEUE"] + if self._backend == "kafka": + self.topic = self.app.config["KAFKA"]["waiverdb_topic"] + else: + self.topic = self.app.config["LISTENER_WAIVERDB_QUEUE"] self.koji_base_url = self.app.config["KOJI_BASE_URL"] def _consume_message(self, msg): diff --git a/greenwave/tests/test_listeners.py b/greenwave/tests/test_listeners.py index c24b89aa..a849a8fd 100644 --- a/greenwave/tests/test_listeners.py +++ b/greenwave/tests/test_listeners.py @@ -879,6 +879,8 @@ def test_listener_resultsdb_subscribe_after_connect(mock_connection): id=listener.uid, ack="client-individual", ) + listener.listen() + assert len(mock_connection.connect.mock_calls) == 1 def test_listener_waiverdb_subscribe_after_connect(mock_connection): diff --git a/greenwave/tests/test_listeners_kafka.py b/greenwave/tests/test_listeners_kafka.py new file mode 100644 index 00000000..290e1d25 --- /dev/null +++ b/greenwave/tests/test_listeners_kafka.py @@ -0,0 +1,471 @@ +# SPDX-License-Identifier: GPL-2.0+ +import json +from unittest.mock import Mock, patch + +from confluent_kafka import KafkaError, KafkaException +from pytest import fixture, raises + +from greenwave.config import TestingConfig +from greenwave.listeners.kafka import KafkaBus, parse_kafka_config +from greenwave.listeners.resultsdb import ResultsDBListener +from greenwave.listeners.waiverdb import WaiverDBListener + + +class KafkaTestingConfig(TestingConfig): + LISTENER_BACKEND = "kafka" + KAFKA = { + "resultsdb_topic": "qa.eng.resultsdb.result.new", + "waiverdb_topic": "qa.eng.waiverdb.waiver.new", + "decision_topic": "qa.eng.greenwave.decision.update", + "consumer": { + "bootstrap.servers": "localhost:9092", + "client.id": "greenwave-test", + "enable.auto.commit": False, + "auto.offset.reset": "latest", + }, + "producer": { + "bootstrap.servers": "localhost:9092", + "client.id": "greenwave-test", + "retries": 3, + }, + "flush_timeout_seconds": 15.0, + } + + +@fixture +def kafka_env(monkeypatch): + monkeypatch.setenv("GREENWAVE_KAFKA_SASL_USERNAME", "alice") + monkeypatch.setenv("GREENWAVE_KAFKA_SASL_PASSWORD", "secret") + + +def _kafka_config(): + return { + "KAFKA": { + "resultsdb_topic": "qa.eng.resultsdb.result.new", + "waiverdb_topic": "qa.eng.waiverdb.waiver.new", + "decision_topic": "qa.eng.greenwave.decision.update", + "consumer": { + "bootstrap.servers": "localhost:9092", + "client.id": "greenwave-test", + }, + "producer": { + "bootstrap.servers": "localhost:9092", + "client.id": "greenwave-test", + }, + "flush_timeout_seconds": 15.0, + } + } + + +def test_parse_kafka_config_injects_sasl(kafka_env): + kafka_config, consumer_config, producer_config = parse_kafka_config(_kafka_config()) + assert kafka_config["decision_topic"] == "qa.eng.greenwave.decision.update" + assert consumer_config["sasl.username"] == "alice" + assert producer_config["sasl.password"] == "secret" + + +def test_parse_kafka_config_missing_sasl(monkeypatch): + monkeypatch.delenv("GREENWAVE_KAFKA_SASL_USERNAME", raising=False) + monkeypatch.delenv("GREENWAVE_KAFKA_SASL_PASSWORD", raising=False) + with raises(RuntimeError, match="GREENWAVE_KAFKA_SASL_USERNAME"): + parse_kafka_config(_kafka_config()) + + +def test_parse_kafka_config_invalid(): + with raises(RuntimeError, match="Invalid KAFKA configuration"): + parse_kafka_config({"KAFKA": {"producer": {}}}) + + +def test_parse_kafka_config_not_dict(): + with raises(RuntimeError, match="expected a dict"): + parse_kafka_config({"KAFKA": None}) + + +def test_resultsdb_listener_uses_kafka_topics(): + listener = ResultsDBListener(config_obj=KafkaTestingConfig) + assert listener._backend == "kafka" + assert listener.topic == "qa.eng.resultsdb.result.new" + assert listener.destination == "qa.eng.greenwave.decision.update" + + +def test_waiverdb_listener_uses_kafka_topics(): + listener = WaiverDBListener(config_obj=KafkaTestingConfig) + assert listener._backend == "kafka" + assert listener.topic == "qa.eng.waiverdb.waiver.new" + assert listener.destination == "qa.eng.greenwave.decision.update" + + +def test_default_backend_is_stomp(): + listener = ResultsDBListener(config_obj=TestingConfig) + assert listener._backend == "stomp" + assert listener.topic.endswith("VirtualTopic.eng.resultsdb.result.new") + + +def test_listen_kafka_starts_consumer(kafka_env): + with ( + patch("greenwave.listeners.kafka.Consumer") as mock_consumer_cls, + patch("greenwave.listeners.kafka.Producer") as mock_producer_cls, + ): + mock_consumer = Mock() + mock_consumer.poll.return_value = None + mock_consumer_cls.return_value = mock_consumer + mock_producer_cls.return_value = Mock() + + listener = ResultsDBListener(config_obj=KafkaTestingConfig) + try: + listener.listen() + mock_consumer.subscribe.assert_called_once_with( + ["qa.eng.resultsdb.result.new"] + ) + consumer_config = mock_consumer_cls.call_args[0][0] + assert consumer_config["group.id"] == "greenwave-resultsdb" + assert consumer_config["sasl.username"] == "alice" + assert consumer_config["enable.auto.commit"] is False + producer_config = mock_producer_cls.call_args[0][0] + assert producer_config["sasl.password"] == "secret" + assert listener._kafka_thread is not None + assert listener._kafka_thread.daemon + listener.listen() + assert mock_consumer.subscribe.call_count == 1 + finally: + listener.disconnect() + + +def test_handle_kafka_message_commits_after_success(): + listener = ResultsDBListener(config_obj=KafkaTestingConfig) + listener._kafka_bus = Mock() + msg = Mock() + msg.offset.return_value = 12 + msg.value.return_value = json.dumps( + { + "outcome": "QUEUED", + "testcase": {"name": "dist.rpmdeplint"}, + "submit_time": "2019-03-25T16:34:41.882620", + "data": {"item": ["nvr-1.0-1"], "type": ["koji_build"]}, + } + ).encode() + + listener._handle_kafka_message(msg) + + listener._kafka_bus.commit.assert_called_once_with(msg) + + +def test_handle_kafka_message_commits_invalid_json(): + listener = ResultsDBListener(config_obj=KafkaTestingConfig) + listener._kafka_bus = Mock() + msg = Mock() + msg.offset.return_value = 12 + msg.value.return_value = b"not-json" + + listener._handle_kafka_message(msg) + + listener._kafka_bus.commit.assert_called_once_with(msg) + + +def test_handle_kafka_message_does_not_commit_on_failure(): + listener = ResultsDBListener(config_obj=KafkaTestingConfig) + listener._kafka_bus = Mock() + msg = Mock() + msg.offset.return_value = 12 + msg.value.return_value = json.dumps({"outcome": "PASSED"}).encode() + + with patch.object(listener, "_consume_message", side_effect=RuntimeError("boom")): + with raises(RuntimeError, match="boom"): + listener._handle_kafka_message(msg) + + listener._kafka_bus.commit.assert_not_called() + + +def test_publish_decision_update_kafka(): + listener = ResultsDBListener(config_obj=KafkaTestingConfig) + listener._kafka_bus = Mock() + decision = { + "subject_type": "koji_build", + "subject_identifier": "nvr-1.0-1", + "product_version": "fedora-rawhide", + "decision_context": "test_context", + "policies_satisfied": True, + "summary": "All required tests passed", + } + + listener._publish_decision_update(decision) + + listener._kafka_bus.publish.assert_called_once() + topic, body, headers = listener._kafka_bus.publish.call_args[0] + assert topic == "qa.eng.greenwave.decision.update" + payload = json.loads(body) + assert payload["topic"] == topic + assert payload["msg"]["subject_identifier"] == "nvr-1.0-1" + assert headers["policies_satisfied"] == "true" + + +def test_publish_decision_update_kafka_error(): + listener = ResultsDBListener(config_obj=KafkaTestingConfig) + listener._kafka_bus = Mock() + listener._kafka_bus.publish.side_effect = KafkaException("send failed") + decision = { + "subject_type": "koji_build", + "subject_identifier": "nvr-1.0-1", + "product_version": "fedora-rawhide", + "decision_context": "test_context", + "policies_satisfied": False, + "summary": "1 of 1 required tests failed", + } + + with raises(KafkaException): + listener._publish_decision_update(decision) + + +def _simulate_successful_produce(mock_producer): + callbacks = [] + + def capture_produce(topic, value=None, headers=None, on_delivery=None): + callbacks.append(on_delivery) + + def trigger_flush(timeout=None): + for cb in callbacks: + if cb: + cb(None, Mock()) + callbacks.clear() + return 0 + + mock_producer.produce.side_effect = capture_produce + mock_producer.flush.side_effect = trigger_flush + + +def test_kafka_bus_publish_success(kafka_env): + with ( + patch("greenwave.listeners.kafka.Consumer"), + patch("greenwave.listeners.kafka.Producer") as mock_producer_cls, + ): + mock_producer = Mock() + mock_producer_cls.return_value = mock_producer + _simulate_successful_produce(mock_producer) + + bus = KafkaBus(_kafka_config(), group_id="greenwave-resultsdb") + bus.publish( + "qa.eng.greenwave.decision.update", + '{"msg": {}}', + {"summary": "ok"}, + ) + + mock_producer.produce.assert_called_once() + args, kwargs = mock_producer.produce.call_args + assert args[0] == "qa.eng.greenwave.decision.update" + assert kwargs["headers"] == [("summary", b"ok")] + mock_producer.flush.assert_called_with(timeout=15.0) + + +def test_kafka_bus_publish_delivery_error(kafka_env): + with ( + patch("greenwave.listeners.kafka.Consumer"), + patch("greenwave.listeners.kafka.Producer") as mock_producer_cls, + ): + mock_producer = Mock() + mock_producer_cls.return_value = mock_producer + callbacks = [] + + def capture_produce(topic, value=None, headers=None, on_delivery=None): + callbacks.append(on_delivery) + + def trigger_flush(timeout=None): + for cb in callbacks: + if cb: + cb(Mock(), None) + callbacks.clear() + return 0 + + mock_producer.produce.side_effect = capture_produce + mock_producer.flush.side_effect = trigger_flush + + bus = KafkaBus(_kafka_config(), group_id="greenwave-resultsdb") + with raises(KafkaException): + bus.publish("qa.eng.greenwave.decision.update", "{}", {}) + + +def test_parse_kafka_config_consumer_not_dict(kafka_env): + config = _kafka_config() + config["KAFKA"]["consumer"] = "nope" + with raises(RuntimeError, match="consumer and producer must be dicts"): + parse_kafka_config(config) + + +def test_kafka_bus_publish_flush_timeout(kafka_env): + with ( + patch("greenwave.listeners.kafka.Consumer"), + patch("greenwave.listeners.kafka.Producer") as mock_producer_cls, + ): + mock_producer = Mock() + mock_producer.flush.return_value = 1 + mock_producer_cls.return_value = mock_producer + + bus = KafkaBus(_kafka_config(), group_id="greenwave-resultsdb") + with raises(KafkaException) as exc_info: + bus.publish("qa.eng.greenwave.decision.update", "{}", {}) + + err = exc_info.value.args[0] + assert isinstance(err, KafkaError) + assert err.code() == KafkaError._MSG_TIMED_OUT + + +def _queued_result_payload(): + return json.dumps( + { + "outcome": "QUEUED", + "testcase": {"name": "dist.rpmdeplint"}, + "submit_time": "2019-03-25T16:34:41.882620", + "data": {"item": ["nvr-1.0-1"], "type": ["koji_build"]}, + } + ).encode() + + +def _kafka_message(error=None, value=None): + msg = Mock() + msg.error.return_value = error + msg.offset.return_value = 12 + msg.value.return_value = value if value is not None else _queued_result_payload() + return msg + + +def test_kafka_bus_commit_and_close(kafka_env): + with ( + patch("greenwave.listeners.kafka.Consumer") as mock_consumer_cls, + patch("greenwave.listeners.kafka.Producer") as mock_producer_cls, + ): + mock_consumer = Mock() + mock_producer = Mock() + mock_consumer_cls.return_value = mock_consumer + mock_producer_cls.return_value = mock_producer + mock_consumer.close.side_effect = RuntimeError("close failed") + mock_producer.flush.side_effect = RuntimeError("flush failed") + + bus = KafkaBus(_kafka_config(), group_id="greenwave-resultsdb") + msg = Mock() + bus.commit(msg) + mock_consumer.commit.assert_called_once_with(message=msg) + bus.close() + + +def test_handle_kafka_message_processed_ok(): + listener = ResultsDBListener(config_obj=KafkaTestingConfig) + listener._kafka_bus = Mock() + msg = _kafka_message() + + with patch.object(listener, "_consume_message", return_value=True): + listener._handle_kafka_message(msg) + + listener._kafka_bus.commit.assert_called_once_with(msg) + + +def test_handle_kafka_message_stopped(): + listener = ResultsDBListener(config_obj=KafkaTestingConfig) + listener._kafka_bus = Mock() + listener.stop = True + + listener._handle_kafka_message(_kafka_message()) + + listener._kafka_bus.commit.assert_not_called() + + +def test_kafka_loop_processes_message(): + listener = ResultsDBListener(config_obj=KafkaTestingConfig) + listener._kafka_bus = Mock() + msg = _kafka_message() + calls = {"n": 0} + + def poll(_timeout=1.0): + calls["n"] += 1 + if calls["n"] == 1: + return msg + listener.stop = True + return None + + listener._kafka_bus.poll.side_effect = poll + listener._kafka_loop() + listener._kafka_bus.commit.assert_called_once_with(msg) + + +def test_kafka_loop_ignores_partition_eof(): + listener = ResultsDBListener(config_obj=KafkaTestingConfig) + listener._kafka_bus = Mock() + error = Mock() + error.code.return_value = KafkaError._PARTITION_EOF + + def poll(_timeout=1.0): + listener.stop = True + return _kafka_message(error=error) + + listener._kafka_bus.poll.side_effect = poll + listener._kafka_loop() + listener._kafka_bus.commit.assert_not_called() + + +def test_kafka_loop_nonfatal_error(): + listener = ResultsDBListener(config_obj=KafkaTestingConfig) + listener._kafka_bus = Mock() + error = Mock() + error.code.return_value = KafkaError._TRANSPORT + error.fatal.return_value = False + + def poll(_timeout=1.0): + listener.stop = True + return _kafka_message(error=error) + + listener._kafka_bus.poll.side_effect = poll + with patch.object(listener, "_terminate") as terminate: + listener._kafka_loop() + terminate.assert_not_called() + listener._kafka_bus.commit.assert_not_called() + + +def test_kafka_loop_fatal_error(): + listener = ResultsDBListener(config_obj=KafkaTestingConfig) + listener._kafka_bus = Mock() + error = Mock() + error.code.return_value = KafkaError._FATAL + error.fatal.return_value = True + + listener._kafka_bus.poll.return_value = _kafka_message(error=error) + with patch.object(listener, "_terminate") as terminate: + listener._kafka_loop() + terminate.assert_called_once() + + +def test_kafka_loop_exception_terminates(): + listener = ResultsDBListener(config_obj=KafkaTestingConfig) + listener._kafka_bus = Mock() + listener._kafka_bus.poll.side_effect = RuntimeError("broker down") + with patch.object(listener, "_terminate") as terminate: + listener._kafka_loop() + terminate.assert_called_once() + + +def test_disconnect_kafka_closes_bus(): + listener = ResultsDBListener(config_obj=KafkaTestingConfig) + listener._kafka_bus = Mock() + listener._kafka_thread = Mock() + listener._kafka_thread.is_alive.return_value = True + listener.disconnect() + listener._kafka_bus.close.assert_called_once() + listener._kafka_thread.join.assert_called_once_with(timeout=5) + + +def test_listen_kafka_sigterm_disconnects(kafka_env): + with ( + patch("greenwave.listeners.kafka.Consumer") as mock_consumer_cls, + patch("greenwave.listeners.kafka.Producer"), + patch("greenwave.listeners.base.signal.signal") as mock_signal, + ): + mock_consumer = Mock() + mock_consumer.poll.return_value = None + mock_consumer_cls.return_value = mock_consumer + listener = ResultsDBListener(config_obj=KafkaTestingConfig) + try: + listener.listen() + handler = mock_signal.call_args[0][1] + with patch.object(listener, "disconnect") as disconnect: + handler(15, None) + disconnect.assert_called_once() + finally: + listener.stop = True + listener.disconnect() diff --git a/pyproject.toml b/pyproject.toml index 0f7293e2..9ae5c9bd 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -27,6 +27,7 @@ dependencies = [ "python-dateutil>=2.8.2,<3.0", "fedora-messaging>=3.4.1,<4.0", "stomp.py>=9,<9.1", + "confluent-kafka>=2.14.0,<3.0.0", "statsd>=4.0.1,<5.0", "pymemcache>=4.0.0,<5.0", "defusedxml>=0.7.1,<1.0", diff --git a/uv.lock b/uv.lock index 9da8150b..94b27626 100644 --- a/uv.lock +++ b/uv.lock @@ -234,6 +234,29 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/d1/d6/3965ed04c63042e047cb6a3e6ed1a63a35087b6a609aa3a15ed8ac56c221/colorama-0.4.6-py2.py3-none-any.whl", hash = "sha256:4f1d9991f5acc0ca119f9d443620b77f9d6b33703e51011c16baf57afb285fc6", size = 25335, upload-time = "2022-10-25T02:36:20.889Z" }, ] +[[package]] +name = "confluent-kafka" +version = "2.15.1" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/44/fd/e8204b211ce0f32d4d03f9c3423b1244c2dc7fae45babbf98979c00a9e11/confluent_kafka-2.15.1.tar.gz", hash = "sha256:99d1223020e27854e75981333983c5120d64762268f8afcc805fbd62e7de49aa", size = 325065, upload-time = "2026-09-10T15:27:59.367Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/c4/e6/79cc2bb61cf05f00e4b367e2aa4cb940e5eddbfa4adeb4177115791f4a7d/confluent_kafka-2.15.1-cp312-cp312-macosx_10_9_x86_64.whl", hash = "sha256:79749b23a3ad751546f9de5bae3651116cbd0df449afc86d122c80359d6453ca", size = 4381400, upload-time = "2026-09-10T15:27:18.564Z" }, + { url = "https://files.pythonhosted.org/packages/6d/6d/67e59d6807d472ad5f2ac831e735da5ffd321da98adf64cef4cc0bf5f6ed/confluent_kafka-2.15.1-cp312-cp312-macosx_11_0_arm64.whl", hash = "sha256:64fbcb0c6ee95d0baee2b238d177595c7baa8ccd74df96bea086103dfbfd254d", size = 4350483, upload-time = "2026-09-10T15:27:20.019Z" }, + { url = "https://files.pythonhosted.org/packages/25/08/8c404bb5849f98e8468b24f1f336360cc37802bb058ff750255621fea484/confluent_kafka-2.15.1-cp312-cp312-manylinux_2_28_aarch64.whl", hash = "sha256:2d2d876c83214f2651d92c557d5057be1356dd3fbb31f26d50863f610655e2fe", size = 4994340, upload-time = "2026-09-10T15:27:21.431Z" }, + { url = "https://files.pythonhosted.org/packages/d7/a3/31112ce258d74be4636e6d1f73717587e19fd311c70477ae807136f3a48c/confluent_kafka-2.15.1-cp312-cp312-manylinux_2_28_x86_64.whl", hash = "sha256:326e237f6b25587d7dc96854e3a1315979843bacbe48d76a50412f12e49328e1", size = 4806262, upload-time = "2026-09-10T15:27:22.822Z" }, + { url = "https://files.pythonhosted.org/packages/3c/b8/d07ef7c72f55c01d13de4873e86e49eb6ebd8989cb104cb982fa365d389b/confluent_kafka-2.15.1-cp312-cp312-win_amd64.whl", hash = "sha256:b6c73dcb91d5f84ed0d24b37817944f44b5d661f02ec8cc7a3c6d497c03bdc7c", size = 4565312, upload-time = "2026-09-10T15:27:24.321Z" }, + { url = "https://files.pythonhosted.org/packages/fa/c4/1e5222ab41a640c03442a112e55f3d64c1353b0bc9775f6234811e989a80/confluent_kafka-2.15.1-cp313-cp313-macosx_13_0_arm64.whl", hash = "sha256:22e848a79e7e6c208d6e227130139203462ad34c9d27e87c861cff03d3402953", size = 4355493, upload-time = "2026-09-10T15:27:26.209Z" }, + { url = "https://files.pythonhosted.org/packages/6e/5b/af7502dae0ca17f391e3c852222eb849f86e16d30607caa41e566280d598/confluent_kafka-2.15.1-cp313-cp313-macosx_13_0_x86_64.whl", hash = "sha256:8f9ffba2a799d63ed6b3c76eb497307a84fe4509b129f2d50e4e86a38925577e", size = 4385139, upload-time = "2026-09-10T15:27:27.596Z" }, + { url = "https://files.pythonhosted.org/packages/77/4d/9ee22367647259bc61ec6e2466f6cd9bb4a9d1dca36de389af6dafabcb49/confluent_kafka-2.15.1-cp313-cp313-manylinux_2_28_aarch64.whl", hash = "sha256:3d0a053129eda18e15d22a25d62888401b0843b88976ac7356c8d80150860248", size = 4994787, upload-time = "2026-09-10T15:27:29.305Z" }, + { url = "https://files.pythonhosted.org/packages/2e/43/102114edc44314b669904ce51f1473137b2830bf4862dfc6d2b1ad3fcf91/confluent_kafka-2.15.1-cp313-cp313-manylinux_2_28_x86_64.whl", hash = "sha256:4d81bff7c1f169ee674423a0f493a4cde2452d3c352e25d419023fc6133e8316", size = 4806674, upload-time = "2026-09-10T15:27:31.895Z" }, + { url = "https://files.pythonhosted.org/packages/fe/6e/f4ba9fbf64044af8e622ebf03da28696390cacb76d778f0542b7b77424db/confluent_kafka-2.15.1-cp313-cp313-win_amd64.whl", hash = "sha256:894e028a3eac45658ee4c769375b38b0792b2e2a9a737174d320b498d7f35b8e", size = 4624973, upload-time = "2026-09-10T15:27:33.707Z" }, + { url = "https://files.pythonhosted.org/packages/97/1d/7ecedad4c2b67efb795071a0baf84e0352f8e1843761912ae11859df3b93/confluent_kafka-2.15.1-cp314-cp314-macosx_13_0_arm64.whl", hash = "sha256:8e872fd5615ed4d074ddfc0f5c6cb54d4687599b2928219c50165e8738a624e5", size = 4355347, upload-time = "2026-09-10T15:27:35.307Z" }, + { url = "https://files.pythonhosted.org/packages/1f/25/a87e7bd71912435a857a9a3c240290c383f4f8808d50c272afe881907af4/confluent_kafka-2.15.1-cp314-cp314-macosx_13_0_x86_64.whl", hash = "sha256:8c24859213dd3e429ffb24fff8d80d6468068622034f7f32b23c2f8e3d9b728b", size = 4384831, upload-time = "2026-09-10T15:27:37.198Z" }, + { url = "https://files.pythonhosted.org/packages/ac/fb/51e2d0de3761b0b28031ac0b2d0e537e756fe8b6bdc4a9b8981a23c1c2ff/confluent_kafka-2.15.1-cp314-cp314-manylinux_2_28_aarch64.whl", hash = "sha256:9120aeabde16fb758cd0742baee54b8456712ded20f54e18cd5e910629e2ff31", size = 4994508, upload-time = "2026-09-10T15:27:38.805Z" }, + { url = "https://files.pythonhosted.org/packages/86/44/97231c660bade51737044eb1a0f63f5ef4c1d98077a7ef77c32c9c91be9c/confluent_kafka-2.15.1-cp314-cp314-manylinux_2_28_x86_64.whl", hash = "sha256:f1f9877521804e79d5a9a9147dc8cf0f804fc3ebced47d6170024b637197d461", size = 4806334, upload-time = "2026-09-10T15:27:40.672Z" }, + { url = "https://files.pythonhosted.org/packages/b0/fd/db44400d91a0899292a4ea582f44699d41fe90218201b984efd6fb03a3ae/confluent_kafka-2.15.1-cp314-cp314-win_amd64.whl", hash = "sha256:2122c014e61b0a7882632ec39f58f2ee32577e5a069bac28757bf0dc6eaceec0", size = 4759970, upload-time = "2026-09-10T15:27:42.573Z" }, +] + [[package]] name = "constantly" version = "23.10.4" @@ -518,6 +541,7 @@ name = "greenwave" version = "2.3.0" source = { editable = "." } dependencies = [ + { name = "confluent-kafka" }, { name = "defusedxml" }, { name = "dogpile-cache" }, { name = "fedora-messaging" }, @@ -553,6 +577,7 @@ test = [ [package.metadata] requires-dist = [ + { name = "confluent-kafka", specifier = ">=2.14.0,<3.0.0" }, { name = "defusedxml", specifier = ">=0.7.1,<1.0" }, { name = "dogpile-cache", specifier = ">=1.3.3,<2.0" }, { name = "fedora-messaging", specifier = ">=3.4.1,<4.0" },