From 0b77a7bc7e7985bf1273ed6d79583cf9841b934f Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Mon, 7 Jul 2025 10:03:08 -0300 Subject: [PATCH] SIENTIAPDE-1110 Implement store_data_package method in Redis activity for debug data storage. Update CoreScouter workflow to conditionally call store_data_package based on debug flag. Enhance tests for store_data_package functionality and error handling. --- scouter/activities/redis.py | 36 +++++++++- .../workflow/sub_workflows/core_scouter.py | 16 ++++- tests/activities/test_redis.py | 68 +++++++++++++++++++ .../sub_workflows/test_core_scouter.py | 18 ++++- 4 files changed, 134 insertions(+), 4 deletions(-) diff --git a/scouter/activities/redis.py b/scouter/activities/redis.py index 0336cb1..42482bb 100644 --- a/scouter/activities/redis.py +++ b/scouter/activities/redis.py @@ -93,7 +93,7 @@ class Redis(RedisBase): @activity.defn(name="group_and_hold_data") async def group_and_hold_data(self, input_data: dict[str, Any]): """ - Groups and holds data in redis. Keep a copy of the most recent + Groups and holds data in redis. Keep a copy of the most recent received data for a given pipeline and schedule. This activity updates the data in redis and return the full keeped data. @@ -169,3 +169,37 @@ class Redis(RedisBase): ) return data_hold_melted.to_dict() + + async def store_data_package(self, input_data: dict[str, Any]): + """ + Stores the data package in redis. It's a debug feature and must be toggled on. + input_data: + metadata: The metadata of the workflow. + workflow_name: The name of the workflow. + schedule_name: The name of the schedule. + held_data: The final scouter output. + data: The data used to collect the data. + """ + metadata = input_data['metadata'] + key = f"data_package_{input_data['workflow_name']}_{input_data['schedule_name']}_{datetime.now().strftime('%Y-%m-%d_%H-%M-%S')}" + + data = DataFrame(input_data['data']) + held_data = DataFrame(input_data['held_data']) + + cache = { + 'data': data.to_dict(), + 'held_data': held_data.to_dict() + } + + try: + self.set(key, cache, ttl=120) + except Exception as e: + self.send_notification( + metadata=metadata, + notification_id="REDIS_SET_ERROR", + message=f"Error setting data package: {e}", + block="store_data_package", + level=NotificationLevel.ERROR, + attachment_content=traceback.format_exc() + ) + raise e diff --git a/scouter/workflow/sub_workflows/core_scouter.py b/scouter/workflow/sub_workflows/core_scouter.py index de0e79a..21a6120 100644 --- a/scouter/workflow/sub_workflows/core_scouter.py +++ b/scouter/workflow/sub_workflows/core_scouter.py @@ -74,7 +74,7 @@ class CoreScouter: if held_data == {}: return - async_export = workflow.execute_activity_method( + await workflow.execute_activity_method( Activities.export_data_to_postgres, { **metadata, @@ -86,4 +86,16 @@ class CoreScouter: start_to_close_timeout=timedelta(seconds=60) ) - await async_export + if input_data.get('debug_data_package', False): + await workflow.execute_activity_method( + Activities.store_data_package, + { + **metadata, + 'data': input_data['data'], + 'held_data': held_data, + 'workflow_name': input_data['workflow_name'], + 'schedule_name': input_data['schedule_name'] + }, + retry_policy=retry_policy, + start_to_close_timeout=timedelta(seconds=60) + ) diff --git a/tests/activities/test_redis.py b/tests/activities/test_redis.py index 672d035..a833e8f 100644 --- a/tests/activities/test_redis.py +++ b/tests/activities/test_redis.py @@ -4,6 +4,7 @@ import pytest import numpy as np from pandas import DataFrame from sientia_do.notifications.handlers import NotificationHandler +from sientia_do.notifications.models import NotificationLevel from scouter.activities.redis import Redis @@ -280,3 +281,70 @@ async def test_group_and_hold_data_empty_dataframe(redis_activity): result = await redis_activity.group_and_hold_data(test_data) assert result == {} + + +@pytest.mark.asyncio +async def test_store_data_package(redis_activity): + """Test store_data_package""" + redis_activity.set = MagicMock() + + test_data = { + **metadata, + 'workflow_name': 'test_workflow', + 'schedule_name': 'test_schedule', + 'held_data': DataFrame({ + 'name': ['sensor1', 'sensor2'], + 'value': [25.5, 30.0], + 'timestamp': ['2023-01-01 12:00:00'] * 2 + }).to_dict(), + 'data': DataFrame({ + 'name': ['sensor1', 'sensor2'], + 'value': [25.5, 30.0], + 'timestamp': ['2023-01-01 12:00:00'] * 2 + }).to_dict() + } + + await redis_activity.store_data_package(test_data) + + redis_activity.set.assert_called_once_with( + ANY, + { + 'data': test_data['data'], + 'held_data': test_data['held_data'] + }, + ttl=120) + + +@pytest.mark.asyncio +async def test_store_data_package_error(redis_activity): + """Test store_data_package error""" + redis_activity.set = MagicMock(side_effect=Exception('test')) + redis_activity.send_notification = MagicMock() + + test_data = { + **metadata, + 'workflow_name': 'test_workflow', + 'schedule_name': 'test_schedule', + 'held_data': DataFrame({ + 'name': ['sensor1', 'sensor2'], + 'value': [25.5, 30.0], + 'timestamp': ['2023-01-01 12:00:00'] * 2 + }).to_dict(), + 'data': DataFrame({ + 'name': ['sensor1', 'sensor2'], + 'value': [25.5, 30.0], + 'timestamp': ['2023-01-01 12:00:00'] * 2 + }).to_dict() + } + + with pytest.raises(Exception): + await redis_activity.store_data_package(test_data) + + redis_activity.send_notification.assert_called_once_with( + metadata=metadata['metadata'], + notification_id="REDIS_SET_ERROR", + message="Error setting data package: test", + block="store_data_package", + level=NotificationLevel.ERROR, + attachment_content=ANY + ) diff --git a/tests/workflow/sub_workflows/test_core_scouter.py b/tests/workflow/sub_workflows/test_core_scouter.py index c92a175..e1be284 100644 --- a/tests/workflow/sub_workflows/test_core_scouter.py +++ b/tests/workflow/sub_workflows/test_core_scouter.py @@ -34,7 +34,8 @@ async def test_core_scouter_workflow_success(mock_workflow, core_scouter): 'schema': 'test_schema', 'table_name': 'test_table', 'retention_time': 3600, - 'model_tags': {} + 'model_tags': {}, + 'debug_data_package': True } ) @@ -98,6 +99,21 @@ async def test_core_scouter_workflow_success(mock_workflow, core_scouter): ) ]) + mock_workflow.execute_activity_method.assert_has_calls([ + call( + Activities.store_data_package, + { + **expected_metadata, + 'workflow_name': 'test_workflow', + 'schedule_name': 'test_schedule', + 'held_data': 'held_data', + 'data': 'test_data' + }, + retry_policy=ANY, + start_to_close_timeout=ANY + ) + ]) + @pytest.mark.asyncio @patch('scouter.workflow.sub_workflows.core_scouter.workflow', new_callable=AsyncMock)