From b17eaf972c94cb6f12683b027af97e706af0cf5d Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Tue, 30 Dec 2025 08:40:07 -0300 Subject: [PATCH] SIENTIAPDE-1445 Update requirements-dev.txt to add E2E testing dependencies: fakeredis and mongomock for in-memory testing, and include testcontainers for PostgreSQL support. --- e2e/__init__.py | 0 e2e/conftest.py | 365 +++++++++++ e2e/fixtures/__init__.py | 0 e2e/fixtures/fake_mongodb_repository.py | 241 +++++++ e2e/fixtures/fake_redis_repository.py | 170 +++++ e2e/scenarios.md | 804 ++++++++++++++++++++++++ e2e/test_pi_web_api_scouter.py | 106 ++++ requirements-dev.txt | 7 +- 8 files changed, 1692 insertions(+), 1 deletion(-) create mode 100644 e2e/__init__.py create mode 100644 e2e/conftest.py create mode 100644 e2e/fixtures/__init__.py create mode 100644 e2e/fixtures/fake_mongodb_repository.py create mode 100644 e2e/fixtures/fake_redis_repository.py create mode 100644 e2e/scenarios.md create mode 100644 e2e/test_pi_web_api_scouter.py diff --git a/e2e/__init__.py b/e2e/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/e2e/conftest.py b/e2e/conftest.py new file mode 100644 index 0000000..b4f0f13 --- /dev/null +++ b/e2e/conftest.py @@ -0,0 +1,365 @@ +""" +Pytest configuration and fixtures for E2E tests. +""" + +from typing import Any +from unittest.mock import AsyncMock, MagicMock, patch + +import pandas as pd +import pytest +import pytest_asyncio +from pandas import DataFrame +from sqlalchemy import create_engine, text +from sqlalchemy.orm import sessionmaker +from testcontainers.postgres import PostgresContainer +from temporalio.testing import WorkflowEnvironment +from temporalio.worker import Worker + +from scouter.activities.activities import Activities +from scouter.workflow.pi_web_api_scouter import PIWebAPIScouter +from scouter.workflow.sub_workflows.core_scouter import CoreScouter +from sientia_do.notifications.handlers import CoreNotificationHandler +from sientia_do.observability.logger import Logger +from sientia_do.observability.metrics_controller import MetricsController + +from e2e.fixtures.fake_mongodb_repository import FakeMongoDBRepository +from e2e.fixtures.fake_redis_repository import FakeRedisRepository + +# Test constants +TEST_MONGODB_CONNECTION_STRING = 'mongodb://localhost:27017' +TEST_DATABASE_NAME = 'test_db' + + +@pytest.fixture(scope='session') +def postgres_container(): + """ + Create a PostgreSQL container using testcontainers. + + This fixture creates a real PostgreSQL database in a Docker container + that will be used for all tests in the session. + """ + postgres = PostgresContainer('postgres:15') + postgres.start() + yield postgres + postgres.stop() + + +@pytest.fixture +def postgres_engine(postgres_container): + """ + Create SQLAlchemy engine for PostgreSQL test database. + + This fixture creates a connection to the PostgreSQL container + created by the postgres_container fixture. + """ + # Get connection URL from container + connection_string = postgres_container.get_connection_url() + engine = create_engine(connection_string) + + yield engine + + engine.dispose() + + +def _create_schema_and_table(engine): + """ + Helper function to create schema and table in the given engine. + + This is used by both the autouse fixture and test_activities to ensure + the schema exists before Activities tries to use it. + + Note: For tests, we create a non-partitioned table to avoid issues + with pandas to_sql recognizing partitioned tables. + """ + schema_name = 'sientia_data' + table_name = 'laborious_data' + + # Use begin() to ensure transaction is properly committed + with engine.begin() as conn: + # Create schema + conn.execute(text(f"CREATE SCHEMA IF NOT EXISTS {schema_name}")) + + # Create table WITHOUT partitioning (simpler for tests) + # Same structure as production, but without PARTITION BY RANGE + # Use UNIQUE constraint directly since table is not partitioned + create_table_sql = f""" + CREATE TABLE IF NOT EXISTS {schema_name}.{table_name} ( + id SERIAL NOT NULL, + model_id int4 NOT NULL, + variable text NOT NULL, + value numeric NULL, + "timestamp" timestamptz NOT NULL, + created_at timestamptz DEFAULT CURRENT_TIMESTAMP NOT NULL, + PRIMARY KEY (id, created_at), + UNIQUE (model_id, timestamp, variable) + ); + """ + + conn.execute(text(create_table_sql)) + # Transaction is automatically committed when exiting the 'with' block + + +@pytest.fixture(autouse=True) +def setup_postgres_schema_and_table(postgres_engine): + """ + Automatically create necessary schema and table before each test. + + This fixture runs automatically (autouse=True) and ensures + that the sientia_data schema and laborious_data table exist + with the correct structure before tests execute. + + Note: For tests, we use a non-partitioned table with a UNIQUE constraint + directly in the table definition, which is simpler and avoids issues + with pandas to_sql recognizing partitioned tables. + """ + _create_schema_and_table(postgres_engine) + yield + + +@pytest.fixture +def mock_logger(): + """Mock logger for testing.""" + logger = MagicMock(spec=Logger) + logger.info = MagicMock() + logger.debug = MagicMock() + logger.error = MagicMock() + logger.warning = MagicMock() + logger.custom_info = MagicMock() + return logger + + +@pytest.fixture +def mock_mongo_client(): + """ + Mock MongoDB client to avoid real connections. + + This fixture mocks the pymongo.MongoClient used by CoreNotificationHandler, + allowing us to use a real NotificationHandler instance without connecting to MongoDB. + """ + mock_client = MagicMock() + mock_db = MagicMock() + mock_collection = MagicMock() + + # Configure the mock chain: client[database] -> db[collection] -> collection + mock_client.__getitem__.return_value = mock_db + mock_db.__getitem__.return_value = mock_collection + + # Mock server_info() to avoid connection attempts + mock_client.server_info = MagicMock() + + # Mock insert_one for notifications + mock_collection.insert_one = MagicMock() + + return mock_client + + +@pytest.fixture +def notification_handler(mock_logger, mock_mongo_client): + """ + Create a real NotificationHandler instance with mocked MongoDB client. + + This fixture creates a real CoreNotificationHandler instance but mocks + the underlying MongoDB connection to avoid real database connections. + """ + # Patch MongoClient where it's imported in the handlers module + with patch('sientia_do.notifications.handlers.MongoClient', return_value=mock_mongo_client): + handler = CoreNotificationHandler( + connection_string=TEST_MONGODB_CONNECTION_STRING, + database=TEST_DATABASE_NAME, + logger=mock_logger, + project_name='scouter', + ) + yield handler + handler.shutdown() + + +@pytest.fixture +def metrics_controller(mock_logger): + """ + Create a real MetricsController instance. + + MetricsController doesn't require external services, so we can use + a real instance without mocking anything. + """ + controller = MetricsController(logger=mock_logger) + yield controller + # MetricsController might have cleanup, but it's optional + + +@pytest.fixture +def mock_pi_web_api_client(): + """Mock PI Web API client.""" + mock_client = MagicMock() + + # Mock DataFrame response similar to real API + # The real PIWebAPIClient returns timestamp as datetime, so we need to match that + mock_df = DataFrame({ + 'timestamp': [ + '2024-01-01 12:00:00+0000', + '2024-01-01 12:01:00+0000', + '2024-01-01 12:02:00+0000', + ], + 'name': ['tag1', 'tag2', 'tag3'], + 'value': [10.5, 20.3, 30.7], + 'tag': ['webid1', 'webid2', 'webid3'], + }) + + # Convert timestamp to datetime (UTC, floored to seconds) to match real client behavior + mock_df['timestamp'] = pd.to_datetime(mock_df['timestamp'], utc=True).dt.floor('s') + + mock_client.get_latest_values_df = AsyncMock(return_value=mock_df) + mock_client.close = MagicMock() + return mock_client + + +@pytest_asyncio.fixture +async def test_activities( + postgres_engine, + postgres_container, + mock_logger, + notification_handler, + metrics_controller, + mock_pi_web_api_client, +): + """ + Create Activities instance with test dependencies. + + This fixture creates a real Activities instance with: + - PostgreSQL database (via testcontainers) + - FakeRedis instead of real Redis + - FakeMongoDB instead of real MongoDB + - Mocked PI Web API client + - Real NotificationHandler and MetricsController (with mocked underlying services) + """ + # Get connection details from container + connection_string = postgres_container.get_connection_url() + + # Parse connection string to get individual components + # Format: postgresql://testuser:testpass@localhost:5432/test + from urllib.parse import urlparse + parsed = urlparse(connection_string) + + # Ensure schema and table exist BEFORE creating Activities + # This ensures the schema exists when Activities initializes its engine + _create_schema_and_table(postgres_engine) + + # Create fake repositories that will be used from the start + fake_redis_repo = FakeRedisRepository( + host='localhost', + port=6379, + username='', + password='', + logger=mock_logger, + notification_handler=notification_handler, + metrics_controller=metrics_controller, + ) + + fake_mongo_repo = FakeMongoDBRepository( + connection_string=TEST_MONGODB_CONNECTION_STRING, + database_name=TEST_DATABASE_NAME, + logger=mock_logger, + notification_handler=notification_handler, + metrics_controller=metrics_controller, + ) + + # Patch create_engine to return our postgres_engine instead of creating a new one + # This ensures Activities uses the same engine from the start + original_create_engine = create_engine + + def patched_create_engine(connection_string, *args, **kwargs): + # Check if this is the connection string that Activities would create + # Activities creates: postgresql://user:password@host:port/dbname + expected_conn_str = f"postgresql://{parsed.username or 'test'}:{parsed.password or 'test'}@localhost:{postgres_container.get_exposed_port(5432)}/{parsed.path.lstrip('/') if parsed.path else 'test'}" + + # If it matches our test container connection, return our engine + if connection_string == expected_conn_str: + return postgres_engine + # Otherwise, use the original create_engine + return original_create_engine(connection_string, *args, **kwargs) + + # Patch RedisRepository to return our fake repository from the start + def patched_redis_repository(*args, **kwargs): + return fake_redis_repo + + # Patch MongoDBRepository to return our fake repository from the start + def patched_mongodb_repository(*args, **kwargs): + return fake_mongo_repo + + # Patch PIWebAPIClient to return our mock from the start + def patched_pi_web_api_client(*args, **kwargs): + return mock_pi_web_api_client + + # Create Activities with test configurations + # All patches ensure it uses our test instances from the start + with patch('sientia_do.temporal.activities.postgres.create_engine', new=patched_create_engine), \ + patch('scouter.activities.redis.RedisRepository', new=patched_redis_repository), \ + patch('scouter.activities.mongodb.MongoDBRepository', new=patched_mongodb_repository), \ + patch('scouter.activities.api.PIWebAPIClient', new=patched_pi_web_api_client): + + activities = Activities( + postgres_config={ + 'host': 'localhost', # Container exposes to localhost + 'port': postgres_container.get_exposed_port(5432), + 'user': parsed.username or 'test', + 'password': parsed.password or 'test', + 'dbname': parsed.path.lstrip('/') if parsed.path else 'test', + 'min_connections': 1, + 'max_connections': 5, + }, + redis_config={ + 'host': 'localhost', + 'port': 6379, + 'username': '', + 'password': '', + }, + mongodb_config={ + 'connection_string': TEST_MONGODB_CONNECTION_STRING, + 'database_name': TEST_DATABASE_NAME, + }, + api_config={ + 'base_url': 'http://localhost:8080', + 'auth_type': 'bearer', + 'auth_token': 'test_token', + }, + logger=mock_logger, + notification_handler=notification_handler, + ) + + # Verify that Activities is using our instances from the start + assert activities.engine is postgres_engine, "Activities should use the same engine as postgres_engine" + assert activities.redis_repository is fake_redis_repo, "Activities should use the same fake Redis repository" + assert activities.mongodb_repository is fake_mongo_repo, "Activities should use the same fake MongoDB repository" + assert activities.pi_web_api_client is mock_pi_web_api_client, "Activities should use the same mock PI Web API client" + + yield activities + + # Cleanup + activities.shutdown() + + +@pytest_asyncio.fixture +async def temporal_test_env(): + """Create Temporal test environment.""" + env = await WorkflowEnvironment.start_time_skipping() + async with env: + yield env + + +@pytest_asyncio.fixture +async def temporal_worker(temporal_test_env, test_activities): + """Create Temporal worker with test activities.""" + async with Worker( + temporal_test_env.client, + task_queue='test-queue', + workflows=[PIWebAPIScouter, CoreScouter], + activities=[ + test_activities.get_tag_values, + test_activities.data_quality_gate, + test_activities.aggregate_data, + test_activities.group_and_hold_data, + test_activities.export_data_to_postgres, + test_activities.write_metrics, + test_activities.store_data_package, + ], + ) as worker: + yield worker diff --git a/e2e/fixtures/__init__.py b/e2e/fixtures/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/e2e/fixtures/fake_mongodb_repository.py b/e2e/fixtures/fake_mongodb_repository.py new file mode 100644 index 0000000..087332d --- /dev/null +++ b/e2e/fixtures/fake_mongodb_repository.py @@ -0,0 +1,241 @@ +""" +Fake MongoDB Repository adapter for testing. + +This adapter implements the MongoDBRepository interface using mongomock +to provide an in-memory MongoDB server for testing. +""" + +import time +from typing import Any + +from mongomock import MongoClient +from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler +from sientia_do.observability.logger import Logger +from sientia_do.observability.metrics_controller import MetricsController +from sientia_do.observability.sientia_monitoring import SientiaMonitoring + + +def clear_mongo_id(docs: list) -> list: + """ + Remove MongoDB internal `_id` fields from nested structures. + + Same implementation as in the real MongoDBRepository. + + Args: + docs: The list of documents or nested structures to clean. + + Returns: + The cleaned documents with `_id` fields removed wherever present. + """ + 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 FakeMongoDBRepository(SientiaMonitoring): + """ + Fake MongoDB Repository that uses mongomock for testing. + + Implements the same interface as MongoDBRepository but uses + mongomock for in-memory MongoDB operations. + """ + + def __init__( + self, + connection_string: str, + database_name: str, + logger: Logger, + notification_handler: NotificationHandler, + metrics_controller: MetricsController, + ): + """ + Initialize fake MongoDB repository with mongomock. + + Args: + connection_string: MongoDB connection string (ignored in fake mode) + database_name: Target database name + logger: Logger instance + notification_handler: Notification handler + metrics_controller: Metrics controller + """ + SientiaMonitoring.__init__( + self, + logger=logger, + notification_handler=notification_handler, + metrics_controller=metrics_controller, + ) + + self.database_name = database_name + # Use mongomock instead of real MongoDB + self.mongo_client = MongoClient() + self.database = self.mongo_client[self.database_name] + logger.info('Fake MongoDB connection initialized') + + def close(self): + """ + Closes fake MongoDB connection and shuts down monitoring. + """ + try: + if self.mongo_client: + self.logger.info('Closing fake MongoDB connection...') + self.mongo_client.close() + self.logger.info('Fake MongoDB connection closed successfully') + except Exception as e: + self.logger.error(f'Failed to close fake MongoDB connection: {e}') + SientiaMonitoring.shutdown(self) + + def __del__(self): + """ + Destructor that closes the connection when destroying the instance. + """ + self.close() + + async def find( + self, + collection_name: str, + filters: dict[str, Any], + metadata: dict[str, Any], + ) -> list[dict[str, Any]]: + """ + Finds documents in a MongoDB collection based on provided filters. + + Args: + collection_name: Name of the collection to search + filters: Query filters to apply + metadata: Dictionary with additional metadata + + Return: + List of documents matching the filters (with `_id` removed) + """ + try: + start_time = time.time() + collection = self.database[collection_name] + documents = list(collection.find(filters, {'_id': 0})) + except Exception as e: + self.logger.error(f'Failed to find documents in fake MongoDB: {e}') + raise e + + return clear_mongo_id(documents) + + async def aggregate( + self, + collection_name: str, + pipeline: list[dict[str, Any]], + metadata: dict[str, Any], + ) -> list[dict[str, Any]]: + """ + Executes an aggregation on a MongoDB collection. + + Args: + collection_name: Name of the collection to aggregate + pipeline: MongoDB aggregation pipeline + metadata: Dictionary with additional metadata + + Return: + List of documents resulting from aggregation (with `_id` removed) + """ + try: + start_time = time.time() + collection = self.database[collection_name] + documents = list(collection.aggregate(pipeline)) + except Exception as e: + self.logger.error(f'Failed to aggregate documents in fake MongoDB: {e}') + raise e + + return clear_mongo_id(documents) + + async def update_many( + self, + collection_name: str, + filters: dict[str, Any], + update: dict[str, Any], + metadata: dict[str, Any], + ) -> None: + """ + Updates multiple documents in a MongoDB collection. + + Args: + collection_name: Collection name + filters: Filters to identify documents to update + update: Update operations to apply + metadata: Dictionary with additional metadata + """ + try: + collection = self.database[collection_name] + collection.update_many(filters, update) + except Exception as e: + self.logger.error(f'Failed to update documents in fake MongoDB: {e}') + raise e + + async def insert_many( + self, + collection_name: str, + documents: list[dict[str, Any]], + metadata: dict[str, Any], + ) -> None: + """ + Inserts multiple documents into a MongoDB collection. + + Args: + collection_name: Collection name + documents: List of documents to insert + metadata: Dictionary with additional metadata + """ + try: + collection = self.database[collection_name] + collection.insert_many(documents) + except Exception as e: + self.logger.error(f'Failed to insert documents in fake MongoDB: {e}') + raise e + + async def insert( + self, + collection_name: str, + document: dict[str, Any], + metadata: dict[str, Any], + ) -> None: + """ + Inserts a document into a MongoDB collection. + + Args: + collection_name: Collection name + document: Document to insert + metadata: Dictionary with additional metadata + """ + try: + collection = self.database[collection_name] + collection.insert_one(document) + except Exception as e: + self.logger.error(f'Failed to insert document in fake MongoDB: {e}') + raise e + + async def delete_many( + self, + collection_name: str, + filters: dict[str, Any], + metadata: dict[str, Any], + ) -> None: + """ + Removes multiple documents from a MongoDB collection. + + Args: + collection_name: Collection name + filters: Filters to identify documents to remove + metadata: Dictionary with additional metadata + """ + try: + collection = self.database[collection_name] + collection.delete_many(filters) + except Exception as e: + self.logger.error(f'Failed to delete documents in fake MongoDB: {e}') + raise e + diff --git a/e2e/fixtures/fake_redis_repository.py b/e2e/fixtures/fake_redis_repository.py new file mode 100644 index 0000000..9dbb762 --- /dev/null +++ b/e2e/fixtures/fake_redis_repository.py @@ -0,0 +1,170 @@ +""" +Fake Redis Repository adapter for testing. + +This adapter implements the RedisRepository interface using fakeredis +to provide an in-memory Redis server for testing. +""" + +import json +from typing import Any + +import fakeredis +from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler +from sientia_do.observability.logger import Logger +from sientia_do.observability.metrics_controller import MetricsController +from sientia_do.observability.sientia_monitoring import SientiaMonitoring + + +class FakeRedisRepository(SientiaMonitoring): + """ + Fake Redis Repository that uses fakeredis for testing. + + Implements the same interface as RedisRepository but uses + fakeredis for in-memory Redis operations. + """ + + def __init__( + self, + host: str, + port: int, + username: str, + password: str, + logger: Logger, + notification_handler: NotificationHandler, + metrics_controller: MetricsController, + ): + """ + Initialize fake Redis repository with fakeredis. + + Args: + host: Redis server address (ignored in fake mode) + port: Redis server port (ignored in fake mode) + username: Username (ignored in fake mode) + password: Password (ignored in fake mode) + logger: Logger instance + notification_handler: Notification handler + metrics_controller: Metrics controller + """ + SientiaMonitoring.__init__( + self, + logger=logger, + notification_handler=notification_handler, + metrics_controller=metrics_controller, + ) + # Create fake Redis server + self.redis_client = fakeredis.FakeStrictRedis( + decode_responses=True + ) + + def _get_redis(self): + """Get fakeredis connection.""" + return self.redis_client + + async def get(self, key: str, metadata: dict | None = None): + """ + Gets a value from fake Redis by key. + + Args: + key: Key of the value to retrieve + metadata: Optional dictionary with additional metadata + + Return: + Deserialized value from Redis or None if not found + """ + redis = self._get_redis() + try: + history = redis.get(key) + return json.loads(history) if history else None + except Exception as e: + self.error(f'Error getting data from fake redis: {e}', metadata or {}) + raise e + + async def set( + self, + key: str, + data: dict, + ttl: int = 600, + nx: bool = False, + metadata: dict | None = None, + ) -> bool: + """ + Sets a value in fake Redis with optional TTL. + + Args: + key: Key of the value to set + data: Dictionary with data to store + ttl: Time to live in seconds (default: 600) + nx: If True, only sets if key doesn't exist (default: False) + metadata: Optional dictionary with additional metadata + + Return: + True if value was set, False otherwise + """ + redis = self._get_redis() + try: + result = redis.set( + key, json.dumps(data), ex=ttl, nx=nx + ) + return bool(result) + except Exception as e: + self.error(f'Error setting data in fake redis: {e}', metadata or {}) + raise e + + async def delete(self, key: str, metadata: dict | None = None): + """ + Removes a key from fake Redis. + + Args: + key: Key to be removed + metadata: Optional dictionary with additional metadata + """ + redis = self._get_redis() + try: + redis.delete(key) + except Exception as e: + self.error(f'Error deleting data from fake redis: {e}', metadata or {}) + raise e + + async def expire(self, key: str, ttl: int, metadata: dict | None = None): + """ + Sets the time to live (TTL) of an existing Redis key. + + Args: + key: Key whose TTL will be set + ttl: Time to live in seconds + metadata: Optional dictionary with additional metadata + """ + redis = self._get_redis() + try: + redis.expire(key, ttl) + except Exception as e: + self.error(f'Error expiring data from fake redis: {e}', metadata or {}) + raise e + + async def keys(self, pattern: str, metadata: dict | None = None): + """ + Gets all keys matching the specified pattern. + + Args: + pattern: Search pattern for keys + metadata: Optional dictionary with additional metadata + + Return: + List of keys matching the pattern + """ + redis = self._get_redis() + try: + keys = redis.keys(pattern) + return keys + except Exception as e: + self.error(f'Error getting keys from fake redis: {e}', metadata or {}) + raise e + + def close(self): + """ + Closes fake Redis connection and shuts down monitoring. + """ + if self.redis_client: + self.redis_client.close() + SientiaMonitoring.shutdown(self) + diff --git a/e2e/scenarios.md b/e2e/scenarios.md new file mode 100644 index 0000000..99087ee --- /dev/null +++ b/e2e/scenarios.md @@ -0,0 +1,804 @@ +# Test Scenarios for PI Web API Scouter Workflow + +This document describes all possible test scenarios for the `pi_web_api_scouter` workflow and its child workflow `core_scouter`. + +## Workflow Overview + +The `pi_web_api_scouter` workflow: +1. Retrieves tag values from PI Web API +2. Delegates processing to `core_scouter` child workflow which: + - Applies data quality gates + - Aggregates data + - Groups and holds data in Redis + - Exports to PostgreSQL + - Writes metrics + - Optionally stores debug data package + +--- + +## 1. PI Web API Scouter - Main Workflow Scenarios + +### 1.1 Success Scenarios + +#### Scenario 1.1.1: Happy Path - Complete Success +**Description**: Workflow completes successfully with valid data from PI Web API + +**Input**: +- Valid `model_name`, `model_id`, `schedule_name` +- Valid `pi_web_api_query` with endpoint, period, max_count, api_timeout +- Valid `model_tags` with webids and configurations +- Valid filters, schema, table_name, retention_time + +**Expected Behavior**: +- `get_tag_values` returns non-empty list of records +- Workflow proceeds to `core_scouter` +- All activities execute successfully +- Data is stored in PostgreSQL +- Metrics are written +- Workflow completes without errors + +**Assertions**: +- PI Web API client called once with correct parameters +- Data exists in SQLite (PostgreSQL substitute) +- Data cached in Redis +- Metrics written +- No errors raised + +--- + +#### Scenario 1.1.2: Success with Multiple Tags +**Description**: Workflow processes multiple tags successfully + +**Input**: +- Multiple tags in `model_tags` (3+ tags) +- Each tag has valid webid, aggr_function, data_range, frequency + +**Expected Behavior**: +- All tags retrieved from PI Web API +- All tags processed through quality gates +- All tags aggregated correctly +- All tags stored in database + +**Assertions**: +- Number of records matches number of tags +- All tags present in final data +- Aggregation applied per tag configuration + +--- + +#### Scenario 1.1.3: Success with Debug Data Package Enabled +**Description**: Workflow completes with `debug_data_package=True` + +**Input**: +- All standard input +- `debug_data_package: True` + +**Expected Behavior**: +- Normal workflow execution +- `store_data_package` activity called +- Data package stored in Redis + +**Assertions**: +- `store_data_package` called once +- Data package key exists in Redis +- Package contains both `data` and `held_data` + +--- + +### 1.2 Early Exit Scenarios + +#### Scenario 1.2.1: Empty Data from PI Web API +**Description**: PI Web API returns empty data + +**Input**: +- Valid configuration +- PI Web API returns empty DataFrame or empty list + +**Expected Behavior**: +- `get_tag_values` returns empty list `[]` +- Workflow checks `if not data:` and returns early +- `core_scouter` is NOT called +- Workflow completes without error + +**Assertions**: +- PI Web API called once +- `core_scouter` NOT called +- No data in PostgreSQL +- No data in Redis (except possibly from previous runs) + +--- + +#### Scenario 1.2.2: None Returned from PI Web API +**Description**: PI Web API returns None + +**Input**: +- Valid configuration +- PI Web API returns None + +**Expected Behavior**: +- `get_tag_values` returns None +- Workflow checks `if not data:` and returns early +- `core_scouter` is NOT called + +**Assertions**: +- PI Web API called once +- `core_scouter` NOT called +- Workflow completes without error + +--- + +### 1.3 Error Scenarios + +#### Scenario 1.3.1: PI Web API Connection Error +**Description**: PI Web API client raises connection error + +**Input**: +- Valid configuration +- PI Web API client raises `PIMSRequestError` or connection exception + +**Expected Behavior**: +- `get_tag_values` catches exception +- Sends notification with `PI_WEB_API_REQUEST_ERROR` +- Raises exception (workflow fails after retries) + +**Assertions**: +- Notification sent with correct error details +- Exception propagated to workflow +- Workflow fails (after retry policy exhausted) +- `core_scouter` NOT called + +--- + +#### Scenario 1.3.2: PI Web API Timeout +**Description**: PI Web API request times out + +**Input**: +- Valid configuration +- `api_timeout` set to low value +- PI Web API takes longer than timeout + +**Expected Behavior**: +- Request times out +- Exception raised +- Notification sent +- Workflow fails after retries + +**Assertions**: +- Timeout exception caught +- Notification sent +- Workflow fails + +--- + +#### Scenario 1.3.3: Invalid Endpoint +**Description**: Invalid PI Web API endpoint provided + +**Input**: +- Invalid endpoint path in `pi_web_api_query` + +**Expected Behavior**: +- PI Web API client raises error +- Notification sent +- Workflow fails + +**Assertions**: +- Error notification sent +- Workflow fails + +--- + +#### Scenario 1.3.4: Missing Required Input Fields +**Description**: Missing required input fields + +**Input**: +- Missing `model_id`, `model_name`, `schedule_name`, or `pi_web_api_query` + +**Expected Behavior**: +- KeyError raised when accessing missing fields +- Workflow fails immediately + +**Assertions**: +- KeyError or similar exception +- Workflow fails before any activity execution + +--- + +## 2. CoreScouter - Child Workflow Scenarios + +### 2.1 Success Scenarios + +#### Scenario 2.1.1: Complete Processing Success +**Description**: All stages complete successfully + +**Input**: +- Valid data from parent workflow +- Valid filters, model_tags, schema, table_name +- `fill_missing_tags: False` +- `debug_data_package: False` + +**Expected Behavior**: +- `data_quality_gate` filters data +- `aggregate_data` aggregates by tag +- `group_and_hold_data` stores in Redis +- `export_data_to_postgres` writes to database +- `write_metrics` records metrics +- Workflow completes + +**Assertions**: +- All activities called in correct order +- Data in PostgreSQL +- Data in Redis +- Metrics written +- `store_data_package` NOT called + +--- + +#### Scenario 2.1.2: Success with Data Quality Filters +**Description**: Data quality filters applied successfully + +**Input**: +- Data with some quality issues +- Filters configured with `NULL_VALUES_FILTER` or `OUT_OF_BOUNDS_FILTER` +- Policy set to `DISCARD` or `WARN` + +**Expected Behavior**: +- Quality gate identifies issues +- Notification sent (WARNING level) +- If policy is `DISCARD`, bad rows removed +- Remaining data processed normally + +**Assertions**: +- Quality issues detected +- Notification sent +- Bad data discarded if policy is `DISCARD` +- Good data processed + +--- + +#### Scenario 2.1.3: Success with Different Aggregation Functions +**Description**: Different aggregation functions applied correctly + +**Input**: +- Multiple tags with different `aggr_function`: `avg`, `mdn`, `max`, `min`, `lts` +- Time-series data with multiple points per tag + +**Expected Behavior**: +- Each tag aggregated with its configured function +- Aggregated values correct for each function type + +**Assertions**: +- `avg` calculates mean correctly +- `mdn` calculates median correctly +- `max` returns maximum value +- `min` returns minimum value +- `lts` returns latest value + +--- + +#### Scenario 2.1.4: Success with Fill Missing Tags +**Description**: Missing tags filled with None + +**Input**: +- `fill_missing_tags: True` +- Some tags missing from data + +**Expected Behavior**: +- Missing tags added to `data_hold` with value `None` +- All expected tags present in final data + +**Assertions**: +- Missing tags present with `None` value +- All model_tags represented in output + +--- + +### 2.2 Early Exit Scenarios + +#### Scenario 2.2.1: Empty Data After Grouping +**Description**: `group_and_hold_data` returns empty dict + +**Input**: +- Data that results in empty `held_data` after grouping + +**Expected Behavior**: +- `group_and_hold_data` returns `{}` +- Workflow checks `if held_data == {}:` and returns early +- `export_data_to_postgres` NOT called +- `write_metrics` NOT called +- `store_data_package` NOT called + +**Assertions**: +- Early return after grouping +- No database export +- No metrics written +- Workflow completes without error + +--- + +#### Scenario 2.2.2: Zero Affected Rows After Export +**Description**: PostgreSQL export returns zero affected rows + +**Input**: +- Data that results in `affected_rows: 0` from export + +**Expected Behavior**: +- `export_data_to_postgres` returns `{'affected_rows': 0}` +- Workflow checks `if data_exported.get('affected_rows', 0) <= 0:` and returns early +- `write_metrics` NOT called +- `store_data_package` NOT called + +**Assertions**: +- Early return after export +- No metrics written +- Workflow completes without error + +--- + +### 2.3 Error Scenarios + +#### Scenario 2.3.1: Data Quality Gate Error +**Description**: Error during quality gate processing + +**Input**: +- Invalid filter configuration +- Filter function raises exception + +**Expected Behavior**: +- Exception caught in quality gate +- Notification sent with `DATA_QUALITY_GATE_ISSUES` +- Exception propagated (workflow fails after retries) + +**Assertions**: +- Error notification sent +- Workflow fails + +--- + +#### Scenario 2.3.2: Aggregation Error +**Description**: Error during data aggregation + +**Input**: +- Invalid aggregation function +- Data format issues + +**Expected Behavior**: +- Invalid function sends notification with `AGGREGATION_ISSUES` +- Returns `'continue'` for invalid function (skips that tag) +- Other errors raise exception + +**Assertions**: +- Invalid function handled gracefully +- Other errors cause workflow failure + +--- + +#### Scenario 2.3.3: Redis Connection Error +**Description**: Redis unavailable during `group_and_hold_data` + +**Input**: +- Valid data +- Redis connection fails + +**Expected Behavior**: +- `redis_repository.get()` or `redis_repository.set()` raises exception +- Notification sent with `REDIS_GET_ERROR` or `REDIS_SET_ERROR` +- Exception propagated (workflow fails after retries) + +**Assertions**: +- Error notification sent +- Workflow fails + +--- + +#### Scenario 2.3.4: PostgreSQL Connection Error +**Description**: PostgreSQL unavailable during export + +**Input**: +- Valid data +- PostgreSQL connection fails + +**Expected Behavior**: +- `export_data_to_postgres` raises exception +- Notification sent with `ERROR_EXPORTING_DATA_TO_POSTGRES` +- Exception propagated (workflow fails after retries) + +**Assertions**: +- Error notification sent +- Workflow fails + +--- + +#### Scenario 2.3.5: PostgreSQL Unique Constraint Violation +**Description**: Duplicate data violates unique constraint + +**Input**: +- Data with duplicate `model_id`, `timestamp`, `variable` combination +- `on_conflict: 'ignore'` configured + +**Expected Behavior**: +- PostgreSQL handles conflict with `ON CONFLICT DO NOTHING` +- `affected_rows` may be 0 for duplicates +- Workflow continues normally + +**Assertions**: +- No exception raised +- Duplicates ignored +- Workflow continues + +--- + +## 3. Activity-Specific Scenarios + +### 3.1 get_tag_values Activity + +#### Scenario 3.1.1: Success with Valid WebIds +**Input**: All webids valid and present +**Expected**: Returns list of records with timestamp, name, value, tag + +#### Scenario 3.1.2: Some WebIds are None +**Input**: Some webids in `model_tags` are `None` +**Expected**: None webids filtered out, only valid webids queried + +#### Scenario 3.1.3: DataFrame with NaN Values +**Input**: PI Web API returns DataFrame with NaN values +**Expected**: NaN values handled, data normalized correctly + +#### Scenario 3.1.4: Timestamp Normalization +**Input**: Multiple timestamps in response +**Expected**: All timestamps normalized to max timestamp value + +--- + +### 3.2 data_quality_gate Activity + +#### Scenario 3.2.1: No Filters Configured +**Input**: Empty `filters: {}` +**Expected**: Data passes through unchanged, filtered by model_tags only + +#### Scenario 3.2.2: NULL_VALUES_FILTER with DISCARD Policy +**Input**: Data with null values, policy `DISCARD` +**Expected**: Null rows removed, notification sent + +#### Scenario 3.2.3: OUT_OF_BOUNDS_FILTER with WARN Policy +**Input**: Data outside range, policy `WARN` +**Expected**: Notification sent, data kept + +#### Scenario 3.2.4: Unknown Filter Type +**Input**: Filter name not in `quality_gate_filters` +**Expected**: Warning logged, filter skipped, processing continues + +#### Scenario 3.2.5: Filter Removes All Data +**Input**: Filter that removes all rows +**Expected**: Empty DataFrame returned, processing continues + +--- + +### 3.3 aggregate_data Activity + +#### Scenario 3.3.1: Single Value Per Tag +**Input**: One data point per tag +**Expected**: Fast path returns value directly + +#### Scenario 3.3.2: Multiple Values - Latest (lts) +**Input**: Multiple points, `aggr_function: 'lts'` +**Expected**: Returns last value in sorted order + +#### Scenario 3.3.3: Multiple Values with NaN +**Input**: Some NaN values in series +**Expected**: NaN values dropped before aggregation + +#### Scenario 3.3.4: All NaN Values +**Input**: All values are NaN +**Expected**: Returns `None`, tag skipped + +#### Scenario 3.3.5: Invalid Aggregation Function +**Input**: Unknown `aggr_function` +**Expected**: Notification sent, returns `'continue'`, tag skipped + +#### Scenario 3.3.6: Empty DataFrame After Filtering +**Input**: No data after quality gate +**Expected**: Returns empty DataFrame dict + +--- + +### 3.4 group_and_hold_data Activity + +#### Scenario 3.4.1: First Run - No Existing Data +**Input**: No existing data in Redis for key +**Expected**: Creates new `data_hold` dict, stores in Redis + +#### Scenario 3.4.2: Subsequent Run - Existing Data +**Input**: Existing `data_hold` in Redis +**Expected**: Merges new data with existing, updates timestamp + +#### Scenario 3.4.3: Removed Tags Cleanup +**Input**: Tags removed from `model_tags` +**Expected**: Removed tags deleted from `data_hold` + +#### Scenario 3.4.4: Empty Input Data +**Input**: Empty DataFrame +**Expected**: Returns empty dict, warning logged + +#### Scenario 3.4.5: Redis Get Error +**Input**: Redis get operation fails +**Expected**: Notification sent, exception raised + +#### Scenario 3.4.6: Redis Set Error +**Input**: Redis set operation fails +**Expected**: Notification sent, exception raised + +--- + +### 3.5 export_data_to_postgres Activity + +#### Scenario 3.5.1: Successful Insert +**Input**: Valid data, no conflicts +**Expected**: Data inserted, `affected_rows > 0` + +#### Scenario 3.5.2: Conflict with Ignore Policy +**Input**: Duplicate data, `on_conflict: 'ignore'` +**Expected**: Duplicates ignored, `affected_rows` may be less than total + +#### Scenario 3.5.3: Conflict with Replace Policy +**Input**: Duplicate data, `on_conflict: 'replace'` +**Expected**: Duplicates updated, `affected_rows` includes updates + +#### Scenario 3.5.4: Timestamp Conversion +**Input**: String timestamps in data +**Expected**: Timestamps converted to datetime format + +#### Scenario 3.5.5: Database Connection Error +**Input**: Database unavailable +**Expected**: Exception raised, notification sent + +--- + +### 3.6 write_metrics Activity + +#### Scenario 3.6.1: Success with Valid Values +**Input**: Data with non-None values +**Expected**: Metrics written for all non-None values + +#### Scenario 3.6.2: Some None Values +**Input**: Some values are None +**Expected**: None values skipped, only non-None values written + +#### Scenario 3.6.3: All None Values +**Input**: All values are None +**Expected**: No metrics written, activity completes + +--- + +### 3.7 store_data_package Activity + +#### Scenario 3.7.1: Success +**Input**: Valid data and held_data +**Expected**: Package stored in Redis with TTL 120 + +#### Scenario 3.7.2: Redis Error +**Input**: Redis set fails +**Expected**: Notification sent, exception raised + +--- + +## 4. Integration Scenarios + +### 4.1 End-to-End Scenarios + +#### Scenario 4.1.1: Complete Happy Path +**Description**: Full workflow from API to database + +**Flow**: +1. PI Web API returns data +2. Quality gate passes +3. Aggregation succeeds +4. Redis storage succeeds +5. PostgreSQL export succeeds +6. Metrics written +7. Debug package stored (if enabled) + +**Assertions**: +- All activities called +- Data in all storage layers +- No errors + +--- + +#### Scenario 4.1.2: Partial Failure with Retry +**Description**: Activity fails, retries succeed + +**Flow**: +1. First attempt fails (e.g., Redis timeout) +2. Retry policy triggers +3. Second attempt succeeds +4. Workflow continues + +**Assertions**: +- Retry policy applied +- Workflow eventually succeeds +- Error logged but not fatal + +--- + +#### Scenario 4.1.3: Complete Failure After Retries +**Description**: Activity fails after all retries exhausted + +**Flow**: +1. Activity fails repeatedly +2. Retry policy exhausted +3. Workflow fails + +**Assertions**: +- All retries attempted +- Workflow fails with error +- Error notification sent + +--- + +## 5. Edge Cases and Boundary Conditions + +### 5.1 Data Edge Cases + +#### Scenario 5.1.1: Very Large Dataset +**Input**: Thousands of data points +**Expected**: Handles efficiently, all processed + +#### Scenario 5.1.2: Single Data Point +**Input**: One tag, one data point +**Expected**: Processes correctly + +#### Scenario 5.1.3: Extreme Values +**Input**: Very large or very small numeric values +**Expected**: Handled correctly, no overflow + +#### Scenario 5.1.4: Special Characters in Tag Names +**Input**: Tag names with special characters +**Expected**: Handled correctly + +--- + +### 5.2 Configuration Edge Cases + +#### Scenario 5.2.1: Very Short Retention Time +**Input**: `retention_time: 1` (1 second) +**Expected**: Data expires quickly but workflow completes + +#### Scenario 5.2.2: Very Long Retention Time +**Input**: `retention_time: 86400` (1 day) +**Expected**: Data persists for full duration + +#### Scenario 5.2.3: Max Count = 1 +**Input**: `max_count: 1` +**Expected**: Only latest value retrieved + +#### Scenario 5.2.4: Max Count = Large Number +**Input**: `max_count: 10000` +**Expected**: Many values retrieved and processed + +--- + +### 5.3 Concurrent Execution Scenarios + +#### Scenario 5.3.1: Multiple Workflows Same Schedule +**Input**: Two workflows with same `schedule_name` running concurrently +**Expected**: Both complete, data merged correctly in Redis + +#### Scenario 5.3.2: Multiple Workflows Different Schedules +**Input**: Multiple workflows with different `schedule_name` +**Expected**: Each uses separate Redis keys, no interference + +--- + +## 6. Performance Scenarios + +### 6.1 Load Scenarios + +#### Scenario 6.1.1: High Throughput +**Input**: Many tags, frequent execution +**Expected**: Handles load efficiently + +#### Scenario 6.1.2: Large Payload +**Input**: Large amount of data per tag +**Expected**: Processes within timeout limits + +--- + +## 7. Test Data Requirements + +### 7.1 Valid Test Data Structure + +```python +{ + 'model_name': 'test_model', + 'model_id': 'test_model_id', + 'schedule_name': 'test_schedule', + 'model_tags': { + 'tag1': { + 'webid': 'webid1', + 'aggr_function': 'avg', + 'data_range': [0, 100], + 'frequency': 60000, + }, + }, + 'trigger_laborious': False, + 'filters': {}, + 'schema': 'test_schema', + 'table_name': 'test_table', + 'retention_time': 3600, + 'fill_missing_tags': False, + 'debug_data_package': False, + 'pi_web_api_query': { + 'endpoint': '/streamsets/recorded', + 'period': '*-1d', + 'max_count': 10, + 'api_timeout': 30, + }, +} +``` + +### 7.2 Mock PI Web API Response + +```python +DataFrame({ + 'timestamp': ['2024-01-01 12:00:00+0000', ...], + 'name': ['tag1', 'tag2', ...], + 'value': [10.5, 20.3, ...], + 'tag': ['webid1', 'webid2', ...], +}) +``` + +--- + +## 8. Test Implementation Notes + +### 8.1 Test Organization + +- Group tests by scenario category +- Use descriptive test names matching scenario IDs +- Share fixtures for common setup +- Use parametrized tests for similar scenarios + +### 8.2 Assertions Checklist + +For each scenario, verify: +- [ ] Correct activities called +- [ ] Correct parameters passed +- [ ] Expected data in storage (SQLite/Redis) +- [ ] Expected notifications sent +- [ ] Expected metrics written +- [ ] No unexpected errors +- [ ] Workflow state correct + +### 8.3 Mock Configuration + +- Mock PI Web API client responses +- Use fake Redis (fakeredis) +- Use fake MongoDB (mongomock) +- Use SQLite for PostgreSQL +- Mock notification handler +- Mock metrics controller + +--- + +## 9. Priority Scenarios + +### High Priority (Must Test) +1. Scenario 1.1.1: Happy Path +2. Scenario 1.2.1: Empty Data +3. Scenario 1.3.1: API Connection Error +4. Scenario 2.1.1: Complete Processing +5. Scenario 2.2.1: Empty After Grouping +6. Scenario 2.3.4: PostgreSQL Error + +### Medium Priority (Should Test) +1. Scenario 1.1.3: Debug Package +2. Scenario 2.1.2: Quality Filters +3. Scenario 2.1.3: Different Aggregations +4. Scenario 3.3.5: Invalid Aggregation +5. Scenario 3.5.2: Conflict Ignore + +### Low Priority (Nice to Have) +1. Scenario 4.1.2: Retry Success +2. Scenario 5.1.1: Large Dataset +3. Scenario 5.3.1: Concurrent Execution + diff --git a/e2e/test_pi_web_api_scouter.py b/e2e/test_pi_web_api_scouter.py new file mode 100644 index 0000000..71f002f --- /dev/null +++ b/e2e/test_pi_web_api_scouter.py @@ -0,0 +1,106 @@ +""" +End-to-end tests for PI Web API Scouter workflow. +""" + +from datetime import datetime + +import pytest +from sqlalchemy import inspect, text +from temporalio.testing import WorkflowEnvironment +from temporalio.worker import Worker + +from scouter.activities.activities import Activities +from scouter.workflow.pi_web_api_scouter import PIWebAPIScouter + + +@pytest.mark.asyncio +@pytest.mark.integration +async def test_pi_web_api_scouter_e2e( + temporal_test_env: WorkflowEnvironment, + temporal_worker: Worker, + test_activities: Activities, + mock_pi_web_api_client, + postgres_engine, +): + """ + End-to-end test for PI Web API Scouter workflow. + + This test: + 1. Starts the workflow with test data + 2. Verifies PI Web API is called + 3. Verifies data flows through CoreScouter + 4. Verifies data is stored in PostgreSQL (schema: sientia_data, table: laborious_data) + 5. Verifies data is cached in Redis + """ + client = temporal_test_env.client + + # Prepare test input + input_data = { + 'model_name': 'PI Web API Scouter Test Model', + 'model_id': '1', + 'schedule_name': 'pi-web-api-scouter-test', + 'model_tags': { + 'tag1': { + 'webid': 'webid1', + 'aggr_function': 'avg', + 'data_range': [0, 100], + 'frequency': 60000, + }, + 'tag2': { + 'webid': 'webid2', + 'aggr_function': 'avg', + 'data_range': [0, 100], + 'frequency': 60000, + }, + }, + 'trigger_laborious': False, + 'filters': {}, + 'schema': 'sientia_data', + 'table_name': 'laborious_data', + 'retention_time': 3600, + 'fill_missing_tags': False, + 'pi_web_api_query': { + 'endpoint': '/streamsets/recorded', + 'period': '*-1d', + 'max_count': 10, + 'api_timeout': 30, + }, + } + + # Start workflow + handle = await client.start_workflow( + PIWebAPIScouter.run, + input_data, + id=f'test-workflow-{datetime.now().timestamp()}', + task_queue='test-queue', + ) + + # Wait for workflow completion + await handle.result() + + # Verify PI Web API was called + mock_pi_web_api_client.get_latest_values_df.assert_called_once() + + # Verify data was stored in PostgreSQL + inspector = inspect(postgres_engine) + + # Schema and table are created by the setup_postgres_schema_and_table fixture + schema_name = 'sientia_data' + table_name = 'laborious_data' + full_table_name = f"{schema_name}.{table_name}" + + # Check if table exists in the schema + table_exists = inspector.has_table(table_name, schema=schema_name) + + assert table_exists, f"Expected table {full_table_name} to exist in PostgreSQL" + + # Verify data was inserted + with postgres_engine.connect() as conn: + result = conn.execute(text(f"SELECT COUNT(*) FROM {full_table_name}")) + row_count = result.scalar() + + assert row_count > 0, f"Expected data in PostgreSQL table {full_table_name}, got {row_count} rows" + + # Verify data was cached in Redis + keys = await test_activities.redis_repository.keys('*') + assert len(keys) > 0, "Expected data in Redis" diff --git a/requirements-dev.txt b/requirements-dev.txt index 56ab376..352add1 100644 --- a/requirements-dev.txt +++ b/requirements-dev.txt @@ -14,6 +14,11 @@ pytest>=7.4.0 # Testing framework pytest-cov>=4.1.0 # Coverage plugin for pytest pytest-asyncio>=0.21.0 # Async test support (already in main requirements) +# E2E Testing Dependencies +fakeredis>=2.20.0 # In-memory Redis server for testing +mongomock>=4.1.2 # In-memory MongoDB for testing + # Development Tools ipython>=8.12.0 # Enhanced Python shell -ipdb>=0.13.13 # IPython debugger \ No newline at end of file +ipdb>=0.13.13 # IPython debugger +testcontainers[postgres] \ No newline at end of file