diff --git a/laborious/activities/gates.py b/laborious/activities/gates.py index c4f9dd5..1cf3fb9 100644 --- a/laborious/activities/gates.py +++ b/laborious/activities/gates.py @@ -73,7 +73,7 @@ class Gates(BaseActivity): filter_output = [] - self.logger.debug(f"Input data:\n {data.to_string()}") + self.logger.debug(f"Input data:\n {data}") self.logger.debug(f"Filters: {filters}") for fil, config in filters.items(): diff --git a/laborious/activities/mlflow.py b/laborious/activities/mlflow.py index ed96192..ae141a2 100644 --- a/laborious/activities/mlflow.py +++ b/laborious/activities/mlflow.py @@ -41,6 +41,7 @@ class MLFlow(BaseActivity): model_name = input_data['model_name'] model_retention = input_data['model_retention'] + self.logger.debug("Raw input data:") self.logger.debug(data) data = data.pivot( @@ -50,9 +51,13 @@ class MLFlow(BaseActivity): data.reset_index(inplace=True) data.columns.name = None + self.logger.debug("Processed input data:") + self.logger.debug(data) + response_data = self.model_monitoring_repository.transform( model_name, data, model_retention) + self.logger.debug("Response data:") self.logger.debug(response_data) return response_data diff --git a/laborious/utils/repository/model_repository.py b/laborious/utils/repository/model_repository.py index 3896411..f2eaeb9 100644 --- a/laborious/utils/repository/model_repository.py +++ b/laborious/utils/repository/model_repository.py @@ -13,7 +13,6 @@ import traceback import mlflow import pandas as pd from sientia.ModelServing import ModelServing -from pathlib import Path class MLFlowRepository(): @@ -260,7 +259,8 @@ class MLFlowRepository(): try: return { 'success': True, - 'content': self.model_serving.get_cached_transform(model_name, data, model_retention).to_dict() + 'content': self.model_serving.get_cached_transform( + model_name, data, model_retention).to_dict() } except Exception as e: @@ -276,7 +276,8 @@ class MLFlowRepository(): try: start_time = datetime.now() data = self.model_serving.get_cached_predict( - model_name, data, model_retention) + model_name, data, model_retention)[-1:] + end_time = datetime.now() data = pd.DataFrame(data, columns=['prediction']) data['response_time'] = (end_time - start_time).total_seconds() diff --git a/laborious/workflows/sub_workflows/format_and_export_prediction.py b/laborious/workflows/sub_workflows/format_and_export_prediction.py index 67cad49..3b1fad7 100644 --- a/laborious/workflows/sub_workflows/format_and_export_prediction.py +++ b/laborious/workflows/sub_workflows/format_and_export_prediction.py @@ -40,8 +40,6 @@ class FormatAndExportPrediction(): data = input_data['data'] prediction_confidence = input_data['prediction_confidence'] - print(f"Input data: {input_data}") - if path_flag is None: # proceed with formatting and exporting prediction = await workflow.execute_local_activity_method( diff --git a/laborious/workflows/sub_workflows/prediction_process.py b/laborious/workflows/sub_workflows/prediction_process.py index de84328..3164639 100644 --- a/laborious/workflows/sub_workflows/prediction_process.py +++ b/laborious/workflows/sub_workflows/prediction_process.py @@ -97,11 +97,13 @@ class PredictionProcess(): ): return + transformed_data = response_data['content'] + path_flag, confidence, comment = await workflow.execute_local_activity_method( Activities.mlflow_content_gate, { 'filters': input_data['mlflow_transform_filters'], - 'data': response_data, + 'data': transformed_data, 'type': 'transform', 'path_priority': input_data['path_priority'] }, @@ -117,7 +119,7 @@ class PredictionProcess(): response_data = await workflow.execute_local_activity_method( Activities.request_predict, { - 'data': response_data, + 'data': transformed_data, 'model_name': model_name, 'model_retention': model_retention }, @@ -148,11 +150,14 @@ class PredictionProcess(): 'path_flag': path_flag, 'data': response_data['content'], 'prediction_confidence': confidence, - 'timestamp': response_data['timestamp'], + 'timestamp': last_timestamp, 'model_id': model_id, 'model_name': model_name, 'model_retention': model_retention, - 'opc_output_config': input_data['opc_output_config'] + 'opc_output_config': input_data['opc_output_config'], + 'schema': input_data['schema'], + 'table_name': input_data['table_name'], + 'comment': comment } ) diff --git a/simulator/redis-feeder.py b/simulator/redis-feeder.py deleted file mode 100644 index 3f7e36a..0000000 --- a/simulator/redis-feeder.py +++ /dev/null @@ -1,55 +0,0 @@ -import redis -import json -import os - -# Redis connection settings -redis_host = "localhost" -redis_port = 6379 - -# Connect to Redis -r = redis.Redis(host=redis_host, port=redis_port, - decode_responses=True, username='default', password='bdnZOpcyiL') - -# Define the key pattern to target -pattern = "slot:opc_tags:*" - -# Step 1: Find and delete matching keys -print("🔍 Searching for keys matching:", pattern) -for key in r.scan_iter(match=pattern): - r.delete(key) - print(f"❌ Deleted: {key}") - -# Step 2: Insert new data -# Example new OPC tag data -new_data = { - "slot:opc_tags:1": { - "server1": { - "name": "server1", - "url": "opc.tcp://sientia-opc-simulator-service.sientia-opc.svc.cluster.local:4840", - "server_uri": "http://opcua-server.simulator", - "tags": { - 'ns=2;i=2': { - 'tag_name': 'Counter', - 'frequency': 1000, - 'topics': ['opcua', 'counter'], - }, - 'ns=2;i=3': { - 'tag_name': 'Rollout', - 'frequency': 1000, - "topics": ['opcua', 'rollout'], - }, - 'ns=2;i=4': { - 'tag_name': 'Square', - 'frequency': 1000, - "topics": ['opcua'], - }, - } - } - } -} - -for key, val in new_data.items(): - r.set(key, json.dumps(val)) - print(f"✅ Set: {key} -> {val}") - -print("🚀 OPC tag keys replaced successfully.")