SIENTIAPDE-1312
Enhance OpcManager to accept subscription period from configuration - Updated OpcManager to include subscription_period_ms as a parameter for better configurability. - Refactored create_subscription method to utilize the new subscription period parameter. - Adjusted unit tests to reflect changes in subscription handling and ensure correct functionality.
This commit is contained in:
@@ -141,6 +141,7 @@ class IngestorManager(BaseActivity):
|
|||||||
manager = OpcManager(
|
manager = OpcManager(
|
||||||
name=server_config['name'],
|
name=server_config['name'],
|
||||||
url=server_config['url'],
|
url=server_config['url'],
|
||||||
|
subscription_period_ms=server_config['subscription_period_ms'],
|
||||||
data_manager=self.data_manager,
|
data_manager=self.data_manager,
|
||||||
logger=self.logger,
|
logger=self.logger,
|
||||||
server_uri=server_config['server_uri'],
|
server_uri=server_config['server_uri'],
|
||||||
|
|||||||
@@ -59,6 +59,7 @@ class OpcManager(BaseActivity):
|
|||||||
self,
|
self,
|
||||||
name: str,
|
name: str,
|
||||||
url: str,
|
url: str,
|
||||||
|
subscription_period_ms: int,
|
||||||
data_manager: DataManager,
|
data_manager: DataManager,
|
||||||
logger: Logger,
|
logger: Logger,
|
||||||
server_uri: str,
|
server_uri: str,
|
||||||
@@ -74,6 +75,7 @@ class OpcManager(BaseActivity):
|
|||||||
self.data_queue: dict = {}
|
self.data_queue: dict = {}
|
||||||
self.non_receive_count = 0
|
self.non_receive_count = 0
|
||||||
self.client: Client | None = None
|
self.client: Client | None = None
|
||||||
|
self.subscription_period_ms = subscription_period_ms
|
||||||
self.cert_path = cert_path
|
self.cert_path = cert_path
|
||||||
self.private_key_path = private_key_path
|
self.private_key_path = private_key_path
|
||||||
self.server_cert_path = server_cert_path
|
self.server_cert_path = server_cert_path
|
||||||
@@ -82,6 +84,7 @@ class OpcManager(BaseActivity):
|
|||||||
self.data_manager = data_manager
|
self.data_manager = data_manager
|
||||||
self.metadata = metadata
|
self.metadata = metadata
|
||||||
|
|
||||||
|
|
||||||
BaseActivity.__init__(
|
BaseActivity.__init__(
|
||||||
self, logger=logger, notification_handler=notification_handler, set_error_counter=True
|
self, logger=logger, notification_handler=notification_handler, set_error_counter=True
|
||||||
)
|
)
|
||||||
@@ -189,7 +192,13 @@ class OpcManager(BaseActivity):
|
|||||||
try:
|
try:
|
||||||
self.client = Client(self.url, watchdog_intervall=3600000)
|
self.client = Client(self.url, watchdog_intervall=3600000)
|
||||||
assert self.client is not None # Informa ao mypy que client não é None
|
assert self.client is not None # Informa ao mypy que client não é None
|
||||||
|
|
||||||
self.client.name = self.pod_id
|
self.client.name = self.pod_id
|
||||||
|
self.client.application_name = self.pod_id
|
||||||
|
pod_uri = f'urn:sientia-do:pod:{self.pod_id}'
|
||||||
|
self.client.application_uri = pod_uri
|
||||||
|
self.client.product_uri = pod_uri
|
||||||
|
|
||||||
if self.cert_path:
|
if self.cert_path:
|
||||||
await self.set_security()
|
await self.set_security()
|
||||||
self.logger.info(f'Starting connection to {self.name}...')
|
self.logger.info(f'Starting connection to {self.name}...')
|
||||||
@@ -199,6 +208,15 @@ class OpcManager(BaseActivity):
|
|||||||
).set(1)
|
).set(1)
|
||||||
self.logger.info(f'Connection to {self.name} successful.')
|
self.logger.info(f'Connection to {self.name} successful.')
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
|
if self.client:
|
||||||
|
try:
|
||||||
|
await self.client.disconnect()
|
||||||
|
except Exception as internal_e:
|
||||||
|
self.logger.error(f'And error occurred while creating connection from {self.name}, and an exception occurred while disconnecting: {internal_e}')
|
||||||
|
|
||||||
|
del self.client
|
||||||
|
self.client = None
|
||||||
|
|
||||||
metrics.OPC_CONNECTION_STATUS.labels(
|
metrics.OPC_CONNECTION_STATUS.labels(
|
||||||
pod_id=self.pod_id, server_name=self.name, server_url=self.url
|
pod_id=self.pod_id, server_name=self.name, server_url=self.url
|
||||||
).set(0)
|
).set(0)
|
||||||
@@ -206,7 +224,7 @@ class OpcManager(BaseActivity):
|
|||||||
self.logger.error(f'Failed to connect to {self.name}: {e}')
|
self.logger.error(f'Failed to connect to {self.name}: {e}')
|
||||||
raise
|
raise
|
||||||
|
|
||||||
async def create_subscription(self, name: str, period: int = 500):
|
async def create_subscription(self, name: str):
|
||||||
"""
|
"""
|
||||||
Creates a subscription with the specified monitoring period.
|
Creates a subscription with the specified monitoring period.
|
||||||
|
|
||||||
@@ -230,8 +248,8 @@ class OpcManager(BaseActivity):
|
|||||||
if not self.client:
|
if not self.client:
|
||||||
raise ValueError('Client not connected. Call connect first.')
|
raise ValueError('Client not connected. Call connect first.')
|
||||||
try:
|
try:
|
||||||
p = period if period is not None else 500
|
self.subscriptions[name] = await self.client.create_subscription(
|
||||||
self.subscriptions[name] = await self.client.create_subscription(p, self)
|
self.subscription_period_ms, self)
|
||||||
self.logger.info(f'Subscription {name} created on {self.name}.')
|
self.logger.info(f'Subscription {name} created on {self.name}.')
|
||||||
metrics.OPC_SUBSCRIPTIONS_CREATED.labels(
|
metrics.OPC_SUBSCRIPTIONS_CREATED.labels(
|
||||||
pod_id=self.pod_id, server_name=self.name, slot_name=name
|
pod_id=self.pod_id, server_name=self.name, slot_name=name
|
||||||
@@ -345,10 +363,12 @@ class OpcManager(BaseActivity):
|
|||||||
except Exception as sub_error:
|
except Exception as sub_error:
|
||||||
self.logger.error(f'Failed to clean up subscription: {sub_error}')
|
self.logger.error(f'Failed to clean up subscription: {sub_error}')
|
||||||
|
|
||||||
|
# Inserir backoff para disconnect, emitindo métrica e alerta
|
||||||
try:
|
try:
|
||||||
await self.client.disconnect()
|
await self.client.disconnect()
|
||||||
except Exception as conn_error:
|
except Exception as conn_error:
|
||||||
self.logger.error(f'Failed to disconnect from OPC UA server: {conn_error}')
|
self.logger.error(f'Failed to disconnect from OPC UA server: {conn_error}')
|
||||||
|
|
||||||
finally:
|
finally:
|
||||||
del self.client
|
del self.client
|
||||||
self.client = None
|
self.client = None
|
||||||
|
|||||||
@@ -50,6 +50,7 @@ def raw_opc_manager(mock_metrics):
|
|||||||
name='TestConnector',
|
name='TestConnector',
|
||||||
url='opc.tcp://localhost:4840',
|
url='opc.tcp://localhost:4840',
|
||||||
data_manager=MagicMock(),
|
data_manager=MagicMock(),
|
||||||
|
subscription_period_ms=1000,
|
||||||
logger=MagicMock(),
|
logger=MagicMock(),
|
||||||
server_uri='opc.tcp://localhost:4840',
|
server_uri='opc.tcp://localhost:4840',
|
||||||
notification_handler=MagicMock(),
|
notification_handler=MagicMock(),
|
||||||
@@ -223,24 +224,26 @@ async def test_create_subscription_no_client(raw_opc_manager):
|
|||||||
|
|
||||||
@mark.asyncio
|
@mark.asyncio
|
||||||
async def test_create_subscription_success_has_period(opc_manager):
|
async def test_create_subscription_success_has_period(opc_manager):
|
||||||
await opc_manager.create_subscription('sub1', 1000)
|
await opc_manager.create_subscription('sub1')
|
||||||
|
|
||||||
opc_manager.client.create_subscription.assert_called_once_with(1000, opc_manager)
|
opc_manager.client.create_subscription.assert_called_once_with(
|
||||||
|
opc_manager.subscription_period_ms, opc_manager)
|
||||||
assert opc_manager.subscriptions['sub1'] is not None
|
assert opc_manager.subscriptions['sub1'] is not None
|
||||||
|
|
||||||
|
|
||||||
@mark.asyncio
|
@mark.asyncio
|
||||||
async def test_create_subscription_success_no_period(opc_manager):
|
async def test_create_subscription_success_no_period(opc_manager):
|
||||||
await opc_manager.create_subscription('sub1', None)
|
await opc_manager.create_subscription('sub1')
|
||||||
|
|
||||||
opc_manager.client.create_subscription.assert_called_once_with(500, opc_manager)
|
opc_manager.client.create_subscription.assert_called_once_with(
|
||||||
|
opc_manager.subscription_period_ms, opc_manager)
|
||||||
assert opc_manager.subscriptions['sub1'] is not None
|
assert opc_manager.subscriptions['sub1'] is not None
|
||||||
|
|
||||||
|
|
||||||
@patch('ingestor.managers.opc_manager.metrics')
|
@patch('ingestor.managers.opc_manager.metrics')
|
||||||
@mark.asyncio
|
@mark.asyncio
|
||||||
async def test_create_subscription_with_metrics(metrics, opc_manager):
|
async def test_create_subscription_with_metrics(metrics, opc_manager):
|
||||||
await opc_manager.create_subscription('sub1', 1000)
|
await opc_manager.create_subscription('sub1')
|
||||||
|
|
||||||
metrics.OPC_SUBSCRIPTIONS_CREATED.labels.assert_called_once_with(
|
metrics.OPC_SUBSCRIPTIONS_CREATED.labels.assert_called_once_with(
|
||||||
pod_id=opc_manager.pod_id, server_name=opc_manager.name, slot_name='sub1'
|
pod_id=opc_manager.pod_id, server_name=opc_manager.name, slot_name='sub1'
|
||||||
@@ -572,6 +575,7 @@ def test_init_metrics_calls_correct_metric_methods(metrics):
|
|||||||
name='TestInitConnector',
|
name='TestInitConnector',
|
||||||
url='opc.tcp://init.test:4840',
|
url='opc.tcp://init.test:4840',
|
||||||
data_manager=MagicMock(),
|
data_manager=MagicMock(),
|
||||||
|
subscription_period_ms=1000,
|
||||||
logger=MagicMock(),
|
logger=MagicMock(),
|
||||||
server_uri='opc.tcp://init.test:4840/uri',
|
server_uri='opc.tcp://init.test:4840/uri',
|
||||||
notification_handler=MagicMock(),
|
notification_handler=MagicMock(),
|
||||||
|
|||||||
Reference in New Issue
Block a user