From 10b331067792fa613c50d1a5440e2e43e454bf16 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Wed, 29 Oct 2025 16:01:54 -0300 Subject: [PATCH 01/10] SIENTIAPDE-1325 Refactor MongoDB, Redis, and Gates classes to extend SientiaMonitoring instead of BaseActivity. Update requirements to point to local dataops library path. Implement repository pattern for MongoDB and Redis operations, enhancing code organization and maintainability. --- requirements.txt | 2 +- scouter/activities/gates.py | 12 +++-- scouter/activities/mongodb.py | 87 +++++++++-------------------------- scouter/activities/redis.py | 32 +++++++++---- 4 files changed, 55 insertions(+), 78 deletions(-) diff --git a/requirements.txt b/requirements.txt index 8589717..ce4f826 100644 --- a/requirements.txt +++ b/requirements.txt @@ -5,6 +5,6 @@ asyncua redis aiokafka pymongo -git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.4.7 +/home/grezewave/Documents/projects/sientia/sientia-dataops-library pydruid[pandas] prometheus-client \ No newline at end of file diff --git a/scouter/activities/gates.py b/scouter/activities/gates.py index ca9e2ce..cb9ebf6 100644 --- a/scouter/activities/gates.py +++ b/scouter/activities/gates.py @@ -10,7 +10,7 @@ with workflow.unsafe.imports_passed_through(): from sientia_do.notifications.handlers import NotificationHandler from sientia_do.notifications.models import NotificationLevel from sientia_do.observability.logger import Logger - from sientia_do.temporal.activities.base import BaseActivity + from sientia_do.observability.sientia_monitoring import SientiaMonitoring from scouter import metrics from scouter.utils.quality.filters import null_values_filter, out_of_bounds_filter @@ -21,7 +21,7 @@ quality_gate_filters = { } -class Gates(BaseActivity): +class Gates(SientiaMonitoring): """ Data quality gates and filtering operations. @@ -44,7 +44,13 @@ class Gates(BaseActivity): logger (Logger): Logger instance for operation logging notification_handler (NotificationHandler): Handler for system notifications """ - BaseActivity.__init__(self, logger, notification_handler, set_error_counter=True) + SientiaMonitoring.__init__(self, logger, notification_handler) + + def close(self): + """ + Close the Gates class. + """ + SientiaMonitoring.close(self) def apply_aggregation( self, values: DataFrame, aggr_function: str, metadata: dict[str, Any] diff --git a/scouter/activities/mongodb.py b/scouter/activities/mongodb.py index a20e356..586d1a6 100644 --- a/scouter/activities/mongodb.py +++ b/scouter/activities/mongodb.py @@ -13,45 +13,13 @@ with workflow.unsafe.imports_passed_through(): from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler from sientia_do.notifications.models import NotificationLevel from sientia_do.observability.logger import Logger - from sientia_do.temporal.activities.base import BaseActivity + from sientia_do.repository.mongodb_repository import MongoDBRepository + from sientia_do.observability.sientia_monitoring import SientiaMonitoring from sientia_do.temporal.constants import DATETIME_FORMAT_MS_WITH_TZ -def clear_mongo_id(docs: list) -> list: - """ - Remove MongoDB internal `_id` fields from documents. - This utility function recursively removes the MongoDB `_id` field from - documents and nested structures. It's used to clean data before - processing or export operations. - - Args: - docs (list): List of documents to clean - - Returns: - list: Documents with `_id` fields removed - - Note: - This function modifies the input list in-place and returns the same reference - """ - for doc in docs: - if isinstance(doc, list): - clear_mongo_id(doc) - - elif isinstance(doc, dict): - if '_id' in doc: - del doc['_id'] - - for _key, value in doc.items(): - if isinstance(value, list): - clear_mongo_id(value) - elif isinstance(value, dict): - clear_mongo_id([value]) - - return docs - - -class MongoDB(BaseActivity): +class MongoDB(SientiaMonitoring): """ MongoDB operations for data retrieval and storage. @@ -85,37 +53,22 @@ class MongoDB(BaseActivity): Raises: ConnectionError: If MongoDB connection fails """ - self.connection_string = connection_string - self.database_name = database_name - - self.client: MongoClient = MongoClient( - self.connection_string, serverSelectionTimeoutMS=5000 - ) - self.client.server_info() # Trigger an exception if connection fails - - self.database = self.client[self.database_name] - - # Initialize MongoDB client here (omitted for brevity) - logger.info('MongoDB connection initialized') - - BaseActivity.__init__( - self, logger=logger, notification_handler=notification_handler, set_error_counter=True + + self.mongodb_repository = MongoDBRepository( + connection_string=connection_string, + database_name=database_name, + logger=logger, + notification_handler=notification_handler, ) - def shutdown(self): - """ - Gracefully close MongoDB client connection. + SientiaMonitoring.__init__(self, logger=logger, notification_handler=notification_handler) - This method ensures proper cleanup of MongoDB connections to prevent - connection leaks and ensure graceful application termination. + def close(self): """ - try: - if self.client: - self.logger.info('Closing MongoDB connection...') - self.client.close() - self.logger.info('MongoDB connection closed successfully') - except Exception as e: - self.logger.error(f'Failed to close MongoDB connection: {e}') + Close the MongoDB connection. + """ + self.mongodb_repository.close() + SientiaMonitoring.shutdown(self) def __del__(self): """ @@ -124,7 +77,7 @@ class MongoDB(BaseActivity): This destructor ensures that MongoDB connections are properly closed when the object is garbage collected, preventing resource leaks. """ - self.shutdown() + self.close() @activity.defn(name='load_latest_data') async def load_latest_data(self, input_data: dict[str, Any]) -> dict[Hashable, Any]: @@ -166,9 +119,11 @@ class MongoDB(BaseActivity): self.debug(f'Data filter: {data_filter}', metadata=metadata) - data = list(self.database[collection_name].find(data_filter, {'_id': 0})) - - data = clear_mongo_id(data) + data = self.mongodb_repository.find( + collection_name=collection_name, + filters=data_filter, + metadata=metadata, + ) self.debug(f'Collected: {data}', metadata=metadata) diff --git a/scouter/activities/redis.py b/scouter/activities/redis.py index 9c76fe4..e4513a2 100644 --- a/scouter/activities/redis.py +++ b/scouter/activities/redis.py @@ -10,11 +10,12 @@ with workflow.unsafe.imports_passed_through(): from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler from sientia_do.notifications.models import NotificationLevel from sientia_do.observability.logger import Logger - from sientia_do.temporal.activities.redis_base import Redis as RedisBase + from sientia_do.observability.sientia_monitoring import SientiaMonitoring + from sientia_do.repository.redis_repository import RedisRepository from sientia_do.temporal.constants import DATETIME_FORMAT, now -class Redis(RedisBase): +class Redis(SientiaMonitoring): """ Redis operations for data caching and temporary storage. @@ -49,7 +50,22 @@ class Redis(RedisBase): logger (Logger): Logger instance for operation logging notification_handler (NotificationHandler): Handler for system notifications """ - RedisBase.__init__(self, host, port, username, password, logger, notification_handler) + SientiaMonitoring.__init__(self, logger, notification_handler) + self.redis_repository = RedisRepository( + host=host, + port=port, + username=username, + password=password, + logger=logger, + notification_handler=notification_handler, + ) + + def close(self): + """ + Close the Redis connection. + """ + self.redis_repository.close() + SientiaMonitoring.shutdown(self) @activity.defn(name='get_last_data_timestamp') async def get_last_data_timestamp(self, input_data: dict[str, Any]) -> str | None: @@ -79,7 +95,7 @@ class Redis(RedisBase): self.info(f'Getting last data timestamp for {key}') try: - data_hold = self.get(key) + data_hold = self.redis_repository.get(key) except Exception as e: self.send_notification( metadata=metadata, @@ -137,7 +153,7 @@ class Redis(RedisBase): self.info(f'Last collected timestamp to insert: {last_data_timestamp}', metadata=metadata) try: - self.set(key, last_data_timestamp, ttl=60 * 60 * 5) + self.redis_repository.set(key, last_data_timestamp, ttl=60 * 60 * 5) except Exception as e: self.send_notification( metadata=metadata, @@ -190,7 +206,7 @@ class Redis(RedisBase): self.info(f'Getting held data for {key}') try: - data_hold = self.get(key) + data_hold = self.redis_repository.get(key) except Exception as e: self.send_notification( metadata=metadata, @@ -234,7 +250,7 @@ class Redis(RedisBase): data['timestamp'].max() if not data.empty else data_hold['timestamp'] ) - self.set(key, data_hold, ttl=retention_time) + self.redis_repository.set(key, data_hold, ttl=retention_time) data_hold_df = DataFrame(data_hold, index=[0]) data_hold_melted = data_hold_df.melt( @@ -280,7 +296,7 @@ class Redis(RedisBase): cache = {'data': data.to_dict(), 'held_data': held_data.to_dict()} try: - self.set(key, cache, ttl=120) + self.redis_repository.set(key, cache, ttl=120) except Exception as e: self.send_notification( metadata=metadata, From 0c160758239657e5781919cb938ba3de47740b6a Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Wed, 29 Oct 2025 16:58:22 -0300 Subject: [PATCH 02/10] SIENTIAPDE-1325 Update GITHUB_BRANCH in values.yaml to feature/SIENTIAPDE-1325 for specific external operation metrics. Refactor shutdown methods in Activities class to ensure proper closure of MongoDB, Redis, and Gates connections. Adjust tests to reflect these changes and enhance repository pattern implementation for MongoDB and Redis activities. --- scouter/activities/activities.py | 4 +- scouter/activities/mongodb.py | 9 +- tests/activities/test_activities.py | 8 +- tests/activities/test_gates.py | 3 + tests/activities/test_mongo.py | 126 ++++++++++++---------------- tests/activities/test_redis.py | 115 ++++++++++++++----------- tests/conftest.py | 12 +++ values.yaml | 2 +- 8 files changed, 152 insertions(+), 127 deletions(-) create mode 100644 tests/conftest.py diff --git a/scouter/activities/activities.py b/scouter/activities/activities.py index 79ddfe1..1dd20d8 100644 --- a/scouter/activities/activities.py +++ b/scouter/activities/activities.py @@ -97,4 +97,6 @@ class Activities(Postgres, Redis, Gates, MongoDB): to prevent connection leaks and ensure graceful application termination. """ Postgres.close(self) - MongoDB.shutdown(self) + MongoDB.close(self) + Redis.close(self) + Gates.close(self) diff --git a/scouter/activities/mongodb.py b/scouter/activities/mongodb.py index 586d1a6..55450eb 100644 --- a/scouter/activities/mongodb.py +++ b/scouter/activities/mongodb.py @@ -8,17 +8,14 @@ with workflow.unsafe.imports_passed_through(): from datetime import datetime from typing import Any - from pandas import DataFrame - from pymongo import MongoClient from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler from sientia_do.notifications.models import NotificationLevel from sientia_do.observability.logger import Logger - from sientia_do.repository.mongodb_repository import MongoDBRepository from sientia_do.observability.sientia_monitoring import SientiaMonitoring + from sientia_do.repository.mongodb_repository import MongoDBRepository from sientia_do.temporal.constants import DATETIME_FORMAT_MS_WITH_TZ - class MongoDB(SientiaMonitoring): """ MongoDB operations for data retrieval and storage. @@ -53,7 +50,7 @@ class MongoDB(SientiaMonitoring): Raises: ConnectionError: If MongoDB connection fails """ - + self.mongodb_repository = MongoDBRepository( connection_string=connection_string, database_name=database_name, @@ -136,7 +133,7 @@ class MongoDB(SientiaMonitoring): self.debug(f'Loaded data: {data}', metadata=metadata) - return DataFrame(data).to_dict() + return data except Exception as e: trace = traceback.format_exc() self.send_notification( diff --git a/tests/activities/test_activities.py b/tests/activities/test_activities.py index 2ec9bb5..8864c5f 100644 --- a/tests/activities/test_activities.py +++ b/tests/activities/test_activities.py @@ -88,8 +88,12 @@ def test___init__(mock_gates_init, mock_redis_init, mock_postgres_init, mock_mon @patch('scouter.activities.activities.Gates.__init__') @patch('scouter.activities.activities.MongoDB.__init__') @patch('scouter.activities.activities.Postgres.close') -@patch('scouter.activities.activities.MongoDB.shutdown') +@patch('scouter.activities.activities.MongoDB.close') +@patch('scouter.activities.activities.Redis.close') +@patch('scouter.activities.activities.Gates.close') def test_shutdown( + mock_gates_close, + mock_redis_close, mock_mongodb_close, mock_postgres_close, _mock_mongodb_init, @@ -129,3 +133,5 @@ def test_shutdown( mock_postgres_close.assert_called() mock_mongodb_close.assert_called() + mock_redis_close.assert_called() + mock_gates_close.assert_called() diff --git a/tests/activities/test_gates.py b/tests/activities/test_gates.py index c916ace..726c666 100644 --- a/tests/activities/test_gates.py +++ b/tests/activities/test_gates.py @@ -16,6 +16,9 @@ def gates_fixture(): notification_handler = MagicMock() gates = Gates(logger=logger, notification_handler=notification_handler) gates.send_notification = MagicMock() + gates.logger = logger + gates.notification_handler = notification_handler + gates.pod_id = 'localhost' return gates diff --git a/tests/activities/test_mongo.py b/tests/activities/test_mongo.py index 0d62e0a..268b02c 100644 --- a/tests/activities/test_mongo.py +++ b/tests/activities/test_mongo.py @@ -5,54 +5,38 @@ from pytest import fixture, mark from sientia_do.notifications.models import NotificationLevel from sientia_do.temporal.constants import DATETIME_FORMAT_MS_WITH_TZ -from scouter.activities.mongodb import MongoDB, clear_mongo_id +from scouter.activities.mongodb import MongoDB -def test_clear_mongo_id(): - """Test clear_mongo_id""" - data = [ - {'_id': '1', 'name': 'test1'}, - {'_id': '2', 'name': [{'_id': '3', 'name': 'test3'}]}, - {'_id': '4', 'name': {'_id': '5', 'name': 'test2'}}, - [{'_id': '6', 'name': 'test2'}], - ] - - result = clear_mongo_id(data) - - assert result == [ - {'name': 'test1'}, - {'name': [{'name': 'test3'}]}, - {'name': {'name': 'test2'}}, - [{'name': 'test2'}], - ] - - -@patch('scouter.activities.mongodb.MongoClient') -def test_mongodb___init__(mock_mongo_client): +@patch('scouter.activities.mongodb.MongoDBRepository') +def test_mongodb___init__(mock_mongodb_repository): """Test MongoDB __init__""" + + logger = MagicMock() + notification_handler = MagicMock() + mongo = MongoDB( connection_string='mongodb://localhost:27017', database_name='test_db', - logger=MagicMock(), - notification_handler=MagicMock(), + logger=logger, + notification_handler=notification_handler, ) - mock_mongo_client.assert_called_once_with( - 'mongodb://localhost:27017', serverSelectionTimeoutMS=5000 + mock_mongodb_repository.assert_called_once_with( + connection_string='mongodb://localhost:27017', + database_name='test_db', + logger=logger, + notification_handler=notification_handler, ) - mock_mongo_client.return_value.server_info.assert_called_once() - - mock_mongo_client.return_value.__getitem__.assert_called_once_with('test_db') - - assert mongo.client is not None - assert mongo.database is not None + assert mongo.mongodb_repository is not None @fixture -@patch('scouter.activities.mongodb.MongoClient') -def mongodb_activity(mock_mongo_client): +@patch('scouter.activities.mongodb.MongoDBRepository') +def mongodb_activity(mock_mongodb_repository): """Test MongoDB activity""" + mongo = MongoDB( connection_string='mongodb://localhost:27017', database_name='test_db', @@ -60,32 +44,32 @@ def mongodb_activity(mock_mongo_client): notification_handler=MagicMock(), ) + mongo.logger = MagicMock() + mongo.notification_handler = MagicMock() + mongo.pod_id = 'localhost' return mongo -def test_shutdown_success(mongodb_activity): - """Test shutdown""" - mongodb_activity.shutdown() +def test_close(mongodb_activity): + """Test close""" + mongodb_activity.close() - mongodb_activity.client.close.assert_called_once() + mongodb_activity.mongodb_repository.close.assert_called_once() -def test_shutdown_error(mongodb_activity): - """Test shutdown""" - mongodb_activity.client.close = MagicMock(side_effect=Exception('test')) +def test_del(mongodb_activity): + mongodb_activity.close = MagicMock() - mongodb_activity.shutdown() + mongodb_activity.__del__() - mongodb_activity.client.close.assert_called_once() + mongodb_activity.close.assert_called_once() @mark.asyncio async def test_load_latest_data_none_last_data_timestamp(mongodb_activity): """Test load_latest_data""" - collection = MagicMock() - mongodb_activity.database.__getitem__.return_value = collection - collection.find.return_value = [ + mongodb_activity.mongodb_repository.find.return_value = [ { 'name': 'test1', 'value': 1, @@ -103,24 +87,26 @@ async def test_load_latest_data_none_last_data_timestamp(mongodb_activity): } ) - mongodb_activity.database.__getitem__.assert_called_once_with('test_collection') + mongodb_activity.mongodb_repository.find.assert_called_once_with( + collection_name='test_collection', + filters={}, + metadata={'workflow_name': 'test_pipeline', 'schedule_name': 'test_schedule'}, + ) - collection.find.assert_called_once_with({}, {'_id': 0}) - - assert result == { - 'name': {0: 'test1'}, - 'value': {0: 1}, - 'inserted_at': {0: '2023-01-01 12:00:00.000000+0000'}, - } + assert result == [ + { + 'name': 'test1', + 'value': 1, + 'inserted_at': '2023-01-01 12:00:00.000000+0000', + } + ] @mark.asyncio async def test_load_latest_data_not_none_last_data_timestamp(mongodb_activity): """Test load_latest_data""" - collection = MagicMock() - mongodb_activity.database.__getitem__.return_value = collection - collection.find.return_value = [ + mongodb_activity.mongodb_repository.find.return_value = [ { 'name': 'test1', 'value': 1, @@ -138,34 +124,32 @@ async def test_load_latest_data_not_none_last_data_timestamp(mongodb_activity): } ) - mongodb_activity.database.__getitem__.assert_called_once_with('test_collection') - - collection.find.assert_called_once_with( - { + mongodb_activity.mongodb_repository.find.assert_called_once_with( + collection_name='test_collection', + filters={ 'inserted_at': { '$gt': datetime.strptime( '2023-01-01 12:00:00.000000+0000', DATETIME_FORMAT_MS_WITH_TZ ) } }, - {'_id': 0}, + metadata={'workflow_name': 'test_pipeline', 'schedule_name': 'test_schedule'}, ) - assert result == { - 'name': {0: 'test1'}, - 'value': {0: 1}, - 'inserted_at': {0: '2023-01-01 12:00:00.000000+0000'}, - } + assert result == [ + { + 'name': 'test1', + 'value': 1, + 'inserted_at': '2023-01-01 12:00:00.000000+0000', + } + ] @mark.asyncio async def test_load_latest_data_error(mongodb_activity): """Test load_latest_data""" - collection = MagicMock() + mongodb_activity.mongodb_repository.find.side_effect = Exception('test') mongodb_activity.send_notification = MagicMock() - mongodb_activity.database.__getitem__.return_value = collection - - collection.find.side_effect = Exception('test') try: await mongodb_activity.load_latest_data( diff --git a/tests/activities/test_redis.py b/tests/activities/test_redis.py index b220424..0daf867 100644 --- a/tests/activities/test_redis.py +++ b/tests/activities/test_redis.py @@ -11,8 +11,8 @@ from scouter.activities.redis import Redis @pytest.fixture -@patch('scouter.activities.redis.RedisBase.__init__') -def redis_activity(_mock_redis_init): +@patch('scouter.activities.redis.RedisRepository') +def redis_activity(mock_redis_repository): logger = MagicMock() notification_handler = MagicMock(spec=NotificationHandler) activity = Redis( @@ -31,24 +31,6 @@ def redis_activity(_mock_redis_init): return activity -@patch('scouter.activities.redis.RedisBase.__init__') -def test_redis_initialization(mock_redis_init): - """Test Redis activity initialization""" - logger = MagicMock() - notification_handler = MagicMock(spec=NotificationHandler) - Redis( - host='localhost', - port=6379, - logger=logger, - notification_handler=notification_handler, - username='test', - password='test', - ) - mock_redis_init.assert_called_once_with( - ANY, 'localhost', 6379, 'test', 'test', logger, notification_handler - ) - - metadata = { 'metadata': { 'model_id': 'test_model_id', @@ -59,15 +41,52 @@ metadata = { } +@patch('scouter.activities.redis.RedisRepository') +def test_redis_initialization(mock_redis_repository): + """Test Redis activity initialization""" + logger = MagicMock() + notification_handler = MagicMock(spec=NotificationHandler) + activity = Redis( + host='localhost', + port=6379, + logger=logger, + notification_handler=notification_handler, + username='test', + password='test', + ) + mock_redis_repository.assert_called_once_with( + host='localhost', + port=6379, + username='test', + password='test', + logger=logger, + notification_handler=notification_handler, + ) + + activity.logger = logger + activity.notification_handler = notification_handler + activity.pod_id = 'localhost' + + assert activity.redis_repository is not None + + @pytest.mark.asyncio async def test_get_last_data_timestamp_none(redis_activity): """Test get_last_data_timestamp""" - test_data = {**metadata, 'workflow_name': 'test_pipeline', 'schedule_name': 'test_schedule'} + test_data = { + **metadata, + 'workflow_name': 'test_pipeline', + 'schedule_name': 'test_schedule', + } - redis_activity.get = MagicMock(return_value=None) + redis_activity.redis_repository.get.return_value = None result = await redis_activity.get_last_data_timestamp(test_data) + redis_activity.redis_repository.get.assert_called_once_with( + 'last_data_timestamp:test_pipeline:test_schedule' + ) + assert result is None @@ -76,11 +95,13 @@ async def test_get_last_data_timestamp_not_none(redis_activity): """Test get_last_data_timestamp""" test_data = {**metadata, 'workflow_name': 'test_pipeline', 'schedule_name': 'test_schedule'} - redis_activity.get = MagicMock(return_value='2023-01-01 12:00:00') + redis_activity.redis_repository.get.return_value = '2023-01-01 12:00:00' result = await redis_activity.get_last_data_timestamp(test_data) - redis_activity.get.assert_called_once_with('last_data_timestamp:test_pipeline:test_schedule') + redis_activity.redis_repository.get.assert_called_once_with( + 'last_data_timestamp:test_pipeline:test_schedule' + ) assert result == '2023-01-01 12:00:00' @@ -91,7 +112,7 @@ async def test_get_last_data_timestamp_error(redis_activity): test_data = {**metadata, 'workflow_name': 'test_pipeline', 'schedule_name': 'test_schedule'} redis_activity.send_notification = MagicMock() - redis_activity.get = MagicMock(side_effect=Exception('test')) + redis_activity.redis_repository.get.side_effect = Exception('test') try: await redis_activity.get_last_data_timestamp(test_data) @@ -122,13 +143,13 @@ async def test_put_last_data_timestamp_empty_dataframe(redis_activity): 'data': DataFrame(columns=['name', 'value', 'timestamp']).to_dict('records'), } - redis_activity.set = MagicMock() + redis_activity.redis_repository.set = MagicMock() result = await redis_activity.put_last_data_timestamp(test_data) assert result is None - redis_activity.set.assert_not_called() + redis_activity.redis_repository.set.assert_not_called() @pytest.mark.asyncio @@ -149,13 +170,13 @@ async def test_put_last_data_timestamp_not_empty_dataframe(redis_activity): 'data': data.to_dict('records'), } - redis_activity.set = MagicMock() + redis_activity.redis_repository.set = MagicMock() result = await redis_activity.put_last_data_timestamp(test_data) assert result == '2023-01-01 12:00:01' - redis_activity.set.assert_called_once_with( + redis_activity.redis_repository.set.assert_called_once_with( 'last_data_timestamp:test_pipeline:test_schedule', '2023-01-01 12:00:01', ttl=18000 ) @@ -177,7 +198,7 @@ async def test_put_last_data_timestamp_error(redis_activity): } redis_activity.send_notification = MagicMock() - redis_activity.set = MagicMock(side_effect=Exception('test')) + redis_activity.redis_repository.set = MagicMock(side_effect=Exception('test')) try: await redis_activity.put_last_data_timestamp(test_data) @@ -220,8 +241,8 @@ async def test_group_and_hold_data_new_key(redis_activity): } # Mock get to return None for new key - redis_activity.get = MagicMock(return_value=None) - redis_activity.set = MagicMock() + redis_activity.redis_repository.get = MagicMock(return_value=None) + redis_activity.redis_repository.set = MagicMock() # Call the method result = await redis_activity.group_and_hold_data(test_data) @@ -236,8 +257,8 @@ async def test_group_and_hold_data_new_key(redis_activity): assert result == expected_result # Verify set was called with correct arguments - redis_activity.set.assert_called_once() - args, kwargs = redis_activity.set.call_args + redis_activity.redis_repository.set.assert_called_once() + args, kwargs = redis_activity.redis_repository.set.call_args assert args[0] == 'held_data_test_pipeline_test_schedule' assert args[1] == {'sensor1': 25.5, 'sensor2': 30.0, 'timestamp': '2023-01-01 12:00:00'} assert kwargs['ttl'] == 3600 @@ -273,8 +294,8 @@ async def test_group_and_hold_data_update_existing_fill_missing(redis_activity): } # Mock get to return existing data - redis_activity.get = MagicMock(return_value=existing_data) - redis_activity.set = MagicMock() + redis_activity.redis_repository.get = MagicMock(return_value=existing_data) + redis_activity.redis_repository.set = MagicMock() # Call the method result = await redis_activity.group_and_hold_data(test_data) @@ -294,8 +315,8 @@ async def test_group_and_hold_data_update_existing_fill_missing(redis_activity): assert result == expected_result # Verify set was called with correct arguments - redis_activity.set.assert_called_once() - args, kwargs = redis_activity.set.call_args + redis_activity.redis_repository.set.assert_called_once() + args, kwargs = redis_activity.redis_repository.set.call_args assert args[0] == 'held_data_test_workflow_test_schedule' assert args[1] == { 'sensor1': 25.5, @@ -329,8 +350,8 @@ async def test_group_and_hold_data_with_none_values(redis_activity): } # Mock get to return None for new key - redis_activity.get = MagicMock(return_value=None) - redis_activity.set = MagicMock() + redis_activity.redis_repository.get = MagicMock(return_value=None) + redis_activity.redis_repository.set = MagicMock() # Call the method result = await redis_activity.group_and_hold_data(test_data) @@ -354,7 +375,7 @@ async def test_group_and_hold_data_empty_dataframe(redis_activity): 'fill_missing_tags': False, } - redis_activity.get = MagicMock(return_value=None) + redis_activity.redis_repository.get = MagicMock(return_value=None) # Call the method result = await redis_activity.group_and_hold_data(test_data) @@ -376,7 +397,7 @@ async def test_group_and_hold_data_error_get(redis_activity): 'fill_missing_tags': False, } - redis_activity.get = MagicMock(side_effect=Exception('test')) + redis_activity.redis_repository.get = MagicMock(side_effect=Exception('test')) redis_activity.send_notification = MagicMock() try: @@ -415,8 +436,8 @@ async def test_group_and_hold_data_error_set(redis_activity): existing_data = {'sensor1': 20.0, 'sensor2': 28.0, 'timestamp': '2023-01-01 11:00:00'} # Mock get to return existing data - redis_activity.get = MagicMock(return_value=existing_data) - redis_activity.set = MagicMock(side_effect=Exception('test')) + redis_activity.redis_repository.get = MagicMock(return_value=existing_data) + redis_activity.redis_repository.set = MagicMock(side_effect=Exception('test')) redis_activity.send_notification = MagicMock() try: @@ -429,7 +450,7 @@ async def test_group_and_hold_data_error_set(redis_activity): @pytest.mark.asyncio async def test_store_data_package(redis_activity): """Test store_data_package""" - redis_activity.set = MagicMock() + redis_activity.redis_repository.set = MagicMock() test_data = { **metadata, @@ -454,7 +475,7 @@ async def test_store_data_package(redis_activity): await redis_activity.store_data_package(test_data) - redis_activity.set.assert_called_once_with( + redis_activity.redis_repository.set.assert_called_once_with( ANY, {'data': test_data['data'], 'held_data': test_data['held_data']}, ttl=120 ) @@ -462,7 +483,7 @@ async def test_store_data_package(redis_activity): @pytest.mark.asyncio async def test_store_data_package_error(redis_activity): """Test store_data_package error""" - redis_activity.set = MagicMock(side_effect=ValueError('test')) + redis_activity.redis_repository.set = MagicMock(side_effect=ValueError('test')) redis_activity.send_notification = MagicMock() test_data = { diff --git a/tests/conftest.py b/tests/conftest.py new file mode 100644 index 0000000..a5aa843 --- /dev/null +++ b/tests/conftest.py @@ -0,0 +1,12 @@ +from unittest.mock import patch + +from pytest import fixture + + +@fixture(autouse=True, scope='session') +def sientia_monitoring_fixture(): + with patch( + 'sientia_do.observability.sientia_monitoring.SientiaMonitoring.__init__' + ) as mock_sientia_monitoring: + mock_sientia_monitoring.return_value = None + yield diff --git a/values.yaml b/values.yaml index 7371af9..85e47e5 100644 --- a/values.yaml +++ b/values.yaml @@ -150,7 +150,7 @@ env: - name: GITHUB_REPO_URL value: "git@github.com:Aignosi/sientia-dataops-scouter_temporal.git" - name: GITHUB_BRANCH - value: "fix/SIENTIAPDE-1314-ajustes-nas-camadas-de-monitoramento-do-sientia" + value: "feature/SIENTIAPDE-1325-adicionar-metricas-especificas-de-operacoes-externas" - name: PYTHON_APP value: "scouter.worker.worker" From d4a74ec5eb021b7c04244a34b222127a3a9d3daa Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Thu, 30 Oct 2025 16:49:55 -0300 Subject: [PATCH 03/10] SIENTIAPDE-1325 Update environment configuration and refactor activities to include metrics controller. Remove Kafka settings and adjust Redis and MongoDB initialization. Update tests to reflect changes in initialization and metrics tracking. --- .env.example | 10 +- init_port_forward.sh | 133 ++++++++++++++++++++++++++ pyproject.toml | 6 +- scouter/activities/activities.py | 10 +- scouter/activities/gates.py | 8 +- scouter/activities/mongodb.py | 10 +- scouter/activities/redis.py | 18 ++-- scouter/worker/worker.py | 2 +- scouter/workflow/scouter.py | 2 +- tests/activities/test_activities.py | 10 +- tests/activities/test_gates.py | 4 +- tests/activities/test_mongo.py | 19 ++-- tests/activities/test_redis.py | 65 ++++++------- tests/conftest.py | 12 --- tests/utils/test_connectors_config.py | 5 + values.yaml | 5 - 16 files changed, 225 insertions(+), 94 deletions(-) create mode 100755 init_port_forward.sh delete mode 100644 tests/conftest.py diff --git a/.env.example b/.env.example index 28055e7..a4ba811 100644 --- a/.env.example +++ b/.env.example @@ -6,14 +6,18 @@ POSTGRES_DB=sientia POSTGRES_MIN_CONNECTIONS=5 POSTGRES_MAX_CONNECTIONS=20 -KAFKA_BOOTSTRAP_SERVERS=localhost:9092 -KAFKA_POLLING_TIME=1000 +MONGODB_USERNAME="root" +MONGODB_PASSWORD="password" +MONGODB_URL="my-release-mongodb.mongodb.svc.cluster.local:27017" +MONGODB_DATABASE="sientia" REDIS_HOST=localhost REDIS_PORT=6379 +REDIS_USERNAME="user" +REDIS_PASSWORD="pass" TEMPORAL_HOST=localhost:7233 -TEMPORAL_NAMESPACE=default +TEMPORAL_NAMESPACE=scouter LOG_LEVEL=INFO diff --git a/init_port_forward.sh b/init_port_forward.sh new file mode 100755 index 0000000..e9c0e3d --- /dev/null +++ b/init_port_forward.sh @@ -0,0 +1,133 @@ +#!/bin/bash + +# Usage: +# 1) Edit the PORT_FORWARDS list below with entries of: +# +# 2) Run: ./init_port_forward.sh +# +# The script will start all port-forwards in the background and keep running +# until interrupted (Ctrl+C). On exit, it will clean up started port-forward processes. + +set -euo pipefail + +# Define your namespace/service/port combinations here +# Example entries: +# "default my-service 8080 80" +# "observability grafana 3000 3000" +PORT_FORWARDS=( + "mongodb my-release-mongodb 27017 27017" + "paradedb paradedb-rw 5432 5432" + "redis redis-master 6379 6379" + "temporal temporal-frontend 7233 7233" +) + +if [ ${#PORT_FORWARDS[@]} -eq 0 ]; then + echo "No port-forward entries defined. Edit PORT_FORWARDS in $(basename "$0")." + exit 1 +fi + +PIDS=() + +cleanup() { + echo "\nStopping port-forward processes..." + for pid in "${PIDS[@]}"; do + if kill -0 "$pid" >/dev/null 2>&1; then + kill "$pid" >/dev/null 2>&1 || true + fi + done +} + +trap cleanup EXIT INT TERM + +timestamp() { date '+%Y-%m-%d %H:%M:%S'; } + +# Allow overriding kubectl binary if needed +KUBECTL=${KUBECTL:-kubectl} + +is_port_free() { + local port="$1" + # Consider port free if nothing is listening locally on it + if command -v ss >/dev/null 2>&1; then + ! ss -ltn | awk '{print $4}' | grep -E "(^|:|\\])${port}$" >/dev/null 2>&1 + else + if command -v lsof >/dev/null 2>&1; then + ! lsof -tiTCP:"${port}" -sTCP:LISTEN >/dev/null 2>&1 + else + # Fallback: attempt to open a TCP connection; expect failure when nothing is listening + ! (exec 3<>"/dev/tcp/127.0.0.1/${port}") 2>/dev/null + fi + fi +} + +free_port_if_stuck() { + local port="$1" + # Try multiple tools to free a stuck listener (often old kubectl PF) + if command -v lsof >/dev/null 2>&1; then + local pids + pids=$(lsof -tiTCP:"${port}" -sTCP:LISTEN 2>/dev/null || true) + if [ -n "${pids}" ]; then + echo "[$(timestamp)] Found listeners on ${port}: ${pids}; terminating" + kill ${pids} 2>/dev/null || true + sleep 0.5 + fi + fi + if ! is_port_free "${port}"; then + if command -v fuser >/dev/null 2>&1; then + echo "[$(timestamp)] Forcing free of ${port} via fuser" + fuser -k "${port}/tcp" 2>/dev/null || true + sleep 0.5 + fi + fi +} + +run_port_forward() { + local namespace="$1" + local service_name="$2" + local local_port="$3" + local service_port="$4" + + # simple and robust supervisor loop with gentle backoff on failures + local delay=2 + local max_delay=20 + while true; do + free_port_if_stuck "${local_port}" + if ! is_port_free "${local_port}"; then + echo "[$(timestamp)] ns=${namespace} svc=${service_name} ${local_port}:${service_port} -> local port busy, retrying in 3s" + sleep 3 + continue + fi + + echo "[$(timestamp)] Starting port-forward: ns=${namespace} svc=${service_name} ${local_port}:${service_port}" + ${KUBECTL} -n "${namespace}" port-forward "svc/${service_name}" "${local_port}:${service_port}" \ + --address=127.0.0.1 --pod-running-timeout=2m --request-timeout=0 + rc=$? + + # If kubectl exits (e.g., connection reset by peer), wait a bit and retry + echo "[$(timestamp)] Port-forward exited (rc=${rc}): ns=${namespace} svc=${service_name} ${local_port}:${service_port}" + sleep "${delay}" + # Exponential backoff up to max_delay + if [ ${delay} -lt ${max_delay} ]; then + delay=$(( delay * 2 )) + if [ ${delay} -gt ${max_delay} ]; then + delay=${max_delay} + fi + fi + done +} + +ONLY_SERVICE_NAME="${ONLY_SERVICE_NAME:-}" + +for entry in "${PORT_FORWARDS[@]}"; do + read -r NAMESPACE SERVICE_NAME LOCAL_PORT SERVICE_PORT <<< "$entry" + if [ -n "${ONLY_SERVICE_NAME}" ] && [ "${SERVICE_NAME}" != "${ONLY_SERVICE_NAME}" ]; then + continue + fi + run_port_forward "${NAMESPACE}" "${SERVICE_NAME}" "${LOCAL_PORT}" "${SERVICE_PORT}" & + PIDS+=("$!") +done + +echo "All port-forwards started: ${#PIDS[@]} process(es). Press Ctrl+C to stop." + +# Do not exit the script if one port-forward fails; they self-restart +set +e +wait diff --git a/pyproject.toml b/pyproject.toml index 517fead..14d6e25 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -114,10 +114,6 @@ python_functions = ["test_*"] addopts = [ "-v", "--strict-markers", - "--cov=model_manager", - "--cov-report=term-missing", - "--cov-report=html", - "--cov-report=xml", ] markers = [ "asyncio: marks tests as async", @@ -126,7 +122,7 @@ markers = [ ] [tool.coverage.run] -source = ["model_manager"] +source = ["scouter"] omit = [ "*/tests/*", "*/venv/*", diff --git a/scouter/activities/activities.py b/scouter/activities/activities.py index 1dd20d8..19b1ae9 100644 --- a/scouter/activities/activities.py +++ b/scouter/activities/activities.py @@ -1,3 +1,4 @@ +from sientia_do.observability.metrics_controller import MetricsController from temporalio import workflow with workflow.unsafe.imports_passed_through(): @@ -50,6 +51,10 @@ class Activities(Postgres, Redis, Gates, MongoDB): logger (Logger): Logger instance for application logging notification_handler (NotificationHandler): Handler for system notifications """ + metrics_controller = MetricsController( + logger=logger, + ) + # Initialize Postgres Postgres.__init__( self, @@ -62,6 +67,7 @@ class Activities(Postgres, Redis, Gates, MongoDB): max_connections=postgres_config['max_connections'], logger=logger, notification_handler=notification_handler, + metrics_controller=metrics_controller, ) # Initialize Redis @@ -73,10 +79,11 @@ class Activities(Postgres, Redis, Gates, MongoDB): notification_handler=notification_handler, username=redis_config['username'], password=redis_config['password'], + metrics_controller=metrics_controller, ) # Initialize Gates - Gates.__init__(self, logger=logger, notification_handler=notification_handler) + Gates.__init__(self, logger=logger, notification_handler=notification_handler, metrics_controller=metrics_controller) # Initialize MongoDB MongoDB.__init__( @@ -85,6 +92,7 @@ class Activities(Postgres, Redis, Gates, MongoDB): database_name=mongodb_config['database_name'], logger=logger, notification_handler=notification_handler, + metrics_controller=metrics_controller, ) self.pod_id = getenv('HOSTNAME', 'localhost') diff --git a/scouter/activities/gates.py b/scouter/activities/gates.py index cb9ebf6..6c7ea32 100644 --- a/scouter/activities/gates.py +++ b/scouter/activities/gates.py @@ -1,5 +1,6 @@ from collections.abc import Hashable +from sientia_do.observability.metrics_controller import MetricsController from temporalio import activity, workflow with workflow.unsafe.imports_passed_through(): @@ -36,21 +37,22 @@ class Gates(SientiaMonitoring): ensure data integrity and enable flexible data processing workflows. """ - def __init__(self, logger: Logger, notification_handler: NotificationHandler): + def __init__(self, logger: Logger, notification_handler: NotificationHandler, metrics_controller: MetricsController): """ Initialize the Gates class with logging and notification services. Args: logger (Logger): Logger instance for operation logging notification_handler (NotificationHandler): Handler for system notifications + metrics_controller (MetricsController): Metrics controller instance """ - SientiaMonitoring.__init__(self, logger, notification_handler) + SientiaMonitoring.__init__(self, logger, notification_handler, metrics_controller) def close(self): """ Close the Gates class. """ - SientiaMonitoring.close(self) + SientiaMonitoring.shutdown(self) def apply_aggregation( self, values: DataFrame, aggr_function: str, metadata: dict[str, Any] diff --git a/scouter/activities/mongodb.py b/scouter/activities/mongodb.py index 55450eb..4c73af9 100644 --- a/scouter/activities/mongodb.py +++ b/scouter/activities/mongodb.py @@ -1,6 +1,6 @@ -from collections.abc import Hashable from datetime import UTC +from sientia_do.observability.metrics_controller import MetricsController from temporalio import activity, workflow with workflow.unsafe.imports_passed_through(): @@ -37,6 +37,7 @@ class MongoDB(SientiaMonitoring): database_name: str, logger: Logger, notification_handler: NotificationHandler, + metrics_controller: MetricsController, ): """ Initialize MongoDB connection and services. @@ -56,9 +57,10 @@ class MongoDB(SientiaMonitoring): database_name=database_name, logger=logger, notification_handler=notification_handler, + metrics_controller=metrics_controller, ) - SientiaMonitoring.__init__(self, logger=logger, notification_handler=notification_handler) + SientiaMonitoring.__init__(self, logger=logger, notification_handler=notification_handler, metrics_controller=metrics_controller) def close(self): """ @@ -77,7 +79,7 @@ class MongoDB(SientiaMonitoring): self.close() @activity.defn(name='load_latest_data') - async def load_latest_data(self, input_data: dict[str, Any]) -> dict[Hashable, Any]: + async def load_latest_data(self, input_data: dict[str, Any]) -> list[dict[str, Any]]: """ Load the latest data from MongoDB collection since a specified timestamp. @@ -116,7 +118,7 @@ class MongoDB(SientiaMonitoring): self.debug(f'Data filter: {data_filter}', metadata=metadata) - data = self.mongodb_repository.find( + data = await self.mongodb_repository.find( collection_name=collection_name, filters=data_filter, metadata=metadata, diff --git a/scouter/activities/redis.py b/scouter/activities/redis.py index e4513a2..bcb66f4 100644 --- a/scouter/activities/redis.py +++ b/scouter/activities/redis.py @@ -1,9 +1,9 @@ -from collections.abc import Hashable - +from sientia_do.observability.metrics_controller import MetricsController from temporalio import activity, workflow with workflow.unsafe.imports_passed_through(): import traceback + from collections.abc import Hashable from typing import Any from pandas import DataFrame @@ -38,6 +38,7 @@ class Redis(SientiaMonitoring): password: str, logger: Logger, notification_handler: NotificationHandler, + metrics_controller: MetricsController, ): """ Initialize Redis connection and services. @@ -50,7 +51,7 @@ class Redis(SientiaMonitoring): logger (Logger): Logger instance for operation logging notification_handler (NotificationHandler): Handler for system notifications """ - SientiaMonitoring.__init__(self, logger, notification_handler) + SientiaMonitoring.__init__(self, logger, notification_handler, metrics_controller) self.redis_repository = RedisRepository( host=host, port=port, @@ -58,6 +59,7 @@ class Redis(SientiaMonitoring): password=password, logger=logger, notification_handler=notification_handler, + metrics_controller=metrics_controller, ) def close(self): @@ -95,7 +97,7 @@ class Redis(SientiaMonitoring): self.info(f'Getting last data timestamp for {key}') try: - data_hold = self.redis_repository.get(key) + data_hold = await self.redis_repository.get(key) except Exception as e: self.send_notification( metadata=metadata, @@ -153,7 +155,7 @@ class Redis(SientiaMonitoring): self.info(f'Last collected timestamp to insert: {last_data_timestamp}', metadata=metadata) try: - self.redis_repository.set(key, last_data_timestamp, ttl=60 * 60 * 5) + await self.redis_repository.set(key, last_data_timestamp, ttl=60 * 60 * 5) except Exception as e: self.send_notification( metadata=metadata, @@ -206,7 +208,7 @@ class Redis(SientiaMonitoring): self.info(f'Getting held data for {key}') try: - data_hold = self.redis_repository.get(key) + data_hold = await self.redis_repository.get(key) except Exception as e: self.send_notification( metadata=metadata, @@ -250,7 +252,7 @@ class Redis(SientiaMonitoring): data['timestamp'].max() if not data.empty else data_hold['timestamp'] ) - self.redis_repository.set(key, data_hold, ttl=retention_time) + await self.redis_repository.set(key, data_hold, ttl=retention_time) data_hold_df = DataFrame(data_hold, index=[0]) data_hold_melted = data_hold_df.melt( @@ -296,7 +298,7 @@ class Redis(SientiaMonitoring): cache = {'data': data.to_dict(), 'held_data': held_data.to_dict()} try: - self.redis_repository.set(key, cache, ttl=120) + await self.redis_repository.set(key, cache, ttl=120) except Exception as e: self.send_notification( metadata=metadata, diff --git a/scouter/worker/worker.py b/scouter/worker/worker.py index c842761..a89875b 100644 --- a/scouter/worker/worker.py +++ b/scouter/worker/worker.py @@ -136,7 +136,7 @@ async def main(): except BaseException as e: # NOSONAR logger.custom_error( - 'An unhandled exception occurred: %s', e, exc_info=True, metadata=metadata + 'An unhandled exception occurred: %s', metadata=metadata ) finally: if notification_handler: diff --git a/scouter/workflow/scouter.py b/scouter/workflow/scouter.py index 5888627..52393e2 100644 --- a/scouter/workflow/scouter.py +++ b/scouter/workflow/scouter.py @@ -95,7 +95,7 @@ class Scouter: retry_policy=retry_policy, ) - if data == {}: + if not data: return await workflow.execute_activity_method( diff --git a/tests/activities/test_activities.py b/tests/activities/test_activities.py index 8864c5f..7d07ecb 100644 --- a/tests/activities/test_activities.py +++ b/tests/activities/test_activities.py @@ -12,7 +12,8 @@ from scouter.activities.redis import Redis @patch('scouter.activities.activities.Postgres.__init__') @patch('scouter.activities.activities.Redis.__init__') @patch('scouter.activities.activities.Gates.__init__') -def test___init__(mock_gates_init, mock_redis_init, mock_postgres_init, mock_mongodb_init): +@patch('scouter.activities.activities.MetricsController') +def test___init__(mock_metrics_controller, mock_gates_init, mock_redis_init, mock_postgres_init, mock_mongodb_init): postgres_config = { 'host': 'localhost', 'port': 5432, @@ -38,7 +39,7 @@ def test___init__(mock_gates_init, mock_redis_init, mock_postgres_init, mock_mon redis_config=redis_config, mongodb_config=mongodb_config, logger=logger, - notification_handler=notification_handler, + notification_handler=notification_handler ) assert isinstance(activities, Activities) @@ -58,6 +59,7 @@ def test___init__(mock_gates_init, mock_redis_init, mock_postgres_init, mock_mon max_connections=postgres_config['max_connections'], logger=logger, notification_handler=notification_handler, + metrics_controller=mock_metrics_controller.return_value, ) mock_redis_init.assert_called_once_with( @@ -68,6 +70,7 @@ def test___init__(mock_gates_init, mock_redis_init, mock_postgres_init, mock_mon password=redis_config['password'], logger=logger, notification_handler=notification_handler, + metrics_controller=mock_metrics_controller.return_value, ) mock_mongodb_init.assert_called_once_with( @@ -76,10 +79,11 @@ def test___init__(mock_gates_init, mock_redis_init, mock_postgres_init, mock_mon database_name=mongodb_config['database_name'], logger=logger, notification_handler=notification_handler, + metrics_controller=mock_metrics_controller.return_value, ) mock_gates_init.assert_called_once_with( - ANY, logger=logger, notification_handler=notification_handler + ANY, logger=logger, notification_handler=notification_handler, metrics_controller=mock_metrics_controller.return_value ) diff --git a/tests/activities/test_gates.py b/tests/activities/test_gates.py index 726c666..7e8f29e 100644 --- a/tests/activities/test_gates.py +++ b/tests/activities/test_gates.py @@ -14,10 +14,12 @@ def gates_fixture(): """Fixture to create a Gates instance with mocked dependencies.""" logger = Mock() notification_handler = MagicMock() - gates = Gates(logger=logger, notification_handler=notification_handler) + metrics_controller = MagicMock() + gates = Gates(logger=logger, notification_handler=notification_handler, metrics_controller=metrics_controller) gates.send_notification = MagicMock() gates.logger = logger gates.notification_handler = notification_handler + gates.metrics_controller = metrics_controller gates.pod_id = 'localhost' return gates diff --git a/tests/activities/test_mongo.py b/tests/activities/test_mongo.py index 268b02c..6b812c2 100644 --- a/tests/activities/test_mongo.py +++ b/tests/activities/test_mongo.py @@ -1,5 +1,5 @@ from datetime import datetime -from unittest.mock import ANY, MagicMock, patch +from unittest.mock import ANY, AsyncMock, MagicMock, patch from pytest import fixture, mark from sientia_do.notifications.models import NotificationLevel @@ -14,12 +14,13 @@ def test_mongodb___init__(mock_mongodb_repository): logger = MagicMock() notification_handler = MagicMock() - + metrics_controller = MagicMock() mongo = MongoDB( connection_string='mongodb://localhost:27017', database_name='test_db', logger=logger, notification_handler=notification_handler, + metrics_controller=metrics_controller, ) mock_mongodb_repository.assert_called_once_with( @@ -27,6 +28,7 @@ def test_mongodb___init__(mock_mongodb_repository): database_name='test_db', logger=logger, notification_handler=notification_handler, + metrics_controller=metrics_controller, ) assert mongo.mongodb_repository is not None @@ -42,11 +44,8 @@ def mongodb_activity(mock_mongodb_repository): database_name='test_db', logger=MagicMock(), notification_handler=MagicMock(), + metrics_controller=AsyncMock(), ) - - mongo.logger = MagicMock() - mongo.notification_handler = MagicMock() - mongo.pod_id = 'localhost' return mongo @@ -69,7 +68,7 @@ def test_del(mongodb_activity): async def test_load_latest_data_none_last_data_timestamp(mongodb_activity): """Test load_latest_data""" - mongodb_activity.mongodb_repository.find.return_value = [ + mongodb_activity.mongodb_repository.find = AsyncMock(return_value = [ { 'name': 'test1', 'value': 1, @@ -77,7 +76,7 @@ async def test_load_latest_data_none_last_data_timestamp(mongodb_activity): '2023-01-01 12:00:00.000000+0000', DATETIME_FORMAT_MS_WITH_TZ ), } - ] + ]) result = await mongodb_activity.load_latest_data( { @@ -106,7 +105,7 @@ async def test_load_latest_data_none_last_data_timestamp(mongodb_activity): async def test_load_latest_data_not_none_last_data_timestamp(mongodb_activity): """Test load_latest_data""" - mongodb_activity.mongodb_repository.find.return_value = [ + mongodb_activity.mongodb_repository.find = AsyncMock(return_value = [ { 'name': 'test1', 'value': 1, @@ -114,7 +113,7 @@ async def test_load_latest_data_not_none_last_data_timestamp(mongodb_activity): '2023-01-01 12:00:00.000000+0000', DATETIME_FORMAT_MS_WITH_TZ ), } - ] + ]) result = await mongodb_activity.load_latest_data( { diff --git a/tests/activities/test_redis.py b/tests/activities/test_redis.py index 0daf867..3ec93c8 100644 --- a/tests/activities/test_redis.py +++ b/tests/activities/test_redis.py @@ -1,5 +1,5 @@ from datetime import datetime -from unittest.mock import ANY, MagicMock, patch +from unittest.mock import ANY, AsyncMock, MagicMock, patch import numpy as np import pytest @@ -15,6 +15,7 @@ from scouter.activities.redis import Redis def redis_activity(mock_redis_repository): logger = MagicMock() notification_handler = MagicMock(spec=NotificationHandler) + metrics_controller = MagicMock() activity = Redis( host='localhost', port=6379, @@ -22,6 +23,7 @@ def redis_activity(mock_redis_repository): notification_handler=notification_handler, username='test', password='test', + metrics_controller=metrics_controller, ) activity.redis_client = MagicMock() @@ -46,6 +48,7 @@ def test_redis_initialization(mock_redis_repository): """Test Redis activity initialization""" logger = MagicMock() notification_handler = MagicMock(spec=NotificationHandler) + metrics_controller = MagicMock() activity = Redis( host='localhost', port=6379, @@ -53,6 +56,7 @@ def test_redis_initialization(mock_redis_repository): notification_handler=notification_handler, username='test', password='test', + metrics_controller=metrics_controller, ) mock_redis_repository.assert_called_once_with( host='localhost', @@ -61,12 +65,9 @@ def test_redis_initialization(mock_redis_repository): password='test', logger=logger, notification_handler=notification_handler, + metrics_controller=metrics_controller, ) - activity.logger = logger - activity.notification_handler = notification_handler - activity.pod_id = 'localhost' - assert activity.redis_repository is not None @@ -79,7 +80,7 @@ async def test_get_last_data_timestamp_none(redis_activity): 'schedule_name': 'test_schedule', } - redis_activity.redis_repository.get.return_value = None + redis_activity.redis_repository.get = AsyncMock(return_value = None) result = await redis_activity.get_last_data_timestamp(test_data) @@ -95,7 +96,7 @@ async def test_get_last_data_timestamp_not_none(redis_activity): """Test get_last_data_timestamp""" test_data = {**metadata, 'workflow_name': 'test_pipeline', 'schedule_name': 'test_schedule'} - redis_activity.redis_repository.get.return_value = '2023-01-01 12:00:00' + redis_activity.redis_repository.get = AsyncMock(return_value = '2023-01-01 12:00:00') result = await redis_activity.get_last_data_timestamp(test_data) @@ -170,7 +171,7 @@ async def test_put_last_data_timestamp_not_empty_dataframe(redis_activity): 'data': data.to_dict('records'), } - redis_activity.redis_repository.set = MagicMock() + redis_activity.redis_repository.set = AsyncMock() result = await redis_activity.put_last_data_timestamp(test_data) @@ -241,8 +242,8 @@ async def test_group_and_hold_data_new_key(redis_activity): } # Mock get to return None for new key - redis_activity.redis_repository.get = MagicMock(return_value=None) - redis_activity.redis_repository.set = MagicMock() + redis_activity.redis_repository.get = AsyncMock(return_value=None) + redis_activity.redis_repository.set = AsyncMock() # Call the method result = await redis_activity.group_and_hold_data(test_data) @@ -257,12 +258,9 @@ async def test_group_and_hold_data_new_key(redis_activity): assert result == expected_result # Verify set was called with correct arguments - redis_activity.redis_repository.set.assert_called_once() - args, kwargs = redis_activity.redis_repository.set.call_args - assert args[0] == 'held_data_test_pipeline_test_schedule' - assert args[1] == {'sensor1': 25.5, 'sensor2': 30.0, 'timestamp': '2023-01-01 12:00:00'} - assert kwargs['ttl'] == 3600 - + redis_activity.redis_repository.set.assert_called_once_with( + 'held_data_test_pipeline_test_schedule', {'sensor1': 25.5, 'sensor2': 30.0, 'timestamp': '2023-01-01 12:00:00'}, ttl=3600 + ) @pytest.mark.asyncio async def test_group_and_hold_data_update_existing_fill_missing(redis_activity): @@ -294,8 +292,8 @@ async def test_group_and_hold_data_update_existing_fill_missing(redis_activity): } # Mock get to return existing data - redis_activity.redis_repository.get = MagicMock(return_value=existing_data) - redis_activity.redis_repository.set = MagicMock() + redis_activity.redis_repository.get = AsyncMock(return_value=existing_data) + redis_activity.redis_repository.set = AsyncMock() # Call the method result = await redis_activity.group_and_hold_data(test_data) @@ -315,17 +313,10 @@ async def test_group_and_hold_data_update_existing_fill_missing(redis_activity): assert result == expected_result # Verify set was called with correct arguments - redis_activity.redis_repository.set.assert_called_once() - args, kwargs = redis_activity.redis_repository.set.call_args - assert args[0] == 'held_data_test_workflow_test_schedule' - assert args[1] == { - 'sensor1': 25.5, - 'sensor2': 28.0, - 'sensor3': 42.0, - 'sensor4': None, - 'timestamp': '2023-01-01 12:00:00', - } - assert kwargs['ttl'] == 3600 + redis_activity.redis_repository.set.assert_called_once_with( + 'held_data_test_workflow_test_schedule', {'sensor1': 25.5, 'sensor2': 28.0, 'sensor3': 42.0, 'sensor4': None, 'timestamp': '2023-01-01 12:00:00'}, ttl=3600 + ) + @pytest.mark.asyncio @@ -350,8 +341,8 @@ async def test_group_and_hold_data_with_none_values(redis_activity): } # Mock get to return None for new key - redis_activity.redis_repository.get = MagicMock(return_value=None) - redis_activity.redis_repository.set = MagicMock() + redis_activity.redis_repository.get = AsyncMock(return_value=None) + redis_activity.redis_repository.set = AsyncMock() # Call the method result = await redis_activity.group_and_hold_data(test_data) @@ -375,7 +366,7 @@ async def test_group_and_hold_data_empty_dataframe(redis_activity): 'fill_missing_tags': False, } - redis_activity.redis_repository.get = MagicMock(return_value=None) + redis_activity.redis_repository.get = AsyncMock(return_value=None) # Call the method result = await redis_activity.group_and_hold_data(test_data) @@ -397,7 +388,7 @@ async def test_group_and_hold_data_error_get(redis_activity): 'fill_missing_tags': False, } - redis_activity.redis_repository.get = MagicMock(side_effect=Exception('test')) + redis_activity.redis_repository.get = AsyncMock(side_effect=Exception('test')) redis_activity.send_notification = MagicMock() try: @@ -436,8 +427,8 @@ async def test_group_and_hold_data_error_set(redis_activity): existing_data = {'sensor1': 20.0, 'sensor2': 28.0, 'timestamp': '2023-01-01 11:00:00'} # Mock get to return existing data - redis_activity.redis_repository.get = MagicMock(return_value=existing_data) - redis_activity.redis_repository.set = MagicMock(side_effect=Exception('test')) + redis_activity.redis_repository.get = AsyncMock(return_value=existing_data) + redis_activity.redis_repository.set = AsyncMock(side_effect=Exception('test')) redis_activity.send_notification = MagicMock() try: @@ -450,7 +441,7 @@ async def test_group_and_hold_data_error_set(redis_activity): @pytest.mark.asyncio async def test_store_data_package(redis_activity): """Test store_data_package""" - redis_activity.redis_repository.set = MagicMock() + redis_activity.redis_repository.set = AsyncMock() test_data = { **metadata, @@ -483,7 +474,7 @@ async def test_store_data_package(redis_activity): @pytest.mark.asyncio async def test_store_data_package_error(redis_activity): """Test store_data_package error""" - redis_activity.redis_repository.set = MagicMock(side_effect=ValueError('test')) + redis_activity.redis_repository.set = AsyncMock(side_effect=ValueError('test')) redis_activity.send_notification = MagicMock() test_data = { diff --git a/tests/conftest.py b/tests/conftest.py deleted file mode 100644 index a5aa843..0000000 --- a/tests/conftest.py +++ /dev/null @@ -1,12 +0,0 @@ -from unittest.mock import patch - -from pytest import fixture - - -@fixture(autouse=True, scope='session') -def sientia_monitoring_fixture(): - with patch( - 'sientia_do.observability.sientia_monitoring.SientiaMonitoring.__init__' - ) as mock_sientia_monitoring: - mock_sientia_monitoring.return_value = None - yield diff --git a/tests/utils/test_connectors_config.py b/tests/utils/test_connectors_config.py index 61d2b0a..0eb7012 100644 --- a/tests/utils/test_connectors_config.py +++ b/tests/utils/test_connectors_config.py @@ -122,6 +122,11 @@ def test_build_redis_config_with_env_vars(): def test_build_mongodb_config_defaults(): """Test that build_mongodb_config returns default values when no env vars are set""" + os.environ['MONGODB_URL'] = 'localhost:27017' + os.environ['MONGODB_DATABASE_NAME'] = 'sientia' + os.environ['MONGODB_USERNAME'] = 'sientia' + os.environ['MONGODB_PASSWORD'] = 'sientia' + config = build_mongodb_config() assert config == { diff --git a/values.yaml b/values.yaml index 85e47e5..804b505 100644 --- a/values.yaml +++ b/values.yaml @@ -170,11 +170,6 @@ env: - name: POSTGRES_MAX_CONNECTIONS value: "40" - - name: KAFKA_BOOTSTRAP_SERVERS - value: "kafka.kafka.svc.cluster.local:9092" - - name: KAFKA_POLLING_TIME - value: "10000" - - name: REDIS_HOST value: "redis-master.redis.svc.cluster.local" - name: REDIS_PORT From 8350a3a0c3c24e97dd7bba7573f4fecb439557d1 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Thu, 30 Oct 2025 16:50:39 -0300 Subject: [PATCH 04/10] SIENTIAPDE-1325 Refactor initialization of Activities, Gates, and MongoDB classes for improved readability by using multi-line argument formatting. Update related tests to match the new initialization style. --- scouter/activities/activities.py | 7 ++++- scouter/activities/gates.py | 7 ++++- scouter/activities/mongodb.py | 7 ++++- scouter/worker/worker.py | 6 ++-- tests/activities/test_activities.py | 11 ++++++-- tests/activities/test_gates.py | 6 +++- tests/activities/test_mongo.py | 40 +++++++++++++++------------ tests/activities/test_redis.py | 20 ++++++++++---- tests/utils/test_connectors_config.py | 2 +- 9 files changed, 71 insertions(+), 35 deletions(-) diff --git a/scouter/activities/activities.py b/scouter/activities/activities.py index 19b1ae9..2f3f8f9 100644 --- a/scouter/activities/activities.py +++ b/scouter/activities/activities.py @@ -83,7 +83,12 @@ class Activities(Postgres, Redis, Gates, MongoDB): ) # Initialize Gates - Gates.__init__(self, logger=logger, notification_handler=notification_handler, metrics_controller=metrics_controller) + Gates.__init__( + self, + logger=logger, + notification_handler=notification_handler, + metrics_controller=metrics_controller, + ) # Initialize MongoDB MongoDB.__init__( diff --git a/scouter/activities/gates.py b/scouter/activities/gates.py index 6c7ea32..f539b4f 100644 --- a/scouter/activities/gates.py +++ b/scouter/activities/gates.py @@ -37,7 +37,12 @@ class Gates(SientiaMonitoring): ensure data integrity and enable flexible data processing workflows. """ - def __init__(self, logger: Logger, notification_handler: NotificationHandler, metrics_controller: MetricsController): + def __init__( + self, + logger: Logger, + notification_handler: NotificationHandler, + metrics_controller: MetricsController, + ): """ Initialize the Gates class with logging and notification services. diff --git a/scouter/activities/mongodb.py b/scouter/activities/mongodb.py index 4c73af9..4e61d86 100644 --- a/scouter/activities/mongodb.py +++ b/scouter/activities/mongodb.py @@ -60,7 +60,12 @@ class MongoDB(SientiaMonitoring): metrics_controller=metrics_controller, ) - SientiaMonitoring.__init__(self, logger=logger, notification_handler=notification_handler, metrics_controller=metrics_controller) + SientiaMonitoring.__init__( + self, + logger=logger, + notification_handler=notification_handler, + metrics_controller=metrics_controller, + ) def close(self): """ diff --git a/scouter/worker/worker.py b/scouter/worker/worker.py index a89875b..75e6715 100644 --- a/scouter/worker/worker.py +++ b/scouter/worker/worker.py @@ -134,10 +134,8 @@ async def main(): try: await asyncio.gather(*handlers) - except BaseException as e: # NOSONAR - logger.custom_error( - 'An unhandled exception occurred: %s', metadata=metadata - ) + except BaseException: # NOSONAR + logger.custom_error('An unhandled exception occurred: %s', metadata=metadata) finally: if notification_handler: notification_handler.shutdown() diff --git a/tests/activities/test_activities.py b/tests/activities/test_activities.py index 7d07ecb..bc31dee 100644 --- a/tests/activities/test_activities.py +++ b/tests/activities/test_activities.py @@ -13,7 +13,9 @@ from scouter.activities.redis import Redis @patch('scouter.activities.activities.Redis.__init__') @patch('scouter.activities.activities.Gates.__init__') @patch('scouter.activities.activities.MetricsController') -def test___init__(mock_metrics_controller, mock_gates_init, mock_redis_init, mock_postgres_init, mock_mongodb_init): +def test___init__( + mock_metrics_controller, mock_gates_init, mock_redis_init, mock_postgres_init, mock_mongodb_init +): postgres_config = { 'host': 'localhost', 'port': 5432, @@ -39,7 +41,7 @@ def test___init__(mock_metrics_controller, mock_gates_init, mock_redis_init, moc redis_config=redis_config, mongodb_config=mongodb_config, logger=logger, - notification_handler=notification_handler + notification_handler=notification_handler, ) assert isinstance(activities, Activities) @@ -83,7 +85,10 @@ def test___init__(mock_metrics_controller, mock_gates_init, mock_redis_init, moc ) mock_gates_init.assert_called_once_with( - ANY, logger=logger, notification_handler=notification_handler, metrics_controller=mock_metrics_controller.return_value + ANY, + logger=logger, + notification_handler=notification_handler, + metrics_controller=mock_metrics_controller.return_value, ) diff --git a/tests/activities/test_gates.py b/tests/activities/test_gates.py index 7e8f29e..814c14e 100644 --- a/tests/activities/test_gates.py +++ b/tests/activities/test_gates.py @@ -15,7 +15,11 @@ def gates_fixture(): logger = Mock() notification_handler = MagicMock() metrics_controller = MagicMock() - gates = Gates(logger=logger, notification_handler=notification_handler, metrics_controller=metrics_controller) + gates = Gates( + logger=logger, + notification_handler=notification_handler, + metrics_controller=metrics_controller, + ) gates.send_notification = MagicMock() gates.logger = logger gates.notification_handler = notification_handler diff --git a/tests/activities/test_mongo.py b/tests/activities/test_mongo.py index 6b812c2..bb8e9b2 100644 --- a/tests/activities/test_mongo.py +++ b/tests/activities/test_mongo.py @@ -68,15 +68,17 @@ def test_del(mongodb_activity): async def test_load_latest_data_none_last_data_timestamp(mongodb_activity): """Test load_latest_data""" - mongodb_activity.mongodb_repository.find = AsyncMock(return_value = [ - { - 'name': 'test1', - 'value': 1, - 'inserted_at': datetime.strptime( - '2023-01-01 12:00:00.000000+0000', DATETIME_FORMAT_MS_WITH_TZ - ), - } - ]) + mongodb_activity.mongodb_repository.find = AsyncMock( + return_value=[ + { + 'name': 'test1', + 'value': 1, + 'inserted_at': datetime.strptime( + '2023-01-01 12:00:00.000000+0000', DATETIME_FORMAT_MS_WITH_TZ + ), + } + ] + ) result = await mongodb_activity.load_latest_data( { @@ -105,15 +107,17 @@ async def test_load_latest_data_none_last_data_timestamp(mongodb_activity): async def test_load_latest_data_not_none_last_data_timestamp(mongodb_activity): """Test load_latest_data""" - mongodb_activity.mongodb_repository.find = AsyncMock(return_value = [ - { - 'name': 'test1', - 'value': 1, - 'inserted_at': datetime.strptime( - '2023-01-01 12:00:00.000000+0000', DATETIME_FORMAT_MS_WITH_TZ - ), - } - ]) + mongodb_activity.mongodb_repository.find = AsyncMock( + return_value=[ + { + 'name': 'test1', + 'value': 1, + 'inserted_at': datetime.strptime( + '2023-01-01 12:00:00.000000+0000', DATETIME_FORMAT_MS_WITH_TZ + ), + } + ] + ) result = await mongodb_activity.load_latest_data( { diff --git a/tests/activities/test_redis.py b/tests/activities/test_redis.py index 3ec93c8..80b08fc 100644 --- a/tests/activities/test_redis.py +++ b/tests/activities/test_redis.py @@ -80,7 +80,7 @@ async def test_get_last_data_timestamp_none(redis_activity): 'schedule_name': 'test_schedule', } - redis_activity.redis_repository.get = AsyncMock(return_value = None) + redis_activity.redis_repository.get = AsyncMock(return_value=None) result = await redis_activity.get_last_data_timestamp(test_data) @@ -96,7 +96,7 @@ async def test_get_last_data_timestamp_not_none(redis_activity): """Test get_last_data_timestamp""" test_data = {**metadata, 'workflow_name': 'test_pipeline', 'schedule_name': 'test_schedule'} - redis_activity.redis_repository.get = AsyncMock(return_value = '2023-01-01 12:00:00') + redis_activity.redis_repository.get = AsyncMock(return_value='2023-01-01 12:00:00') result = await redis_activity.get_last_data_timestamp(test_data) @@ -259,9 +259,12 @@ async def test_group_and_hold_data_new_key(redis_activity): # Verify set was called with correct arguments redis_activity.redis_repository.set.assert_called_once_with( - 'held_data_test_pipeline_test_schedule', {'sensor1': 25.5, 'sensor2': 30.0, 'timestamp': '2023-01-01 12:00:00'}, ttl=3600 + 'held_data_test_pipeline_test_schedule', + {'sensor1': 25.5, 'sensor2': 30.0, 'timestamp': '2023-01-01 12:00:00'}, + ttl=3600, ) + @pytest.mark.asyncio async def test_group_and_hold_data_update_existing_fill_missing(redis_activity): """Test updating existing data with group_and_hold_data""" @@ -314,11 +317,18 @@ async def test_group_and_hold_data_update_existing_fill_missing(redis_activity): # Verify set was called with correct arguments redis_activity.redis_repository.set.assert_called_once_with( - 'held_data_test_workflow_test_schedule', {'sensor1': 25.5, 'sensor2': 28.0, 'sensor3': 42.0, 'sensor4': None, 'timestamp': '2023-01-01 12:00:00'}, ttl=3600 + 'held_data_test_workflow_test_schedule', + { + 'sensor1': 25.5, + 'sensor2': 28.0, + 'sensor3': 42.0, + 'sensor4': None, + 'timestamp': '2023-01-01 12:00:00', + }, + ttl=3600, ) - @pytest.mark.asyncio async def test_group_and_hold_data_with_none_values(redis_activity): """Test handling of None values in group_and_hold_data""" diff --git a/tests/utils/test_connectors_config.py b/tests/utils/test_connectors_config.py index 0eb7012..d57b335 100644 --- a/tests/utils/test_connectors_config.py +++ b/tests/utils/test_connectors_config.py @@ -126,7 +126,7 @@ def test_build_mongodb_config_defaults(): os.environ['MONGODB_DATABASE_NAME'] = 'sientia' os.environ['MONGODB_USERNAME'] = 'sientia' os.environ['MONGODB_PASSWORD'] = 'sientia' - + config = build_mongodb_config() assert config == { From 43512891113d58807e3b54a428ada688efa16c2c Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Mon, 3 Nov 2025 14:57:44 -0300 Subject: [PATCH 05/10] SIENTIAPDE-1325 Refactor notification methods in Gates, MongoDB, and Redis classes to use asynchronous send_notification_async. Update related tests to ensure proper async handling and verification of notifications. --- scouter/activities/gates.py | 12 +++++------ scouter/activities/mongodb.py | 2 +- scouter/activities/redis.py | 10 ++++----- tests/activities/test_gates.py | 37 ++++++++++++++++++++++------------ tests/activities/test_mongo.py | 4 +++- tests/activities/test_redis.py | 28 ++++++++++++++++--------- 6 files changed, 57 insertions(+), 36 deletions(-) diff --git a/scouter/activities/gates.py b/scouter/activities/gates.py index f539b4f..b5ab579 100644 --- a/scouter/activities/gates.py +++ b/scouter/activities/gates.py @@ -59,7 +59,7 @@ class Gates(SientiaMonitoring): """ SientiaMonitoring.shutdown(self) - def apply_aggregation( + async def apply_aggregation( self, values: DataFrame, aggr_function: str, metadata: dict[str, Any] ) -> float | None | str: """ @@ -106,7 +106,7 @@ class Gates(SientiaMonitoring): if aggr_function in aggregation_map: return aggregation_map[aggr_function](clean_values) else: - self.send_notification( + await self.send_notification_async( metadata=metadata, notification_id='AGGREGATION_ISSUES', message=f'Invalid aggregation function: {aggr_function}', @@ -165,7 +165,7 @@ class Gates(SientiaMonitoring): # Get the latest timestamp (last row since data is sorted) latest_timestamp = group['timestamp'].iloc[-1] - aggr_value = self.apply_aggregation(group, aggr_function, metadata) + aggr_value = await self.apply_aggregation(group, aggr_function, metadata) if aggr_value == 'continue': continue @@ -197,7 +197,7 @@ class Gates(SientiaMonitoring): except Exception as e: trace = traceback.format_exc() - self.send_notification( + await self.send_notification_async( metadata=metadata, notification_id='AGGREGATION_ISSUES', message=f'Error aggregating data: {e}', @@ -255,7 +255,7 @@ class Gates(SientiaMonitoring): except Exception as e: trace = traceback.format_exc() - self.send_notification( + await self.send_notification_async( metadata=metadata, notification_id='DATA_QUALITY_GATE_ISSUES', message=f'Error applying filter {filter_name}: {e}', @@ -273,7 +273,7 @@ class Gates(SientiaMonitoring): message = f'{len(filtered_data)} rows has quality issues: {filter_name}: {policy}' attachment = filtered_data.to_string() - self.send_notification( + await self.send_notification_async( metadata=metadata, notification_id=f'DATA_QUALITY_GATE_ISSUES__{filter_name}', message=message, diff --git a/scouter/activities/mongodb.py b/scouter/activities/mongodb.py index 4e61d86..55d269e 100644 --- a/scouter/activities/mongodb.py +++ b/scouter/activities/mongodb.py @@ -143,7 +143,7 @@ class MongoDB(SientiaMonitoring): return data except Exception as e: trace = traceback.format_exc() - self.send_notification( + await self.send_notification_async( metadata=metadata, notification_id='MONGO_LOAD_ERROR', message=f'Error loading data from MongoDB: {e}', diff --git a/scouter/activities/redis.py b/scouter/activities/redis.py index bcb66f4..6e2b885 100644 --- a/scouter/activities/redis.py +++ b/scouter/activities/redis.py @@ -99,7 +99,7 @@ class Redis(SientiaMonitoring): try: data_hold = await self.redis_repository.get(key) except Exception as e: - self.send_notification( + await self.send_notification_async( metadata=metadata, notification_id='REDIS_GET_ERROR', message=f'Error getting last data timestamp: {e}', @@ -157,7 +157,7 @@ class Redis(SientiaMonitoring): try: await self.redis_repository.set(key, last_data_timestamp, ttl=60 * 60 * 5) except Exception as e: - self.send_notification( + await self.send_notification_async( metadata=metadata, notification_id='REDIS_SET_ERROR', message=f'Error setting last data timestamp: {e}', @@ -210,7 +210,7 @@ class Redis(SientiaMonitoring): try: data_hold = await self.redis_repository.get(key) except Exception as e: - self.send_notification( + await self.send_notification_async( metadata=metadata, notification_id='REDIS_GET_ERROR', message=f'Error getting held data: {e}', @@ -262,7 +262,7 @@ class Redis(SientiaMonitoring): data_hold_melted.reset_index(drop=True, inplace=True) except Exception as e: - self.send_notification( + await self.send_notification_async( metadata=metadata, notification_id='REDIS_SET_ERROR', message=f'Error setting held data: {e}', @@ -300,7 +300,7 @@ class Redis(SientiaMonitoring): try: await self.redis_repository.set(key, cache, ttl=120) except Exception as e: - self.send_notification( + await self.send_notification_async( metadata=metadata, notification_id='REDIS_SET_ERROR', message=f'Error setting data package: {e}', diff --git a/tests/activities/test_gates.py b/tests/activities/test_gates.py index 814c14e..d02ed0a 100644 --- a/tests/activities/test_gates.py +++ b/tests/activities/test_gates.py @@ -1,5 +1,5 @@ from typing import Any -from unittest.mock import ANY, MagicMock, Mock, call, patch +from unittest.mock import ANY, AsyncMock, MagicMock, Mock, call, patch import numpy as np import pandas as pd @@ -21,6 +21,9 @@ def gates_fixture(): metrics_controller=metrics_controller, ) gates.send_notification = MagicMock() + gates.send_notification_async = AsyncMock() + gates.emit_metric = AsyncMock() + gates.logger = logger gates.notification_handler = notification_handler gates.metrics_controller = metrics_controller @@ -38,6 +41,13 @@ metadata = { } +@patch('scouter.activities.gates.SientiaMonitoring') +def test_close(mock_sientia_monitoring, gates_fixture): + """Test close method.""" + gates_fixture.close() + mock_sientia_monitoring.shutdown.assert_called_once() + + @pytest.mark.asyncio async def test_data_quality_gate_with_null_values_filter_discard(gates_fixture): """Test data_quality_gate with NULL_VALUES_FILTER and DISCARD policy.""" @@ -64,7 +74,7 @@ async def test_data_quality_gate_with_null_values_filter_discard(gates_fixture): # Verify assert len(result['tag']) == 2 assert 'tag2' not in result['tag'] - gates_fixture.send_notification.assert_called_once() + gates_fixture.send_notification_async.assert_called_once() @pytest.mark.asyncio @@ -97,7 +107,7 @@ async def test_data_quality_gate_with_out_of_bounds_filter_keep(gates_fixture): # Verify data is kept but notification is sent assert len(result['tag']) == 3 # All rows kept - gates_fixture.send_notification.assert_called_once() + gates_fixture.send_notification_async.assert_called_once() @pytest.mark.asyncio @@ -134,7 +144,7 @@ async def test_data_quality_gate_with_multiple_filters(gates_fixture): 'timestamp': {0: '2023-01-01', 3: '2023-01-04'}, } # Should be called twice (once for each filter) - assert gates_fixture.send_notification.call_count == 2 + assert gates_fixture.send_notification_async.call_count == 2 @pytest.mark.asyncio @@ -182,8 +192,8 @@ async def test_data_quality_gate_with_filter_error(gates_fixture): # Verify error notification is sent and data is unchanged assert len(result['tag']) == 1 - gates_fixture.send_notification.assert_called_once() - call_args = gates_fixture.send_notification.call_args[1] + gates_fixture.send_notification_async.assert_called_once() + call_args = gates_fixture.send_notification_async.call_args[1] assert call_args['notification_id'] == 'DATA_QUALITY_GATE_ISSUES' assert call_args['level'] == NotificationLevel.ERROR assert 'Filter error' in call_args['message'] @@ -246,16 +256,17 @@ async def test_data_quality_gate_with_no_filters(gates_fixture): (pd.DataFrame({'value': [1.0, 2.0]}), 'invalid', 'continue'), ], ) -def test_apply_aggregation(gates_fixture, group_data, aggr_function, expected_result): +@pytest.mark.asyncio +async def test_apply_aggregation(gates_fixture, group_data, aggr_function, expected_result): """Test apply_aggregation method with various scenarios.""" - result = gates_fixture.apply_aggregation(group_data, aggr_function, metadata) + result = await gates_fixture.apply_aggregation(group_data, aggr_function, metadata) assert result == expected_result # Check notification was sent for invalid function if aggr_function == 'invalid': - gates_fixture.send_notification.assert_called_once() + gates_fixture.send_notification_async.assert_called_once() else: - gates_fixture.send_notification.assert_not_called() + gates_fixture.send_notification_async.assert_not_called() @pytest.mark.asyncio @@ -297,7 +308,7 @@ async def test_aggregate_data(gates_fixture): @pytest.mark.asyncio async def test_aggregate_data_with_continue(gates_fixture): - gates_fixture.apply_aggregation = MagicMock(return_value='continue') + gates_fixture.apply_aggregation = AsyncMock(return_value='continue') input_data = { 'data': [ @@ -324,7 +335,7 @@ async def test_aggregate_data_with_continue(gates_fixture): # Verify assert result == expected_result - gates_fixture.send_notification.assert_not_called() + gates_fixture.send_notification_async.assert_not_called() @pytest.mark.asyncio @@ -352,7 +363,7 @@ async def test_aggregate_data_raise_exception(gates_fixture): await gates_fixture.aggregate_data(input_data) except Exception as e: assert str(e) == 'Test exception' - gates_fixture.send_notification.assert_called_once_with( + gates_fixture.send_notification_async.assert_called_once_with( metadata=metadata['metadata'], notification_id='AGGREGATION_ISSUES', message='Error aggregating data: Test exception', diff --git a/tests/activities/test_mongo.py b/tests/activities/test_mongo.py index bb8e9b2..d0b00d7 100644 --- a/tests/activities/test_mongo.py +++ b/tests/activities/test_mongo.py @@ -153,6 +153,8 @@ async def test_load_latest_data_error(mongodb_activity): """Test load_latest_data""" mongodb_activity.mongodb_repository.find.side_effect = Exception('test') mongodb_activity.send_notification = MagicMock() + mongodb_activity.send_notification_async = AsyncMock() + mongodb_activity.emit_metric = AsyncMock() try: await mongodb_activity.load_latest_data( @@ -165,7 +167,7 @@ async def test_load_latest_data_error(mongodb_activity): except Exception as e: assert str(e) == 'test' - mongodb_activity.send_notification.assert_called_once_with( + mongodb_activity.send_notification_async.assert_called_once_with( metadata={'workflow_name': 'test_pipeline', 'schedule_name': 'test_schedule'}, notification_id='MONGO_LOAD_ERROR', message='Error loading data from MongoDB: test', diff --git a/tests/activities/test_redis.py b/tests/activities/test_redis.py index 80b08fc..8c2b7ec 100644 --- a/tests/activities/test_redis.py +++ b/tests/activities/test_redis.py @@ -43,6 +43,14 @@ metadata = { } +@patch('scouter.activities.redis.SientiaMonitoring') +def test_close(mock_sientia_monitoring, redis_activity): + """Test close method.""" + redis_activity.close() + redis_activity.redis_repository.close.assert_called_once() + mock_sientia_monitoring.shutdown.assert_called_once() + + @patch('scouter.activities.redis.RedisRepository') def test_redis_initialization(mock_redis_repository): """Test Redis activity initialization""" @@ -112,7 +120,7 @@ async def test_get_last_data_timestamp_error(redis_activity): """Test get_last_data_timestamp error""" test_data = {**metadata, 'workflow_name': 'test_pipeline', 'schedule_name': 'test_schedule'} - redis_activity.send_notification = MagicMock() + redis_activity.send_notification_async = AsyncMock() redis_activity.redis_repository.get.side_effect = Exception('test') try: @@ -121,7 +129,7 @@ async def test_get_last_data_timestamp_error(redis_activity): except Exception as e: assert str(e) == 'test' - redis_activity.send_notification.assert_called_once_with( + redis_activity.send_notification_async.assert_called_once_with( metadata=metadata['metadata'], notification_id='REDIS_GET_ERROR', message='Error getting last data timestamp: test', @@ -198,8 +206,8 @@ async def test_put_last_data_timestamp_error(redis_activity): ).to_dict('records'), } - redis_activity.send_notification = MagicMock() - redis_activity.redis_repository.set = MagicMock(side_effect=Exception('test')) + redis_activity.send_notification_async = AsyncMock() + redis_activity.redis_repository.set = AsyncMock(side_effect=Exception('test')) try: await redis_activity.put_last_data_timestamp(test_data) @@ -207,7 +215,7 @@ async def test_put_last_data_timestamp_error(redis_activity): except Exception as e: assert str(e) == 'test' - redis_activity.send_notification.assert_called_once_with( + redis_activity.send_notification_async.assert_called_once_with( metadata=metadata['metadata'], notification_id='REDIS_SET_ERROR', message='Error setting last data timestamp: test', @@ -399,7 +407,7 @@ async def test_group_and_hold_data_error_get(redis_activity): } redis_activity.redis_repository.get = AsyncMock(side_effect=Exception('test')) - redis_activity.send_notification = MagicMock() + redis_activity.send_notification_async = AsyncMock() try: await redis_activity.group_and_hold_data(test_data) @@ -407,7 +415,7 @@ async def test_group_and_hold_data_error_get(redis_activity): except Exception as e: assert str(e) == 'test' - redis_activity.send_notification.assert_called_once_with( + redis_activity.send_notification_async.assert_called_once_with( metadata=metadata['metadata'], notification_id='REDIS_GET_ERROR', message='Error getting held data: test', @@ -439,7 +447,7 @@ async def test_group_and_hold_data_error_set(redis_activity): # Mock get to return existing data redis_activity.redis_repository.get = AsyncMock(return_value=existing_data) redis_activity.redis_repository.set = AsyncMock(side_effect=Exception('test')) - redis_activity.send_notification = MagicMock() + redis_activity.send_notification_async = AsyncMock() try: await redis_activity.group_and_hold_data(test_data) @@ -485,7 +493,7 @@ async def test_store_data_package(redis_activity): async def test_store_data_package_error(redis_activity): """Test store_data_package error""" redis_activity.redis_repository.set = AsyncMock(side_effect=ValueError('test')) - redis_activity.send_notification = MagicMock() + redis_activity.send_notification_async = AsyncMock() test_data = { **metadata, @@ -511,7 +519,7 @@ async def test_store_data_package_error(redis_activity): with pytest.raises(ValueError): await redis_activity.store_data_package(test_data) - redis_activity.send_notification.assert_called_once_with( + redis_activity.send_notification_async.assert_called_once_with( metadata=metadata['metadata'], notification_id='REDIS_SET_ERROR', message='Error setting data package: test', From 32882d0b6fc7824dcf5ab1f92b7fc2ee76a549e7 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Tue, 4 Nov 2025 08:49:38 -0300 Subject: [PATCH 06/10] SIENTIAPDE-1325 Update metric labels from 'pipeline_name' to 'workflow_name' in metrics and related classes, ensuring consistency across the codebase and tests. --- scouter/activities/gates.py | 4 ++-- scouter/metrics.py | 2 +- tests/activities/test_gates.py | 6 +++--- tests/test_metrics.py | 4 ++-- 4 files changed, 8 insertions(+), 8 deletions(-) diff --git a/scouter/activities/gates.py b/scouter/activities/gates.py index b5ab579..0c2df2e 100644 --- a/scouter/activities/gates.py +++ b/scouter/activities/gates.py @@ -304,7 +304,7 @@ class Gates(SientiaMonitoring): metrics.LABORIOUS_DATA_WRITTEN_COUNT.labels( pod_id=self.pod_id, model_name=metadata['model_name'], - pipeline_name=metadata['workflow_name'], + workflow_name=metadata['workflow_name'], ).inc() # Register metrics @@ -312,7 +312,7 @@ class Gates(SientiaMonitoring): metrics.TAG_CHANGES_MONITOR.labels( pod_id=self.pod_id, model_name=metadata['model_name'], - pipeline_name=metadata['workflow_name'], + workflow_name=metadata['workflow_name'], tag_name=row['variable'], ).set(row['value']) diff --git a/scouter/metrics.py b/scouter/metrics.py index ac249c7..b25fa46 100644 --- a/scouter/metrics.py +++ b/scouter/metrics.py @@ -8,7 +8,7 @@ APP_UP = Gauge( ) # Core labels for consistent metric labeling -CORE_LABELS = ['pod_id', 'model_name', 'pipeline_name'] +CORE_LABELS = ['pod_id', 'model_name', 'workflow_name'] # Data processing metrics LABORIOUS_DATA_WRITTEN_COUNT = Counter( diff --git a/tests/activities/test_gates.py b/tests/activities/test_gates.py index d02ed0a..32d173e 100644 --- a/tests/activities/test_gates.py +++ b/tests/activities/test_gates.py @@ -390,7 +390,7 @@ async def test_write_metrics(mock_metrics, gates_fixture): mock_metrics.LABORIOUS_DATA_WRITTEN_COUNT.labels.assert_called_once_with( pod_id=gates_fixture.pod_id, model_name=metadata['metadata']['model_name'], - pipeline_name=metadata['metadata']['workflow_name'], + workflow_name=metadata['metadata']['workflow_name'], ) mock_metrics.LABORIOUS_DATA_WRITTEN_COUNT.labels.return_value.inc.assert_called_once() @@ -407,13 +407,13 @@ async def test_write_metrics(mock_metrics, gates_fixture): call( pod_id=gates_fixture.pod_id, model_name=metadata['metadata']['model_name'], - pipeline_name=metadata['metadata']['workflow_name'], + workflow_name=metadata['metadata']['workflow_name'], tag_name='tag1', ), call( pod_id=gates_fixture.pod_id, model_name=metadata['metadata']['model_name'], - pipeline_name=metadata['metadata']['workflow_name'], + workflow_name=metadata['metadata']['workflow_name'], tag_name='tag2', ), ], diff --git a/tests/test_metrics.py b/tests/test_metrics.py index 79a96d9..be7b4f3 100644 --- a/tests/test_metrics.py +++ b/tests/test_metrics.py @@ -15,7 +15,7 @@ def test_scouter_laborious_data_written_count(): assert set(metrics.LABORIOUS_DATA_WRITTEN_COUNT._labelnames) == { 'pod_id', 'model_name', - 'pipeline_name', + 'workflow_name', } @@ -27,6 +27,6 @@ def test_scouter_tag_changes_monitor(): assert set(metrics.TAG_CHANGES_MONITOR._labelnames) == { 'pod_id', 'model_name', - 'pipeline_name', + 'workflow_name', 'tag_name', } From d4fbb69b26b5d88393fafee0310d701fa9ad44f4 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Wed, 5 Nov 2025 09:28:17 -0300 Subject: [PATCH 07/10] SIENTIAPDE-1325 Update requirements to use the remote sientia-dataops-library repository and increment image tag in values.yaml from 0.4.9 to 0.5.0 for version consistency. --- requirements.txt | 2 +- values.yaml | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/requirements.txt b/requirements.txt index ce4f826..0f19531 100644 --- a/requirements.txt +++ b/requirements.txt @@ -5,6 +5,6 @@ asyncua redis aiokafka pymongo -/home/grezewave/Documents/projects/sientia/sientia-dataops-library +git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.5.2 pydruid[pandas] prometheus-client \ No newline at end of file diff --git a/values.yaml b/values.yaml index 804b505..5692f2b 100644 --- a/values.yaml +++ b/values.yaml @@ -11,7 +11,7 @@ image: # This sets the pull policy for images. pullPolicy: Always # Overrides the image tag whose default is the chart appVersion. - tag: "0.4.9" + tag: "0.5.0" # This is for the secrets for pulling an image from a private repository more information can be found here: https://kubernetes.io/docs/tasks/configure-pod-container/pull-image-private-registry/ imagePullSecrets: From 082ae682ee0a4b9d67a565eba4109c0cd9fd4cf4 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Thu, 6 Nov 2025 15:19:55 -0300 Subject: [PATCH 08/10] SIENTIAPDE-1325 Update requirements to use sientia-dataops-library version 1.5.3 and enhance Redis activity logging by including metadata in get and set operations. --- requirements.txt | 2 +- scouter/activities/redis.py | 16 ++++++++-------- 2 files changed, 9 insertions(+), 9 deletions(-) diff --git a/requirements.txt b/requirements.txt index 0f19531..1dc88da 100644 --- a/requirements.txt +++ b/requirements.txt @@ -5,6 +5,6 @@ asyncua redis aiokafka pymongo -git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.5.2 +git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.5.3 pydruid[pandas] prometheus-client \ No newline at end of file diff --git a/scouter/activities/redis.py b/scouter/activities/redis.py index 6e2b885..968c320 100644 --- a/scouter/activities/redis.py +++ b/scouter/activities/redis.py @@ -94,10 +94,10 @@ class Redis(SientiaMonitoring): metadata = input_data['metadata'] key = f'last_data_timestamp:{input_data["workflow_name"]}:{input_data["schedule_name"]}' - self.info(f'Getting last data timestamp for {key}') + self.info(f'Getting last data timestamp for {key}', metadata=metadata) try: - data_hold = await self.redis_repository.get(key) + data_hold = await self.redis_repository.get(key, metadata=metadata) except Exception as e: await self.send_notification_async( metadata=metadata, @@ -142,7 +142,7 @@ class Redis(SientiaMonitoring): metadata = input_data['metadata'] key = f'last_data_timestamp:{input_data["workflow_name"]}:{input_data["schedule_name"]}' - self.info(f'Putting last data timestamp for {key}') + self.info(f'Putting last data timestamp for {key}', metadata=metadata) data = DataFrame(input_data['data']) @@ -155,7 +155,7 @@ class Redis(SientiaMonitoring): self.info(f'Last collected timestamp to insert: {last_data_timestamp}', metadata=metadata) try: - await self.redis_repository.set(key, last_data_timestamp, ttl=60 * 60 * 5) + await self.redis_repository.set(key, last_data_timestamp, ttl=60 * 60 * 5, metadata=metadata) except Exception as e: await self.send_notification_async( metadata=metadata, @@ -205,10 +205,10 @@ class Redis(SientiaMonitoring): key = f'held_data_{input_data["workflow_name"]}_{input_data["schedule_name"]}' - self.info(f'Getting held data for {key}') + self.info(f'Getting held data for {key}', metadata=metadata) try: - data_hold = await self.redis_repository.get(key) + data_hold = await self.redis_repository.get(key, metadata=metadata) except Exception as e: await self.send_notification_async( metadata=metadata, @@ -252,7 +252,7 @@ class Redis(SientiaMonitoring): data['timestamp'].max() if not data.empty else data_hold['timestamp'] ) - await self.redis_repository.set(key, data_hold, ttl=retention_time) + await self.redis_repository.set(key, data_hold, ttl=retention_time, metadata=metadata) data_hold_df = DataFrame(data_hold, index=[0]) data_hold_melted = data_hold_df.melt( @@ -298,7 +298,7 @@ class Redis(SientiaMonitoring): cache = {'data': data.to_dict(), 'held_data': held_data.to_dict()} try: - await self.redis_repository.set(key, cache, ttl=120) + await self.redis_repository.set(key, cache, ttl=120, metadata=metadata) except Exception as e: await self.send_notification_async( metadata=metadata, From 3e8ba6c758a1c6a8ba7688b58b45a473a04ffca0 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Mon, 10 Nov 2025 08:23:01 -0300 Subject: [PATCH 09/10] SIENTIAPDE-1325 Update release workflow to trigger only on merged pull requests, ensuring proper release management. --- .github/workflows/release.yml | 1 + 1 file changed, 1 insertion(+) diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index cba1cce..4c48fb6 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -8,6 +8,7 @@ on: jobs: release: + if: github.event.pull_request.merged == true uses: Aignosi/github_workflow_templates/.github/workflows/dataops-module-release.yml@main permissions: write-all with: From 1b7391e8f1a6869b43a5873f73e7ee62b0f8f8b5 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Mon, 10 Nov 2025 08:39:19 -0300 Subject: [PATCH 10/10] SIENTIAPDE-1325 Enhance Redis activity by formatting set and get operations for improved readability and consistency. Update tests to include metadata in assertions for better verification of Redis interactions. --- scouter/activities/redis.py | 4 +++- tests/activities/test_redis.py | 18 ++++++++++++++---- 2 files changed, 17 insertions(+), 5 deletions(-) diff --git a/scouter/activities/redis.py b/scouter/activities/redis.py index 968c320..57083af 100644 --- a/scouter/activities/redis.py +++ b/scouter/activities/redis.py @@ -155,7 +155,9 @@ class Redis(SientiaMonitoring): self.info(f'Last collected timestamp to insert: {last_data_timestamp}', metadata=metadata) try: - await self.redis_repository.set(key, last_data_timestamp, ttl=60 * 60 * 5, metadata=metadata) + await self.redis_repository.set( + key, last_data_timestamp, ttl=60 * 60 * 5, metadata=metadata + ) except Exception as e: await self.send_notification_async( metadata=metadata, diff --git a/tests/activities/test_redis.py b/tests/activities/test_redis.py index 8c2b7ec..9944951 100644 --- a/tests/activities/test_redis.py +++ b/tests/activities/test_redis.py @@ -93,7 +93,8 @@ async def test_get_last_data_timestamp_none(redis_activity): result = await redis_activity.get_last_data_timestamp(test_data) redis_activity.redis_repository.get.assert_called_once_with( - 'last_data_timestamp:test_pipeline:test_schedule' + 'last_data_timestamp:test_pipeline:test_schedule', + metadata=metadata['metadata'], ) assert result is None @@ -109,7 +110,8 @@ async def test_get_last_data_timestamp_not_none(redis_activity): result = await redis_activity.get_last_data_timestamp(test_data) redis_activity.redis_repository.get.assert_called_once_with( - 'last_data_timestamp:test_pipeline:test_schedule' + 'last_data_timestamp:test_pipeline:test_schedule', + metadata=metadata['metadata'], ) assert result == '2023-01-01 12:00:00' @@ -186,7 +188,10 @@ async def test_put_last_data_timestamp_not_empty_dataframe(redis_activity): assert result == '2023-01-01 12:00:01' redis_activity.redis_repository.set.assert_called_once_with( - 'last_data_timestamp:test_pipeline:test_schedule', '2023-01-01 12:00:01', ttl=18000 + 'last_data_timestamp:test_pipeline:test_schedule', + '2023-01-01 12:00:01', + ttl=18000, + metadata=metadata['metadata'], ) @@ -270,6 +275,7 @@ async def test_group_and_hold_data_new_key(redis_activity): 'held_data_test_pipeline_test_schedule', {'sensor1': 25.5, 'sensor2': 30.0, 'timestamp': '2023-01-01 12:00:00'}, ttl=3600, + metadata=metadata['metadata'], ) @@ -334,6 +340,7 @@ async def test_group_and_hold_data_update_existing_fill_missing(redis_activity): 'timestamp': '2023-01-01 12:00:00', }, ttl=3600, + metadata=metadata['metadata'], ) @@ -485,7 +492,10 @@ async def test_store_data_package(redis_activity): await redis_activity.store_data_package(test_data) redis_activity.redis_repository.set.assert_called_once_with( - ANY, {'data': test_data['data'], 'held_data': test_data['held_data']}, ttl=120 + ANY, + {'data': test_data['data'], 'held_data': test_data['held_data']}, + ttl=120, + metadata=metadata['metadata'], )