Merge pull request #26 from Aignosi/fix/SIENTIAPDE-1314-ajustes-nas-camadas-de-monitoramento-do-sientia

SIENTIAPDE-1314: Refactor IngestorManager and OpcManager for Improved Server Management
This commit is contained in:
Matheus Demoner
2025-10-27 14:27:56 -03:00
committed by GitHub
6 changed files with 56 additions and 43 deletions

View File

@@ -1,7 +1,8 @@
import os import os
import argparse import argparse
from pathspec import PathSpec from pathspec import PathSpec
import yaml import yaml # type: ignore
from typing import Any
''' '''
Usage: Usage:
@@ -33,7 +34,7 @@ def encode_file_tree_to_yaml(directory, ignore_file, include_library):
"""Encode the file tree into a single YAML file.""" """Encode the file tree into a single YAML file."""
ignore_patterns = load_ignore_patterns( ignore_patterns = load_ignore_patterns(
ignore_file, include_library) if ignore_file else None ignore_file, include_library) if ignore_file else None
file_tree = {} file_tree: dict[str, Any] = {}
for root, dirs, files in os.walk(directory): for root, dirs, files in os.walk(directory):
# Skip ignored directories # Skip ignored directories

View File

@@ -189,7 +189,17 @@ class IngestorManager(BaseActivity):
def __del__(self): def __del__(self):
asyncio.run(self.shutdown()) asyncio.run(self.shutdown())
async def update_opc_servers(self): # NOSONAR async def remove_server(self, server: str):
"""
Removes an OPC server from the ingestor.
"""
if server in self.opc_managers:
await self.opc_managers[server].shutdown()
del self.opc_managers[server]
for slot, _config in self.managed_tags.items():
self.managed_tags[slot].pop(server, None)
async def update_opc_servers(self):
""" """
Updates the OPC (OLE for Process Control) server connections managed by the ingestor. Updates the OPC (OLE for Process Control) server connections managed by the ingestor.
This method ensures that the OPC servers defined in `self.managed_tags` are properly This method ensures that the OPC servers defined in `self.managed_tags` are properly
@@ -244,11 +254,7 @@ class IngestorManager(BaseActivity):
f'removing server from managed tags.' f'removing server from managed tags.'
) )
for slot, _value in current_managed_tags.items(): await self.remove_server(server)
if server in self.opc_managers:
await self.opc_managers[server].shutdown()
del self.opc_managers[server]
self.managed_tags[slot].pop(server, None)
servers = list(self.opc_managers.keys()) servers = list(self.opc_managers.keys())
for server in servers: for server in servers:
@@ -256,8 +262,7 @@ class IngestorManager(BaseActivity):
self.logger.warning( self.logger.warning(
f'Server {server} not found in managed tags. Desconnecting from server.' f'Server {server} not found in managed tags. Desconnecting from server.'
) )
await self.opc_managers[server].shutdown() await self.remove_server(server)
self.opc_managers.pop(server, None)
metrics.OPC_MANAGERS_ACTIVE.labels(pod_id=self.pod_id).set(len(self.opc_managers)) metrics.OPC_MANAGERS_ACTIVE.labels(pod_id=self.pod_id).set(len(self.opc_managers))
@@ -276,6 +281,7 @@ class IngestorManager(BaseActivity):
- Removes lost servers from managed tags - Removes lost servers from managed tags
- Updates OPC manager metrics - Updates OPC manager metrics
""" """
to_disconnect: list[str] = []
for server, opc_manager in self.opc_managers.items(): for server, opc_manager in self.opc_managers.items():
opc_manager.check_cycles() opc_manager.check_cycles()
@@ -283,11 +289,10 @@ class IngestorManager(BaseActivity):
if is_lost: if is_lost:
self.logger.warning(f'OPC server {server} is lost. Server will be disconnected.') self.logger.warning(f'OPC server {server} is lost. Server will be disconnected.')
for slot, _config in self.managed_tags.items(): to_disconnect.append(server)
if server in self.opc_managers:
await self.opc_managers[server].shutdown() for server in to_disconnect:
del self.opc_managers[server] await self.remove_server(server)
self.managed_tags[slot].pop(server, None)
metrics.OPC_MANAGERS_ACTIVE.labels(pod_id=self.pod_id).set(len(self.opc_managers)) metrics.OPC_MANAGERS_ACTIVE.labels(pod_id=self.pod_id).set(len(self.opc_managers))

