SIENTIAPDE-1314

Refactor IngestorManager and OpcManager for improved server management

- Introduced a new method `remove_server` in IngestorManager to encapsulate server removal logic, enhancing code clarity and reusability.
- Updated `update_opc_servers` to utilize the new `remove_server` method for better maintainability.
- Refactored security settings in OpcManager to use direct attribute assignments instead of method calls, streamlining the code.
- Adjusted unit tests to reflect changes in the IngestorManager and OpcManager, ensuring proper asynchronous behavior and mocking.
This commit is contained in:
vitor-aignosi
2025-10-23 15:11:32 -03:00
parent cc0fde664b
commit 56023dfdf4
4 changed files with 54 additions and 38 deletions

View File

@@ -189,7 +189,17 @@ class IngestorManager(BaseActivity):
def __del__(self):
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.
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.'
)
for slot, _value in current_managed_tags.items():
if server in self.opc_managers:
await self.opc_managers[server].shutdown()
del self.opc_managers[server]
self.managed_tags[slot].pop(server, None)
await self.remove_server(server)
servers = list(self.opc_managers.keys())
for server in servers:
@@ -256,8 +262,7 @@ class IngestorManager(BaseActivity):
self.logger.warning(
f'Server {server} not found in managed tags. Desconnecting from server.'
)
await self.opc_managers[server].shutdown()
self.opc_managers.pop(server, None)
await self.remove_server(server)
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
- Updates OPC manager metrics
"""
to_disconnect: list[str] = []
for server, opc_manager in self.opc_managers.items():
opc_manager.check_cycles()
@@ -284,11 +290,14 @@ class IngestorManager(BaseActivity):
self.logger.warning(f'OPC server {server} is lost. Server will be disconnected.')
for slot, _config in self.managed_tags.items():
if server in self.opc_managers:
await self.opc_managers[server].shutdown()
del self.opc_managers[server]
to_disconnect.append(server)
self.managed_tags[slot].pop(server, None)
for server in to_disconnect:
if server in self.opc_managers:
await self.opc_managers[server].shutdown()
del self.opc_managers[server]
metrics.OPC_MANAGERS_ACTIVE.labels(pod_id=self.pod_id).set(len(self.opc_managers))
def declare_active(self):

View File

@@ -154,21 +154,21 @@ class OpcManager(BaseActivity):
raise ValueError(
'Certificate and private key paths must be provided for secure connection.'
)
cert = Path(self.cert_path) if self.cert_path else None
private_key = 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
cert = str(Path(self.cert_path)) if self.cert_path else None
private_key = str(Path(self.private_key_path)) if self.private_key_path else None
server_cert = str(Path(self.server_cert_path)) if self.server_cert_path else None
if self.client:
await self.client.set_application_uri(self.server_uri)
self.client.application_uri = self.server_uri
self.logger.info('Setting security...')
await self.client.set_security(
SecurityPolicyBasic256,
certificate=str(cert),
private_key=str(private_key),
server_certificate=str(server_cert),
certificate=cert,
private_key=private_key,
server_certificate=server_cert,
)
await self.client.set_secure_channel_timeout(10000000)
await self.client.set_session_timeout(10000000)
self.client.secure_channel_timeout = 10000000
self.client.session_timeout = 10000000
async def connect(self):
"""
@@ -286,9 +286,7 @@ class OpcManager(BaseActivity):
self.logger.debug(f'Addr nodes: {addr_nodes}')
self.nodes.update(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
for node, config in self.nodes.items():
@@ -299,6 +297,10 @@ class OpcManager(BaseActivity):
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):
"""
Unsubscribes from a given subscription.
@@ -437,10 +439,6 @@ class OpcManager(BaseActivity):
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 = {
'tag': tag,
'name': self.nodes[str(node)]['tag_name'],
@@ -451,6 +449,10 @@ class OpcManager(BaseActivity):
for topic in self.nodes[tag]['topics']:
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):
"""
Checks the cycle counts for all monitored nodes and sends