diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..14a1a9e --- /dev/null +++ b/Dockerfile @@ -0,0 +1,16 @@ +from python:3.11-slim + +# Set the working directory +WORKDIR /app + +# Copy the requirements file into the container +COPY requirements.txt . + +# Copy code into the container +COPY ./ingestor . + +# Install the required packages +RUN pip install --no-cache-dir -r requirements.txt + +# Run the application +CMD ["python", "ingestor.py"] \ No newline at end of file diff --git a/README.md b/README.md index 800a28c..2a67e0d 100644 --- a/README.md +++ b/README.md @@ -1,2 +1,40 @@ # sientia-dataops-opc-gateway OPC gateway to manage Scouter pipelines + +## Local tests +### Generate your ssh key to Docker +''' +ssh-keygen -t ed25519 -C "docker-access" -f ~/.ssh/id_ed25519_docker +''' +Add the public key to yout Git SSH keys + +### Enable Docker BuildKit +''' +export DOCKER_BUILDKIT=1 +''' +or make it permanent: +''' +echo '{ "features": { "buildkit": true } }' | sudo tee /etc/docker/daemon.json +sudo systemctl restart docker +''' + +### Run docker compose +''' +docker compose build --ssh default=$HOME/.ssh/id_ed25519_docker +docker compose up -d +''' + +### Populate redis server +Create venv with python3.11 +''' +python3.11 -m venv venv +source ./venv/bin/activate +''' +Install requirements +''' +pip install -r requirements.txt +''' +Run feeder +''' +python simulator/redis-feeder.py +''' \ No newline at end of file diff --git a/__init__.py b/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/docker-compose.yaml b/docker-compose.yaml new file mode 100644 index 0000000..cc53895 --- /dev/null +++ b/docker-compose.yaml @@ -0,0 +1,88 @@ +version: '3.8' + +services: + zookeeper: + image: confluentinc/cp-zookeeper:latest + container_name: zookeeper + environment: + ZOOKEEPER_CLIENT_PORT: 2181 + ZOOKEEPER_TICK_TIME: 2000 + networks: + - kafka-net + env_file: + - .env + + kafka: + image: confluentinc/cp-kafka:latest + container_name: kafka + ports: + - "9092:9092" + environment: + KAFKA_BROKER_ID: 1 + KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 + 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 + depends_on: + - zookeeper + networks: + - kafka-net + env_file: + - .env + + redis: + image: redis:latest + container_name: redis + ports: + - "6379:6379" + networks: + - kafka-net + env_file: + - .env + + redis-commander: + image: rediscommander/redis-commander:latest + container_name: redis-commander + environment: + REDIS_HOSTS: local:redis:6379 + ports: + - "8081:8081" + depends_on: + - redis + networks: + - kafka-net + + + kafdrop: + image: obsidiandynamics/kafdrop:latest + networks: + - kafka-net + depends_on: + - kafka + ports: + - 19000:9000 + environment: + KAFKA_BROKERCONNECT: kafka:29092 + + simulator: + build: + context: . + dockerfile: simulator/Dockerfile + args: + GIT_REPO: ${SIMULATOR_GIT_REPO} + GIT_BRANCH: ${SIMULATOR_GIT_BRANCH} + container_name: simulator + ports: + - "4840:4840" + depends_on: + - kafka + - redis + networks: + - kafka-net + env_file: + - .env + +networks: + kafka-net: + driver: bridge diff --git a/ingestor/ingestor.py b/ingestor/ingestor.py index 798ccb7..1510b41 100644 --- a/ingestor/ingestor.py +++ b/ingestor/ingestor.py @@ -6,14 +6,14 @@ from time import sleep def main(): # Get os parameters - kafka_servers = getenv("KAFKA_SERVERS") - redis_host = getenv("REDIS_HOST") - redis_port = int(getenv("REDIS_PORT")) - lease_ttl = int(getenv("LEASE_TTL")) - heartbeat_ttl = int(getenv("HEARTBEAT_TTL")) - pod_id = getenv("HOSTNAME") - number_of_ingestors = int(getenv("REPLICA_COUNT")) - poll_interval = int(getenv("POLL_INTERVAL")) + 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)) + pod_id = getenv("HOSTNAME", "localhost") + number_of_ingestors = int(getenv("REPLICA_COUNT", 1)) + poll_interval = int(getenv("POLL_INTERVAL", 5)) logger = getLogger(__name__) logger.setLevel(getenv("LOG_LEVEL", "INFO")) @@ -29,22 +29,23 @@ def main(): lease_ttl, heartbeat_ttl, pod_id, number_of_ingestors, 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}") - # Subscribe to acquired slots - ingestor_manager.update_opc_servers() - ingestor_manager.subscribe_to_tags(acquired) + if not acquired: + 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() - - # Update opc servers - ingestor_manager.update_slot_config() - + logger.info("Polling for slot updates...") # Get active ingestors ingestors = ingestor_manager.get_active_ingestors() @@ -53,13 +54,18 @@ def main(): if ingestor_diff > 0: # Some ingestors are innactive, so theres "ingestor_diff" slots available + logger.info(f"Slots available: {ingestor_diff}") # Get slot lease acquired = ingestor_manager.get_slot_leases(ingestor_diff) - # Subscribe to acquired slots - ingestor_manager.update_opc_servers() - ingestor_manager.subscribe_to_tags(acquired) + if not acquired: + logger.info("No slots acquired") + + else: + # Subscribe to acquired slots + ingestor_manager.update_opc_servers() + ingestor_manager.subscribe_to_tags(acquired) elif slot_diff > 0: # Some ingestors are active and without slots, so we need to drop @@ -68,6 +74,17 @@ def main(): ingestor_manager.drop_slot_leases(overleases) + # Update opc servers + ingestor_manager.update_slot_config() + + if not ingestor_manager.managed_tags: + logger.info("No managed tags found") + sleep(poll_interval) + continue + + # Declare ingestor as active + ingestor_manager.declare_active() + # Sleep for poll interval sleep(poll_interval) diff --git a/ingestor/managers/data_manager.py b/ingestor/managers/data_manager.py index 742bca2..e04b480 100644 --- a/ingestor/managers/data_manager.py +++ b/ingestor/managers/data_manager.py @@ -43,9 +43,16 @@ class DataManager(): Exception: If there is an error during message delivery, it will be handled by the `delivery_error` callback. """ - self.kafka_producer.send( - topic=topic, value=data).add_callback( - self.delivery_report).add_errback( - self.delivery_error) + try: - self.kafka_producer.flush() + self.logger.info( + f"Publishing message to topic {topic}: {data}") + self.kafka_producer.send( + topic=topic, value=data).add_callback( + self.delivery_report).add_errback( + self.delivery_error) + + self.kafka_producer.flush(timeout=10) + + except Exception as e: + self.logger.error(f"Failed to publish message: {e}") diff --git a/ingestor/managers/ingestor_manager.py b/ingestor/managers/ingestor_manager.py index 5e7bcab..8579a96 100644 --- a/ingestor/managers/ingestor_manager.py +++ b/ingestor/managers/ingestor_manager.py @@ -1,5 +1,7 @@ from logging import Logger +import traceback from typing import Dict, List +from copy import deepcopy from ingestor.managers.data_manager import DataManager from ingestor.managers.opc_manager import OpcManager from ingestor.managers.resource_manager import ResourceManager @@ -23,37 +25,88 @@ class IngestorManager(): self.managed_tags = {} self.opc_servers = {} + def initialize_opc_from_config(self, server_config: dict, data_manager: DataManager, logger: Logger) -> OpcManager | None: + """ + Initializes an OPC Manager instance using the provided server configuration. + Args: + server_config (dict): A dictionary containing the OPC server configuration. + Expected keys include: + - 'name' (str): The name of the OPC server. + - 'url' (str): The URL of the OPC server. + - 'server_uri' (str): The URI of the OPC server. + - 'cert_path' (str, optional): Path to the client certificate file. + - 'private_key_path' (str, optional): Path to the private key file. + - 'server_cert_path' (str, optional): Path to the server certificate file. + data_manager (DataManager): An instance of the DataManager to handle data operations. + logger (Logger): A logger instance for logging messages. + Returns: + OpcManager | None: An initialized OpcManager instance if successful, + otherwise None if an error occurs during initialization. + """ + + try: + manager = OpcManager( + server_config['name'], server_config['url'], data_manager, logger, server_config['server_uri'], + server_config.get('cert_path'), server_config.get( + 'private_key_path'), server_config.get('server_cert_path') + ) + + manager.config = server_config + manager.connect() + except Exception as e: + logger.error(f"Failed to initialize OpcManager: {e}") + return None + + return manager + def update_opc_servers(self): + registered_servers = [] for slot, slot_config in self.managed_tags.items(): for server, server_config in slot_config.items(): + registered_servers.append(server) server_config = server_config.copy() server_config.pop('tags', None) - if server not in self.opc_servers: - self.opc_managers[server] = OpcManager.initialize_from_config( - server_config, self.data_manager, self.logger - ) - self.opc_managers[server].connect() - elif self.opc_servers[server] != server_config: - del self.opc_managers[server] - self.opc_managers[server] = OpcManager.initialize_from_config( - server_config, self.data_manager, self.logger - ) - self.opc_managers[server].connect() + server_instance = None + if server not in self.opc_managers: - self.opc_servers[server] = server_config + server_instance = self.initialize_opc_from_config( + server_config, self.data_manager, self.logger + ) + + elif self.opc_managers[server].config != server_config: + del self.opc_managers[server] + server_instance = self.initialize_opc_from_config( + server_config, self.data_manager, self.logger + ) + + if server_instance is not None: + self.opc_managers[server] = server_instance + + for server in list(self.opc_managers.keys()): + if server not in registered_servers: + self.logger.warning( + f"Server {server} not found in managed tags. " + f"Desconnecting from server." + ) + self.opc_managers[server].disconnect() + self.opc_managers.pop(server, None) def declare_active(self): self.resource_manager.ingestor_heartbeat() def get_active_ingestors(self) -> List[str]: - self.resource_manager.get_all_ingestors() + ingestors = self.resource_manager.get_all_ingestors() + return ingestors if ingestors else [] def get_slot_leases(self, max_slots: int = 1) -> Dict: acquired = {} for i in range(1, self.number_of_ingestors + 1): if self.resource_manager.lease_tag(str(i)): self.logger.info(f"Leased slot {i}") - acquired[str(i)] = self.resource_manager.get_tag_slot(str(i)) + slots = self.resource_manager.get_tag_slot(str(i)) + if slots is None: + continue + acquired[str(i)] = slots if len(acquired) >= max_slots: self.managed_tags.update(acquired) @@ -66,9 +119,26 @@ class IngestorManager(): self.managed_tags.update(acquired) return acquired - def update_slot_config(self) -> Dict: + def unsubscribe_slot(self, slot: str): + for server in self.managed_tags[slot].keys(): + if server in self.opc_managers: + self.opc_managers[server].unsubscribe(slot) + + def update_slot_config(self): + removed_slots = [] + update = {} for slot, slot_config in self.managed_tags.items(): + self.resource_manager.renew_tag_lease(slot) update = self.resource_manager.get_tag_slot(slot) + if update is None: + self.logger.warning( + f"Slot {slot} configuration not found. " + f"Removing slot from managed tags." + ) + self.unsubscribe_slot(slot) + removed_slots.append(slot) + + continue if update != slot_config: self.logger.info( @@ -77,27 +147,64 @@ class IngestorManager(): ) self.managed_tags[slot] = update + self.unsubscribe_slot(slot) self.update_opc_servers() - self.subscribe_to_tags(update) + self.subscribe_to_tags({slot: update}) - return update + for slot in removed_slots: + del self.managed_tags[slot] + self.update_opc_servers() def drop_slot_leases(self, ids: List[str]) -> None: for id in ids: self.resource_manager.drop_tag_lease(id) - def subscribe_to_tags(self, tags: Dict) -> None: + def subscribe_to_tags(self, tags: Dict) -> None: # NOSONAR + to_remove = [] + self.logger.info(tags) for slot, slot_config in tags.items(): for server, server_config in slot_config.items(): + self.logger.info( + f"Subscribing to tags from {slot}:{server}" + ) + tags_to_sub = server_config.get('tags') if server not in self.opc_managers: self.logger.error( f"Server {server} not found in opc_managers." ) continue if slot not in self.opc_managers[server].subscriptions: - self.opc_managers[server].create_subscription( - slot + try: + self.opc_managers[server].create_subscription( + slot + ) + except Exception as e: + self.logger.error( + f"Failed to create subscription for slot {slot}: {e}" + ) + to_remove.append([slot, server]) + continue + try: + self.logger.info( + tags_to_sub ) - self.opc_managers[server].subscribe( - slot, server_config['tags'], self.poll_interval - ) + self.opc_managers[server].subscribe( + slot, deepcopy(tags_to_sub), self.poll_interval + ) + self.logger.info( + tags_to_sub + ) + except Exception as e: + self.logger.error( + f"Failed to subscribe to tags from {slot}:{server}\n{tags}: {e}" + ) + self.logger.error(traceback.format_exc()) + self.logger.warning( + "Removing subscription from server " + f"{server} for slot {slot}" + ) + self.opc_managers[server].unsubscribe(slot) + to_remove.append([slot, server]) + + for slot, server in to_remove: + self.managed_tags[slot].pop(server, None) diff --git a/ingestor/managers/opc_manager.py b/ingestor/managers/opc_manager.py index e6946cc..ea1f81c 100644 --- a/ingestor/managers/opc_manager.py +++ b/ingestor/managers/opc_manager.py @@ -40,6 +40,7 @@ class OpcManager(): server_config.get('cert_path'), server_config.get( 'private_key_path'), server_config.get('server_cert_path') ) + self.config = server_config def set_security(self): """ @@ -140,6 +141,8 @@ class OpcManager(): raise ValueError( "Subscription not created. Call create_subscription first.") + self.logger.info(f"Subscribing to {subscription}...") + self.logger.info(f"Subscribing to nodes: {nodes}") self.addr_nodes = [self.client.get_node( n) for n in nodes if n not in self.nodes] self.nodes.update(nodes) @@ -156,8 +159,9 @@ class OpcManager(): def unsubscribe(self, subscription: str): if not self.subscriptions.get(subscription): - raise ValueError( - "Subscription not created. Call create_subscription first.") + self.logger.warning( + f"Subscription {subscription} not found. Cannot unsubscribe.") + return self.subscriptions[subscription].delete() del self.subscriptions[subscription] self.logger.info(f"Unsubscribed from {subscription}.") @@ -177,10 +181,14 @@ class OpcManager(): """ self.logger.warning('Disconnecting from OPC server') + if self.client is None: + self.logger.warning("Client already disconnected.") + return try: [self.subscriptions[sub].delete() for sub in self.subscriptions] self.logger.warning("Deleted all subscriptions.") - self.client.disconnect() + del self.client + self.client = None self.logger.warning("Disconnected from OPC UA server.") except Exception as sub_error: self.logger.error(f"Failed to clean up subscription: {sub_error}") @@ -217,7 +225,7 @@ class OpcManager(): data = { 'tag': tag, 'name': self.nodes[str(node)]['tag_name'], - 'timestamp': source_timestamp, + 'timestamp': source_timestamp.strftime('%Y-%m-%d %H:%M:%S'), 'value': value } diff --git a/redis-ui.ipynb b/redis-ui.ipynb new file mode 100644 index 0000000..e69de29 diff --git a/simulator/Dockerfile b/simulator/Dockerfile new file mode 100644 index 0000000..d467676 --- /dev/null +++ b/simulator/Dockerfile @@ -0,0 +1,30 @@ +# syntax=docker/dockerfile:1.4 + +FROM python:3.11-slim + +# Enable use of SSH agent/socket +# This line enables SSH during build +# (don't forget the syntax header above) +RUN apt-get update && apt-get install -y git openssh-client && rm -rf /var/lib/apt/lists/* + +# Use build-time SSH mount for Git clone +# The SSH key will NOT remain in the image +# IMPORTANT: this block requires BuildKit +# and the --ssh flag during docker build + +# SSH config to skip host key check (safe in CI/local dev) +RUN mkdir -p /root/.ssh && echo "StrictHostKeyChecking no" > /root/.ssh/config + +WORKDIR /app + +# Clone using SSH +ARG GIT_REPO +ARG GIT_BRANCH=main + +# Mount SSH key just for this RUN +RUN --mount=type=ssh git clone --branch ${GIT_BRANCH} ${GIT_REPO} . + +# Install requirements if exists +RUN if [ -f requirements.txt ]; then pip install --no-cache-dir -r requirements.txt; fi + +CMD ["python", "server.py"] diff --git a/simulator/redis-feeder.py b/simulator/redis-feeder.py new file mode 100644 index 0000000..21c8833 --- /dev/null +++ b/simulator/redis-feeder.py @@ -0,0 +1,54 @@ +import redis +import json +import os + +# Redis connection settings +redis_host = os.getenv("REDIS_HOST", "localhost") +redis_port = int(os.getenv("REDIS_PORT", 6379)) + +# Connect to Redis +r = redis.Redis(host=redis_host, port=redis_port, decode_responses=True) + +# Define the key pattern to target +pattern = "slot:opc_tags:*" + +# Step 1: Find and delete matching keys +print("🔍 Searching for keys matching:", pattern) +for key in r.scan_iter(match=pattern): + r.delete(key) + print(f"❌ Deleted: {key}") + +# Step 2: Insert new data +# Example new OPC tag data +new_data = { + "slot:opc_tags:1": { + "server1": { + "name": "server1", + "url": "opc.tcp://localhost:4840", + "server_uri": "http://opcua-server.simulator", + "tags": { + 'ns=2;i=2': { + 'tag_name': 'Counter', + 'frequency': 1000, + 'topics': ['opcua', 'counter'], + }, + 'ns=2;i=3': { + 'tag_name': 'Rollout', + 'frequency': 1000, + "topics": ['opcua', 'rollout'], + }, + 'ns=2;i=4': { + 'tag_name': 'Square', + 'frequency': 1000, + "topics": ['opcua'], + }, + } + } + } +} + +for key, val in new_data.items(): + r.set(key, json.dumps(val)) + print(f"✅ Set: {key} -> {val}") + +print("🚀 OPC tag keys replaced successfully.") diff --git a/tests/managers/test_ingestor_manager.py b/tests/managers/test_ingestor_manager.py index 03709f2..cb6ccf7 100644 --- a/tests/managers/test_ingestor_manager.py +++ b/tests/managers/test_ingestor_manager.py @@ -53,11 +53,61 @@ def test___init__(resource_manager_mock, data_manager_mock, opc_manager_mock): assert ingestor.resource_manager == resource_manager_mock.return_value +@patch('ingestor.managers.ingestor_manager.OpcManager') +def test_initialize_opc_from_config(opc_manager, ingestor_manager): + server_config = { + 'name': 'server1', + 'url': 'opc.tcp://localhost:4840', + 'server_uri': 'http://opcua-server.simulator', + 'cert_path': '/path/to/cert', + 'private_key_path': '/path/to/private_key', + 'server_cert_path': '/path/to/server_cert' + } + + opc_manager.return_value = MagicMock() + result = ingestor_manager.initialize_opc_from_config( + server_config, ingestor_manager.data_manager, ingestor_manager.logger) + + opc_manager.assert_called_once_with( + server_config['name'], server_config['url'], ingestor_manager.data_manager, ingestor_manager.logger, + server_config['server_uri'], server_config['cert_path'], server_config['private_key_path'], + server_config['server_cert_path'] + ) + + assert result == opc_manager.return_value + result.connect.assert_called_once() + + +@patch('ingestor.managers.ingestor_manager.OpcManager') +def test_initialize_opc_from_config_exception(opc_manager, ingestor_manager): + server_config = { + 'name': 'server1', + 'url': 'opc.tcp://localhost:4840', + 'server_uri': 'http://opcua-server.simulator', + 'cert_path': '/path/to/cert', + 'private_key_path': '/path/to/private_key', + 'server_cert_path': '/path/to/server_cert' + } + + ingestor_manager.logger.error = MagicMock() + opc_manager.side_effect = Exception("Initialization error") + + result = ingestor_manager.initialize_opc_from_config( + server_config, ingestor_manager.data_manager, ingestor_manager.logger) + + assert result is None + ingestor_manager.logger.error.assert_called_once_with( + "Failed to initialize OpcManager: Initialization error") + + @patch('ingestor.managers.ingestor_manager.OpcManager') def test_update_opc_servers(opc_manager, ingestor_manager): - manager1 = MagicMock() - manager2 = MagicMock() - manager3 = MagicMock() + manager1 = MagicMock( + config={"config": "config1"}) + manager2 = MagicMock( + config={"config": "config2"}) + manager3 = MagicMock( + config={"config": "config3"}) def mock_initialize_from_config(config, data_manager, logger): if config == {"config": "config1"}: @@ -66,15 +116,18 @@ def test_update_opc_servers(opc_manager, ingestor_manager): return manager2 elif config == {"config": "config3"}: return manager3 + else: + return None - opc_manager.initialize_from_config = MagicMock( + ingestor_manager.initialize_opc_from_config = MagicMock( side_effect=mock_initialize_from_config ) ingestor_manager.managed_tags = { "slot1": { "server1": {"config": "config1"}, - "server2": {"config": "config2"} + "server2": {"config": "config2"}, + 'server5': {"config": "config5"} }, "slot2": { "server3": {"config": "config3"}, @@ -82,31 +135,33 @@ def test_update_opc_servers(opc_manager, ingestor_manager): } } - ingestor_manager.opc_servers = { - "server3": {"config": "config3"}, - "server2": {"old_config": "old_config2"} - } - - mock = MagicMock() - ingestor_manager.opc_managers['server3'] = MagicMock() + mock = MagicMock( + config={"config": "old_config2"}) + ingestor_manager.opc_managers['server3'] = MagicMock( + config={"config": "config3"}) ingestor_manager.opc_managers['server2'] = mock + ingestor_manager.opc_managers['server4'] = MagicMock() + ingestor_manager.update_opc_servers() assert len(ingestor_manager.opc_managers) == 3 - assert len(ingestor_manager.opc_servers) == 3 - opc_manager.initialize_from_config.assert_any_call( + ingestor_manager.initialize_opc_from_config.assert_any_call( {"config": "config1"}, ingestor_manager.data_manager, ingestor_manager.logger) - opc_manager.initialize_from_config.assert_any_call( + ingestor_manager.initialize_opc_from_config.assert_any_call( {"config": "config2"}, ingestor_manager.data_manager, ingestor_manager.logger) + ingestor_manager.initialize_opc_from_config.assert_any_call( + {"config": "config5"}, ingestor_manager.data_manager, ingestor_manager.logger) + assert ingestor_manager.initialize_opc_from_config.call_count == 3 - ingestor_manager.opc_managers['server1'].connect.assert_called_once() - ingestor_manager.opc_managers['server2'].connect.assert_called_once() - ingestor_manager.opc_managers['server3'].connect.assert_not_called() - - assert ingestor_manager.opc_servers['server1'] == {"config": "config1"} - assert ingestor_manager.opc_servers['server2'] == {"config": "config2"} - assert ingestor_manager.opc_servers['server3'] == {"config": "config3"} + assert ingestor_manager.opc_managers['server1'].config == { + "config": "config1"} + assert ingestor_manager.opc_managers['server2'].config == { + "config": "config2"} + assert ingestor_manager.opc_managers['server3'].config == { + "config": "config3"} + assert 'server4' not in ingestor_manager.opc_managers + assert 'server5' not in ingestor_manager.opc_managers assert ingestor_manager.opc_managers['server2'] != mock @@ -123,6 +178,14 @@ def test_get_active_ingestors(ingestor_manager): ingestor_manager.resource_manager.get_all_ingestors.assert_called_once() +def test_get_active_ingestors_empty(ingestor_manager): + ingestor_manager.resource_manager.get_all_ingestors = MagicMock( + return_value=None) + result = ingestor_manager.get_active_ingestors() + assert result == [] + ingestor_manager.resource_manager.get_all_ingestors.assert_called_once() + + def test_get_slot_leases_1_success(ingestor_manager): ingestor_manager.resource_manager.lease_tag = MagicMock( return_value=True) @@ -162,21 +225,49 @@ def test_get_slot_leases_1_failure(ingestor_manager): assert result == {} +def test_unsubscribe_slot(ingestor_manager): + ingestor_manager.managed_tags = { + "slot1": { + "server1": {"tags": "config1"}, + "server2": {"tags": "config2"} + }, + "slot2": { + "server3": {"tags": "config3"}, + "server1": {"tags": "config1"} + } + } + ingestor_manager.opc_managers = { + "server1": MagicMock(), + "server2": MagicMock(), + "server3": MagicMock() + } + ingestor_manager.unsubscribe_slot("slot1") + + ingestor_manager.opc_managers["server1"].unsubscribe.assert_called_once_with( + "slot1") + ingestor_manager.opc_managers["server2"].unsubscribe.assert_called_once_with( + "slot1") + ingestor_manager.opc_managers["server3"].unsubscribe.assert_not_called() + + def test_update_slot_config(ingestor_manager): ingestor_manager.managed_tags = { "slot1": {"config": "old_config"}, - "slot2": {"config": "new_config"} + "slot2": {"config": "new_config"}, + "slot3": {"config": "old_config"} } ingestor_manager.resource_manager.get_tag_slot = MagicMock( side_effect=[ {"config": "updated_config"}, - {"config": "new_config"} + {"config": "new_config"}, + None ] ) ingestor_manager.update_opc_servers = MagicMock() ingestor_manager.subscribe_to_tags = MagicMock() + ingestor_manager.unsubscribe_slot = MagicMock() ingestor_manager.update_slot_config() @@ -184,11 +275,20 @@ def test_update_slot_config(ingestor_manager): "config": "updated_config"} assert ingestor_manager.managed_tags["slot2"] == { "config": "new_config"} + assert "slot3" not in ingestor_manager.managed_tags ingestor_manager.update_opc_servers.assert_called_once() ingestor_manager.subscribe_to_tags.assert_called_once_with( {"config": "updated_config"} ) + ingestor_manager.unsubscribe_slot.assert_any_call("slot3") + ingestor_manager.unsubscribe_slot.assert_any_call("slot1") + assert ingestor_manager.unsubscribe_slot.call_count == 2 + + ingestor_manager.resource_manager.renew_tag_lease.assert_any_call("slot1") + ingestor_manager.resource_manager.renew_tag_lease.assert_any_call("slot3") + ingestor_manager.resource_manager.renew_tag_lease.assert_any_call("slot2") + assert ingestor_manager.resource_manager.renew_tag_lease.call_count == 3 def test_drop_slot_leases(ingestor_manager): @@ -211,7 +311,7 @@ def test_subscribe_to_tags(ingestor_manager): 'slot1': { "server1": {"tags": "config1"}, "server2": {"tags": "config2"}, - 'server3': {"tags": "config3"} + 'server3': {"tags": "config3"}, } } diff --git a/tests/managers/test_opc_manager.py b/tests/managers/test_opc_manager.py index 597c621..772ce92 100644 --- a/tests/managers/test_opc_manager.py +++ b/tests/managers/test_opc_manager.py @@ -151,7 +151,8 @@ def test_subscribe_success(opc_manager_subscribed): opc_manager_subscribed.nodes = { 'ns=3;i=1001': 'data' } - opc_manager_subscribed.subscribe(tags, 1000) + + opc_manager_subscribed.subscribe('sub1', tags, 1000) assert opc_manager_subscribed.nodes == tags assert opc_manager_subscribed.addr_nodes == [