Code import - branch 1.4.2

This commit is contained in:
2026-08-05 13:53:41 +00:00
commit 6282d0fc6c
39 changed files with 6234 additions and 0 deletions

View File

View File

@@ -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()

View File

@@ -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."