diff --git a/Dockerfile b/Dockerfile index 14a1a9e..b129853 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,16 +1,30 @@ +# syntax=docker/dockerfile:1.4 + from python:3.11-slim +RUN apt-get update && apt-get install -y git openssh-client && rm -rf /var/lib/apt/lists/* + # Set the working directory WORKDIR /app # Copy the requirements file into the container COPY requirements.txt . +COPY __init__.py . # Copy code into the container -COPY ./ingestor . +COPY ./ingestor ./ingestor # Install the required packages -RUN pip install --no-cache-dir -r requirements.txt +# Add github to known hosts +# This is needed for SSH to work +# The SSH key will NOT remain in the image +# IMPORTANT: this block requires BuildKit +# and the --ssh flag during docker build +RUN --mount=type=ssh \ + mkdir -p ~/.ssh && \ + ssh-keyscan github.com >> ~/.ssh/known_hosts && \ + pip install --no-cache-dir -r requirements.txt + # Run the application -CMD ["python", "ingestor.py"] \ No newline at end of file +CMD ["python", "-m", "ingestor.ingestor"] \ No newline at end of file diff --git a/docker-compose.yaml b/docker-compose.yaml index cc53895..ce7fed1 100644 --- a/docker-compose.yaml +++ b/docker-compose.yaml @@ -1,6 +1,21 @@ version: '3.8' services: + ingestor: + restart: always + build: + context: . + environment: + HOSTNAME: ingestor + container_name: ingestor + depends_on: + - kafka + - redis + networks: + - kafka-net + env_file: + - .env + zookeeper: image: confluentinc/cp-zookeeper:latest container_name: zookeeper @@ -20,7 +35,7 @@ services: environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 - KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:29092,PLAINTEXT_HOST://localhost:9092 + KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:29092,PLAINTEXT_HOST://localhost:9092, KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 diff --git a/ingestor/ingestor.py b/ingestor/ingestor.py index 1510b41..766e615 100644 --- a/ingestor/ingestor.py +++ b/ingestor/ingestor.py @@ -9,12 +9,13 @@ def main(): kafka_servers = getenv("KAFKA_SERVERS", "localhost:9092") redis_host = getenv("REDIS_HOST", "localhost") redis_port = int(getenv("REDIS_PORT", 6379)) - lease_ttl = int(getenv("LEASE_TTL", 20)) - heartbeat_ttl = int(getenv("HEARTBEAT_TTL", 10)) + lease_ttl = int(getenv("LEASE_TTL", 10)) + heartbeat_ttl = int(getenv("HEARTBEAT_TTL", 20)) pod_id = getenv("HOSTNAME", "localhost") - number_of_ingestors = int(getenv("REPLICA_COUNT", 1)) poll_interval = int(getenv("POLL_INTERVAL", 5)) + kafka_servers = kafka_servers.split(",") + logger = getLogger(__name__) logger.setLevel(getenv("LOG_LEVEL", "INFO")) handler = StreamHandler() @@ -27,9 +28,12 @@ def main(): ingestor_manager = IngestorManager( kafka_servers, redis_host, redis_port, lease_ttl, heartbeat_ttl, pod_id, - number_of_ingestors, poll_interval, logger + poll_interval, logger ) + # Declare ingestor ative + ingestor_manager.declare_active() + # Get slot lease acquired = ingestor_manager.get_slot_leases() logger.info(f"Acquired slots: {acquired}") @@ -38,18 +42,33 @@ def main(): logger.warning("No slots available") else: - # Declare ingestor ative - ingestor_manager.declare_active() # Subscribe to acquired slots ingestor_manager.update_opc_servers() ingestor_manager.subscribe_to_tags(acquired) while True: + # Declare ingestor as active + ingestor_manager.declare_active() + logger.info("Polling for slot updates...") # Get active ingestors ingestors = ingestor_manager.get_active_ingestors() + number_of_slots = ingestor_manager.get_number_of_slots() - ingestor_diff = number_of_ingestors - len(ingestors) + if not ingestor_manager.managed_tags and number_of_slots > 0: + # This ingestor is active and has no slots, so we need to try to + # acquire a slot lease + + logger.info("No slots acquired, trying to acquire a slot lease") + acquired = ingestor_manager.get_slot_leases(1) + if not acquired: + logger.info("No slots acquired") + else: + # Subscribe to acquired slots + ingestor_manager.update_opc_servers() + ingestor_manager.subscribe_to_tags(acquired) + + ingestor_diff = number_of_slots - len(ingestors) slot_diff = len(ingestor_manager.managed_tags) - 1 if ingestor_diff > 0: @@ -74,16 +93,19 @@ def main(): ingestor_manager.drop_slot_leases(overleases) - # Update opc servers - ingestor_manager.update_slot_config() + logger.info( + f"Active ingestors: {ingestors}, " + f"Number of slots: {number_of_slots}, " + f"Managed tags: {ingestor_manager.managed_tags}" + f"Managed servers: {ingestor_manager.opc_managers}" + ) if not ingestor_manager.managed_tags: - logger.info("No managed tags found") - sleep(poll_interval) - continue + # No slots acquired + logger.info("No slots acquired in this loop") - # Declare ingestor as active - ingestor_manager.declare_active() + # Update opc servers + ingestor_manager.update_slot_config() # Sleep for poll interval sleep(poll_interval) diff --git a/ingestor/managers/data_manager.py b/ingestor/managers/data_manager.py index e04b480..1833470 100644 --- a/ingestor/managers/data_manager.py +++ b/ingestor/managers/data_manager.py @@ -1,27 +1,47 @@ import json from logging import Logger +from time import sleep from kafka import KafkaProducer +from kafka.errors import NoBrokersAvailable class DataManager(): def __init__(self, kafka_servers: str, logger: Logger) -> None: - self.kafka_producer = KafkaProducer( - bootstrap_servers=kafka_servers, - value_serializer=lambda v: json.dumps(v).encode( - 'utf-8'), # Serialize JSON messages - key_serializer=lambda k: str(k).encode('utf-8') if k else None, - ) + for i in range(0, 3): + logger.info( + f"Trying ({i}) to initializing DataManager with Kafka servers: {kafka_servers}") + try: + self.kafka_producer = KafkaProducer( + bootstrap_servers=kafka_servers, + value_serializer=lambda v: json.dumps(v).encode( + 'utf-8'), # Serialize JSON messages + key_serializer=lambda k: str( + k).encode('utf-8') if k else None, + ) + break + except NoBrokersAvailable: + logger.error( + f"Kafka servers {kafka_servers} are not available. Retrying...") + sleep(5) + else: + logger.error( + f"Failed to connect to Kafka servers {kafka_servers} after 3 attempts.") + raise NoBrokersAvailable( + f"Failed to connect to Kafka servers {kafka_servers} after 3 attempts.") + + logger.info( + f"DataManager initialized with Kafka servers: {kafka_servers}") self.logger = logger def __del__(self): """Destructor to close the producer connection.""" - self.logger.info("Closing Kafka producer...") + print("Closing Kafka producer...") self.kafka_producer.flush() self.kafka_producer.close() def delivery_report(self, msg: str): """Callback for delivery reports from Kafka.""" - self.logger.info( + self.logger.debug( f"Record successfully produced to {msg.topic} [{msg.partition}] at offset {msg.offset}") def delivery_error(self, err: str): @@ -45,7 +65,7 @@ class DataManager(): try: - self.logger.info( + self.logger.debug( f"Publishing message to topic {topic}: {data}") self.kafka_producer.send( topic=topic, value=data).add_callback( diff --git a/ingestor/managers/ingestor_manager.py b/ingestor/managers/ingestor_manager.py index 8579a96..8adf6d1 100644 --- a/ingestor/managers/ingestor_manager.py +++ b/ingestor/managers/ingestor_manager.py @@ -11,7 +11,6 @@ class IngestorManager(): def __init__(self, kafka_servers: str, redis_host: str, redis_port: int, lease_ttl: int, heartbeat_ttl: int, pod_id: str, - number_of_ingestors: int, poll_interval: int, logger: Logger): self.data_manager = DataManager(kafka_servers, logger) @@ -19,7 +18,7 @@ class IngestorManager(): self.resource_manager = ResourceManager( redis_host, redis_port, lease_ttl, heartbeat_ttl, pod_id ) - self.number_of_ingestors = number_of_ingestors + self.number_of_slots = 0 self.poll_interval = poll_interval self.logger = logger self.managed_tags = {} @@ -98,9 +97,14 @@ class IngestorManager(): ingestors = self.resource_manager.get_all_ingestors() return ingestors if ingestors else [] + def get_number_of_slots(self) -> int: + slots = self.resource_manager.get_all_slots() + self.number_of_slots = len(slots) if slots else 0 + return self.number_of_slots + def get_slot_leases(self, max_slots: int = 1) -> Dict: acquired = {} - for i in range(1, self.number_of_ingestors + 1): + for i in range(1, self.number_of_slots + 1): if self.resource_manager.lease_tag(str(i)): self.logger.info(f"Leased slot {i}") slots = self.resource_manager.get_tag_slot(str(i)) diff --git a/ingestor/managers/opc_manager.py b/ingestor/managers/opc_manager.py index ea1f81c..ea9c95f 100644 --- a/ingestor/managers/opc_manager.py +++ b/ingestor/managers/opc_manager.py @@ -27,6 +27,10 @@ class OpcManager(): self.subscriptions = {} self.data_manager = data_manager + def __str__(self): + return f"OpcManager(name={self.name}, url={self.url}, server_uri={self.server_uri})\n" \ + f"nodes={self.nodes}, subscriptions={self.subscriptions}" + def initialize_from_config(self, server_config: dict, data_manager: DataManager, logger: Logger): """ Initializes the OpcManager instance using a server configuration dictionary. diff --git a/ingestor/managers/resource_manager.py b/ingestor/managers/resource_manager.py index 7572d08..2717f3f 100644 --- a/ingestor/managers/resource_manager.py +++ b/ingestor/managers/resource_manager.py @@ -106,3 +106,14 @@ class ResourceManager: """ return self.redis.keys("heartbeat:ingestor:*") + + def get_all_slots(self) -> List[str]: + """ + Retrieves the number of slots available in Redis. + This method counts the number of keys in Redis that match the pattern for OPC tag leases + and returns the count. + Returns: + int: The number of slots available. + """ + + return self.redis.keys("slot:opc_tags:*") diff --git a/requirements.txt b/requirements.txt index ff96ac6..e564b8e 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1,3 +1,3 @@ asyncua==1.1.5 redis -git+https://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@github.com/Aignosi/sientia-dataops-library.git \ No newline at end of file +git+ssh://git@github.com/Aignosi/sientia-dataops-library.git \ No newline at end of file diff --git a/simulator/redis-feeder.py b/simulator/redis-feeder.py index 21c8833..06b9884 100644 --- a/simulator/redis-feeder.py +++ b/simulator/redis-feeder.py @@ -3,8 +3,8 @@ import json import os # Redis connection settings -redis_host = os.getenv("REDIS_HOST", "localhost") -redis_port = int(os.getenv("REDIS_PORT", 6379)) +redis_host = "localhost" +redis_port = 6379 # Connect to Redis r = redis.Redis(host=redis_host, port=redis_port, decode_responses=True) @@ -24,7 +24,7 @@ new_data = { "slot:opc_tags:1": { "server1": { "name": "server1", - "url": "opc.tcp://localhost:4840", + "url": "opc.tcp://simulator:4840", "server_uri": "http://opcua-server.simulator", "tags": { 'ns=2;i=2': { @@ -44,6 +44,30 @@ new_data = { }, } } + }, + "slot:opc_tags:2": { + "server2": { + "name": "server2", + "url": "opc.tcp://simulator:4840", + "server_uri": "http://opcua-server.simulator", + "tags": { + 'ns=2;i=2': { + 'tag_name': 'Counter', + 'frequency': 1000, + 'topics': ['opcua2', 'counter'], + }, + 'ns=2;i=3': { + 'tag_name': 'Rollout', + 'frequency': 1000, + "topics": ['opcua2', 'rollout'], + }, + 'ns=2;i=4': { + 'tag_name': 'Square', + 'frequency': 1000, + "topics": ['opcua2'], + }, + } + } } } diff --git a/tests/__init__.py b/tests/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/tests/functional/__init__.py b/tests/functional/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/tests/functional/conftest.py b/tests/functional/conftest.py new file mode 100644 index 0000000..2a05ff6 --- /dev/null +++ b/tests/functional/conftest.py @@ -0,0 +1,47 @@ +import subprocess +from time import sleep +from typing import Generator +import uuid +from kafka import KafkaConsumer +import pytest +from redis import Redis + + +@pytest.fixture(scope="session", autouse=True) +def docker_compose(): + """Sobe os containers antes dos testes e derruba depois.""" + print("\n🚀 Subindo Docker Compose...") + subprocess.run(["docker", "compose", "up", "-d"], check=True) + + print("⏳ Aguardando containers ficarem prontos...") + sleep(15) # ajuste conforme necessário + + yield # os testes rodam aqui + + print("\n🧹 Derrubando Docker Compose...") + subprocess.run(["docker", "compose", "down"], check=True) + + +@pytest.fixture() +def redis_client(): + redis = Redis(host="localhost", port=6379, decode_responses=True) + redis.flushdb() + + yield redis + + # Limpa o banco de dados após os testes + redis.flushdb() + redis.close() + + +def kafka_searcher(topic) -> Generator[KafkaConsumer, None, None]: + consumer = KafkaConsumer( + topic, + bootstrap_servers="localhost:9092", + group_id=f"test-group-{uuid.uuid4()}", + auto_offset_reset="earliest", # Começa a consumir apenas mensagens novas + enable_auto_commit=True, + ) + + yield consumer + consumer.close() diff --git a/tests/functional/test_single_node.py b/tests/functional/test_single_node.py new file mode 100644 index 0000000..24ee1aa --- /dev/null +++ b/tests/functional/test_single_node.py @@ -0,0 +1,118 @@ +import json +import subprocess +from time import sleep + +from tests.functional.conftest import kafka_searcher + + +new_data = { + "slot:opc_tags:1": { + "server1": { + "name": "server1", + "url": "opc.tcp://simulator:4840", + "server_uri": "http://opcua-server.simulator", + "tags": { + 'ns=2;i=2': { + 'tag_name': 'Counter', + 'frequency': 1000, + 'topics': [], + }, + 'ns=2;i=3': { + 'tag_name': 'Rollout', + 'frequency': 1000, + "topics": [], + }, + 'ns=2;i=4': { + 'tag_name': 'Square', + 'frequency': 1000, + "topics": [], + }, + } + } + }, + "slot:opc_tags:2": { + "server2": { + "name": "server2", + "url": "opc.tcp://simulator:4840", + "server_uri": "http://opcua-server.simulator", + "tags": { + 'ns=2;i=2': { + 'tag_name': 'Counter', + 'frequency': 1000, + 'topics': [], + }, + 'ns=2;i=3': { + 'tag_name': 'Rollout', + 'frequency': 1000, + "topics": [], + }, + 'ns=2;i=4': { + 'tag_name': 'Square', + 'frequency': 1000, + "topics": [], + }, + } + } + } +} + + +def test_simple(redis_client): + new_data['slot:opc_tags:1']['server1']['tags']['ns=2;i=2']['topics'] = [ + 'test_topic_1'] + redis_client.set("slot:opc_tags:1", + json.dumps(new_data['slot:opc_tags:1'])) + + sleep(20) # Espera o Ingestor processar os dados + + # Check if lease is in Redis + assert redis_client.get("lease:opc_tags:1") == 'ingestor' + assert redis_client.get("heartbeat:ingestor:ingestor") == '1' + + # Check if data is in Kafka + + kafka = next(kafka_searcher('test_topic_1')) + sleep(1) + messages = kafka.poll(timeout_ms=10000) + + assert messages, "Expected messages in Kafka, but got none." + + +def test_simple_double_slot(redis_client): + + new_data['slot:opc_tags:1']['server1']['tags']['ns=2;i=2']['topics'] = [ + 'test_topic_double_slot1'] + redis_client.set("slot:opc_tags:1", + json.dumps(new_data['slot:opc_tags:1'])) + + sleep(20) # Espera o Ingestor processar os dados + + assert redis_client.get("lease:opc_tags:1") == 'ingestor' + assert redis_client.get("heartbeat:ingestor:ingestor") == '1' + + # Check if data is in Kafka + kafka1 = next(kafka_searcher('test_topic_double_slot1')) + messages = kafka1.poll(timeout_ms=10000) + + assert messages, "Expected messages in test_topic_double_slot1, but got none." + + new_data['slot:opc_tags:2']['server2']['tags']['ns=2;i=2']['topics'] = [ + 'test_topic_double_slot2'] + redis_client.set("slot:opc_tags:2", + json.dumps(new_data['slot:opc_tags:2'])) + + sleep(20) # Espera o Ingestor processar os dados + + # Check if lease is in Redis + assert redis_client.get("lease:opc_tags:2") == 'ingestor' + assert redis_client.get("lease:opc_tags:1") == 'ingestor' + assert redis_client.get("heartbeat:ingestor:ingestor") == '1' + + # Check if data is in Kafka + kafka2 = next(kafka_searcher('test_topic_double_slot2')) + messages = kafka2.poll(timeout_ms=10000) + + assert messages, "Expected messages in test_topic_double_slot2, but got none." + + messages = kafka1.poll(timeout_ms=10000) + assert messages, "Expected messages in test_topic_double_slot1, but got none." diff --git a/tests/managers/test_data_manager.py b/tests/unit/managers/test_data_manager.py similarity index 100% rename from tests/managers/test_data_manager.py rename to tests/unit/managers/test_data_manager.py diff --git a/tests/managers/test_ingestor_manager.py b/tests/unit/managers/test_ingestor_manager.py similarity index 100% rename from tests/managers/test_ingestor_manager.py rename to tests/unit/managers/test_ingestor_manager.py diff --git a/tests/managers/test_opc_manager.py b/tests/unit/managers/test_opc_manager.py similarity index 100% rename from tests/managers/test_opc_manager.py rename to tests/unit/managers/test_opc_manager.py diff --git a/tests/managers/test_resource_manager.py b/tests/unit/managers/test_resource_manager.py similarity index 100% rename from tests/managers/test_resource_manager.py rename to tests/unit/managers/test_resource_manager.py