diff --git a/scouter/activities/redis.py b/scouter/activities/redis.py index bffc478..f72059d 100644 --- a/scouter/activities/redis.py +++ b/scouter/activities/redis.py @@ -111,6 +111,7 @@ class Redis(RedisBase): metadata=metadata ) data = DataFrame(input_data['data']) + model_tags = input_data['model_tags'] retention_time = input_data['retention_time'] key = f"held_data_{input_data['workflow_name']}_{input_data['schedule_name']}" @@ -138,6 +139,19 @@ class Redis(RedisBase): try: + # Remove possibly removed tags + tags = list(model_tags.keys()) + self.debug( + f"Tags to keep: {tags}", + metadata=metadata + ) + data_hold = {tag: content for tag, + content in data_hold.items() if tag in tags} + self.debug( + f"Data hold after removing removed tags: {data_hold}", + metadata=metadata + ) + to_register_metrics = [] for _, row in data.iterrows(): value = row['value'] @@ -152,6 +166,10 @@ class Redis(RedisBase): self.set(key, data_hold, ttl=retention_time) # Register metrics + self.debug( + f"Metrics to register: {to_register_metrics}", + metadata=metadata + ) for metric in to_register_metrics: metrics.TAG_CHANGES_MONITOR.labels( pod_id=self.pod_id, diff --git a/scouter/worker/worker.py b/scouter/worker/worker.py index aff8341..6aa61ea 100644 --- a/scouter/worker/worker.py +++ b/scouter/worker/worker.py @@ -1,5 +1,6 @@ from temporalio import workflow, client from temporalio.worker import Worker +from temporalio.runtime import Runtime, TelemetryConfig, PrometheusConfig with workflow.unsafe.imports_passed_through(): import sys @@ -21,6 +22,7 @@ with workflow.unsafe.imports_passed_through(): ) POD_ID = os.getenv("HOSTNAME", "localhost") +SDK_METRICS_PORT = int(os.getenv('HTTP_SDK_METRICS_PORT', "9091")) async def main(): @@ -62,11 +64,21 @@ async def main(): 'KAFKA_BOOTSTRAP_SERVERS', 'localhost:9092') ) + logger.info(f'Starting SDK Metrics Server on port {SDK_METRICS_PORT}...') + + new_runtime = Runtime( + telemetry=TelemetryConfig( + metrics=PrometheusConfig( + bind_address=f"0.0.0.0:{SDK_METRICS_PORT}") + ) + ) + logger.info('Starting Temporal Client...') temporal_client = await client.Client.connect( target_host=host, - namespace=os.getenv('TEMPORAL_NAMESPACE', 'default') + namespace=os.getenv('TEMPORAL_NAMESPACE', 'laborious'), + runtime=new_runtime ) logger.info('Starting Workers...') diff --git a/scouter/workflow/sub_workflows/core_scouter.py b/scouter/workflow/sub_workflows/core_scouter.py index b9f4a35..d686d7f 100644 --- a/scouter/workflow/sub_workflows/core_scouter.py +++ b/scouter/workflow/sub_workflows/core_scouter.py @@ -65,6 +65,7 @@ class CoreScouter: 'workflow_name': input_data['workflow_name'], 'data': grouped_data, 'model_id': input_data['model_id'], + 'model_tags': input_data['model_tags'], 'retention_time': input_data['retention_time'] }, retry_policy=retry_policy, diff --git a/tests/activities/test_redis.py b/tests/activities/test_redis.py index bed4c10..f68e48f 100644 --- a/tests/activities/test_redis.py +++ b/tests/activities/test_redis.py @@ -219,7 +219,11 @@ async def test_group_and_hold_data_new_key(redis_activity): 'name': ['sensor1', 'sensor2'], 'value': [25.5, 30.0], 'timestamp': ['2023-01-01 12:00:00'] * 2 - }).to_dict('records') + }).to_dict('records'), + 'model_tags': { + 'sensor1': 'sensor1', + 'sensor2': 'sensor2' + } } # Mock get to return None for new key @@ -271,7 +275,12 @@ async def test_group_and_hold_data_update_existing(redis_activity): 'name': ['sensor1', 'sensor3'], 'value': [25.5, 42.0], 'timestamp': ['2023-01-01 12:00:00'] * 2 - }).to_dict('records') + }).to_dict('records'), + 'model_tags': { + 'sensor1': 'sensor1', + 'sensor2': 'sensor2', + 'sensor3': 'sensor3' + } } # Mock get to return existing data @@ -317,7 +326,11 @@ async def test_group_and_hold_data_with_none_values(redis_activity): 'name': ['sensor1', 'sensor2'], 'value': [None, 30.0], 'timestamp': [datetime(2023, 1, 1, 12, 0, 0)] * 2 - }).to_dict('records') + }).to_dict('records'), + 'model_tags': { + 'sensor1': 'sensor1', + 'sensor2': 'sensor2' + } } # Mock get to return None for new key @@ -341,7 +354,11 @@ async def test_group_and_hold_data_empty_dataframe(redis_activity): 'workflow_name': 'test_workflow', 'schedule_name': 'test_schedule', 'retention_time': 3600, - 'data': DataFrame(columns=['name', 'value', 'timestamp']).to_dict('records') + 'data': DataFrame(columns=['name', 'value', 'timestamp']).to_dict('records'), + 'model_tags': { + 'sensor1': 'sensor1', + 'sensor2': 'sensor2' + } } redis_activity.get = MagicMock(return_value=None) @@ -361,7 +378,11 @@ async def test_group_and_hold_data_error_get(redis_activity): 'schedule_name': 'test_schedule', 'retention_time': 3600, 'model_id': 1, - 'data': DataFrame(columns=['name', 'value', 'timestamp']).to_dict('records') + 'data': DataFrame(columns=['name', 'value', 'timestamp']).to_dict('records'), + 'model_tags': { + 'sensor1': 'sensor1', + 'sensor2': 'sensor2' + } } redis_activity.get = MagicMock(side_effect=Exception('test')) @@ -395,7 +416,11 @@ async def test_group_and_hold_data_error_set(redis_activity): 'schedule_name': 'test_schedule', 'retention_time': 3600, 'model_id': 1, - 'data': DataFrame(columns=['name', 'value', 'timestamp']).to_dict('records') + 'data': DataFrame(columns=['name', 'value', 'timestamp']).to_dict('records'), + 'model_tags': { + 'sensor1': 'sensor1', + 'sensor2': 'sensor2' + } } existing_data = { @@ -434,7 +459,11 @@ async def test_store_data_package(redis_activity): 'name': ['sensor1', 'sensor2'], 'value': [25.5, 30.0], 'timestamp': ['2023-01-01 12:00:00'] * 2 - }).to_dict() + }).to_dict(), + 'model_tags': { + 'sensor1': 'sensor1', + 'sensor2': 'sensor2' + } } await redis_activity.store_data_package(test_data) @@ -467,7 +496,11 @@ async def test_store_data_package_error(redis_activity): 'name': ['sensor1', 'sensor2'], 'value': [25.5, 30.0], 'timestamp': ['2023-01-01 12:00:00'] * 2 - }).to_dict() + }).to_dict(), + 'model_tags': { + 'sensor1': 'sensor1', + 'sensor2': 'sensor2' + } } with pytest.raises(Exception): diff --git a/tests/workflow/sub_workflows/test_core_scouter.py b/tests/workflow/sub_workflows/test_core_scouter.py index e1be284..5e88c5b 100644 --- a/tests/workflow/sub_workflows/test_core_scouter.py +++ b/tests/workflow/sub_workflows/test_core_scouter.py @@ -80,7 +80,8 @@ async def test_core_scouter_workflow_success(mock_workflow, core_scouter): 'schedule_name': 'test_schedule', 'data': 'grouped_data', 'model_id': 'test_model_id', - 'retention_time': 3600 + 'retention_time': 3600, + 'model_tags': {} }, retry_policy=ANY, start_to_close_timeout=ANY @@ -184,7 +185,8 @@ async def test_core_scouter_workflow_with_empty_data(mock_workflow, core_scouter 'schedule_name': 'test_schedule', 'data': {}, 'model_id': 'test_model_id', - 'retention_time': 3600 + 'retention_time': 3600, + 'model_tags': {} }, retry_policy=ANY, start_to_close_timeout=ANY diff --git a/values.yaml b/values.yaml index a729358..b54199c 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.3.2" + tag: "0.4.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: @@ -112,6 +112,13 @@ tolerations: [] affinity: {} services: + sdk-metrics: + enabled: true + type: ClusterIP + port: 9091 + targetPort: 9091 + name: sdk-metrics + metrics: enabled: true type: ClusterIP @@ -125,25 +132,25 @@ serviceMonitor: # Se true, um recurso ServiceMonitor será criado. enabled: true # O intervalo no qual as métricas devem ser coletadas (ex: 30s, 1m). - interval: 30s - # O path do endpoint de métricas na sua aplicação. - path: /metrics - # Labels adicionais para o recurso ServiceMonitor. - # Essencial para que o Prometheus Operator o descubra. Se você usa o helm chart kube-prometheus-stack, - # ele procura por ServiceMonitors com o label "release: kube-prometheus-stack". + endpoints: + - port: metrics + path: /metrics + interval: 30s + relabelings: [] + - port: sdk-metrics + path: /metrics + interval: 30s + relabelings: [] + additionalLabels: release: kube-prometheus-stack - # Configurações de relabeling adicionais, se necessário. - # ref: https://prometheus.io/docs/prometheus/latest/configuration/configuration/#relabel_config - relabelings: [] - port: metrics env: # Entrypoint variables - name: GITHUB_REPO_URL value: "git@github.com:Aignosi/sientia-dataops-scouter_temporal.git" - name: GITHUB_BRANCH - value: "SIENTIAPDE-1174-mapear-e-implementar-metricas-a-serem-criadas" + value: "SIENTIAPDE-1169-pensar-e-projetar-testes-de-breakdown-e-performance" - name: PYTHON_APP value: "scouter.worker.worker" @@ -217,7 +224,7 @@ ssh: # kubectl create secret docker-registry docker-hub-secret --namespace sientia --docker-server=http://aignosi.azurecr.io --docker-username=aignosi --docker-password=5I5zpQ6sRaHqX1hD3dr+2mo647yO3FRc359/wu6gsP+ACRDRz5mp -# helm upgrade --install sientia-scouter-worker sientia/sientia-module -n sientia --create-namespace -f ./values.yaml --version 0.4.0 +# helm upgrade --install sientia-scouter-worker sientia/sientia-module -n sientia --create-namespace -f ./values.yaml --version 0.5.0 # kubectl create secret generic git-ssh-key-sientia-scouter-worker \ # --namespace sientia \