View File

@@ -154,21 +154,21 @@ class OpcManager(BaseActivity):
raise ValueError( raise ValueError(
'Certificate and private key paths must be provided for secure connection.' 'Certificate and private key paths must be provided for secure connection.'
) )
cert = Path(self.cert_path) if self.cert_path else None cert = str(Path(self.cert_path)) if self.cert_path else None
private_key = Path(self.private_key_path) if self.private_key_path else None private_key = str(Path(self.private_key_path)) if self.private_key_path else None
server_cert = Path(self.server_cert_path) if self.server_cert_path else None server_cert = str(Path(self.server_cert_path)) if self.server_cert_path else None
if self.client: if self.client:
await self.client.set_application_uri(self.server_uri) self.client.application_uri = self.server_uri
self.logger.info('Setting security...') self.logger.info('Setting security...')
await self.client.set_security( await self.client.set_security(
SecurityPolicyBasic256, SecurityPolicyBasic256,
certificate=str(cert), certificate=cert,
private_key=str(private_key), private_key=private_key,
server_certificate=str(server_cert), server_certificate=server_cert,
) )
await self.client.set_secure_channel_timeout(10000000) self.client.secure_channel_timeout = 10000000
await self.client.set_session_timeout(10000000) self.client.session_timeout = 10000000
async def connect(self): async def connect(self):
""" """
@@ -286,9 +286,7 @@ class OpcManager(BaseActivity):
self.logger.debug(f'Addr nodes: {addr_nodes}') self.logger.debug(f'Addr nodes: {addr_nodes}')
self.nodes.update(nodes) self.nodes.update(nodes)
self.logger.debug(f'Nodes: {self.nodes}') self.logger.debug(f'Nodes: {self.nodes}')
metrics.OPC_TAGS_SUBSCRIBED.labels(pod_id=self.pod_id, server_name=self.name).set(
len(self.nodes)
)
self.collect_period = collect_period self.collect_period = collect_period
for node, config in self.nodes.items(): for node, config in self.nodes.items():
@@ -299,6 +297,10 @@ class OpcManager(BaseActivity):
await self.subscriptions[subscription].subscribe_data_change(addr_nodes) await self.subscriptions[subscription].subscribe_data_change(addr_nodes)
metrics.OPC_TAGS_SUBSCRIBED.labels(pod_id=self.pod_id, server_name=self.name).set(
len(self.nodes)
)
async def unsubscribe(self, subscription: str): async def unsubscribe(self, subscription: str):
""" """
Unsubscribes from a given subscription. Unsubscribes from a given subscription.
@@ -437,10 +439,6 @@ class OpcManager(BaseActivity):
f'{tag} after {self.nodes[tag]["cycle_rule"]["cycle_count"]} cycles' f'{tag} after {self.nodes[tag]["cycle_rule"]["cycle_count"]} cycles'
) )
self.nodes[tag]['cycle_rule']['cycle_count'] = 0
self.non_receive_count = 0
metrics.OPC_CYCLES_WITHOUT_DATA.labels(pod_id=self.pod_id, server_name=self.name).set(0)
data = { data = {
'tag': tag, 'tag': tag,
'name': self.nodes[str(node)]['tag_name'], 'name': self.nodes[str(node)]['tag_name'],
@@ -451,6 +449,10 @@ class OpcManager(BaseActivity):
for topic in self.nodes[tag]['topics']: for topic in self.nodes[tag]['topics']:
self.data_manager.publish(topic, data) self.data_manager.publish(topic, data)
self.nodes[tag]['cycle_rule']['cycle_count'] = 0
self.non_receive_count = 0
metrics.OPC_CYCLES_WITHOUT_DATA.labels(pod_id=self.pod_id, server_name=self.name).set(0)
def check_cycles(self): def check_cycles(self):
""" """
Checks the cycle counts for all monitored nodes and sends Checks the cycle counts for all monitored nodes and sends

View File

@@ -451,7 +451,11 @@ async def test_subscribe_to_tags(ingestor_manager):
ingestor_manager.manage_server = AsyncMock(side_effect=[0, 1, 2]) ingestor_manager.manage_server = AsyncMock(side_effect=[0, 1, 2])
ingestor_manager.managed_tags = {'slot1': MagicMock(), 'slot2': MagicMock()} ingestor_manager.managed_tags = {'slot1': MagicMock(), 'slot2': MagicMock()}
ingestor_manager.opc_managers = {'server1': AsyncMock(), 'server2': AsyncMock()} ingestor_manager.opc_managers = {
'server1': AsyncMock(),
'server2': AsyncMock(),
'server3': AsyncMock(),
}
ingestor_manager.subscriptions = {'server1': AsyncMock()} ingestor_manager.subscriptions = {'server1': AsyncMock()}
tags = { tags = {
'slot1': { 'slot1': {
@@ -509,21 +513,22 @@ async def test_check_opc_servers_integrity_all_healthy(metrics, ingestor_manager
) )
def test_check_opc_servers_integrity_server_lost(ingestor_manager): @mark.asyncio
async def test_check_opc_servers_integrity_server_lost(ingestor_manager):
# Setup mock OPC manager that will be lost # Setup mock OPC manager that will be lost
opc_manager = MagicMock() opc_manager = AsyncMock()
opc_manager.check_cycles.return_value = None opc_manager.check_cycles.return_value = None
opc_manager.check_opc_listenning.return_value = True # Server is lost opc_manager.check_opc_listenning.return_value = True # Server is lost
opc_manager.config = {'config': 'config1'} opc_manager.config = {'config': 'config1'}
ingestor_manager.opc_managers = {'server1': opc_manager} ingestor_manager.opc_managers = {'server1': opc_manager}
ingestor_manager.managed_tags = {'slot1': MagicMock()}
# Mock the initialize_opc_from_config method to return a new manager
new_manager = MagicMock()
ingestor_manager.initialize_opc_from_config = MagicMock(return_value=new_manager)
# Call the method # Call the method
ingestor_manager.check_opc_servers_integrity() await ingestor_manager.check_opc_servers_integrity()
assert 'server1' not in ingestor_manager.opc_managers
ingestor_manager.managed_tags['slot1'].pop.assert_called_once_with('server1', None)
def test_check_opc_servers_integrity_server_lost_with_tags(ingestor_manager): def test_check_opc_servers_integrity_server_lost_with_tags(ingestor_manager):

View File

@@ -105,7 +105,7 @@ async def test_shutdown_error(opc_manager):
async def test_set_security_success(opc_manager): async def test_set_security_success(opc_manager):
await opc_manager.set_security() await opc_manager.set_security()
opc_manager.client.set_application_uri.assert_called_once_with(opc_manager.server_uri) assert opc_manager.client.application_uri == opc_manager.server_uri
opc_manager.client.set_security.assert_called_once_with( opc_manager.client.set_security.assert_called_once_with(
SecurityPolicyBasic256, SecurityPolicyBasic256,
@@ -114,8 +114,8 @@ async def test_set_security_success(opc_manager):
server_certificate=opc_manager.server_cert_path, server_certificate=opc_manager.server_cert_path,
) )
opc_manager.client.set_secure_channel_timeout.assert_called_once_with(10000000) assert opc_manager.client.secure_channel_timeout == 10000000
opc_manager.client.set_session_timeout.assert_called_once_with(10000000) assert opc_manager.client.session_timeout == 10000000
@mark.asyncio @mark.asyncio

View File

@@ -139,7 +139,7 @@ env:
- name: GITHUB_REPO_URL - name: GITHUB_REPO_URL
value: "git@github.com:Aignosi/sientia-dataops-opc-ingestor.git" value: "git@github.com:Aignosi/sientia-dataops-opc-ingestor.git"
- name: GITHUB_BRANCH - name: GITHUB_BRANCH
value: "SIENTIAPDE-1312-melhorias-e-correcoes-nas-pipelines-de-dados" value: "fix/SIENTIAPDE-1314-ajustes-nas-camadas-de-monitoramento-do-sientia"
- name: PYTHON_APP - name: PYTHON_APP
value: "ingestor.app" value: "ingestor.app"