From b4827935a7775381872e8dd99401aea7a3cc65f6 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Thu, 31 Jul 2025 14:43:34 -0300 Subject: [PATCH] SIENTIAPDE-1174 Integrate write_metrics activity into CoreScouter for enhanced metric logging and update worker.py to include write_metrics in workflow tasks. --- scouter/worker/worker.py | 3 ++- scouter/workflow/sub_workflows/core_scouter.py | 13 ++++++++----- 2 files changed, 10 insertions(+), 6 deletions(-) diff --git a/scouter/worker/worker.py b/scouter/worker/worker.py index 62d2752..aff8341 100644 --- a/scouter/worker/worker.py +++ b/scouter/worker/worker.py @@ -84,7 +84,8 @@ async def main(): activities.aggregate_data, activities.group_and_hold_data, activities.export_data_to_postgres, - activities.store_data_package + activities.write_metrics, + activities.store_data_package, ], max_concurrent_workflow_tasks=100, max_concurrent_activities=100, diff --git a/scouter/workflow/sub_workflows/core_scouter.py b/scouter/workflow/sub_workflows/core_scouter.py index c71f10b..3a9c1e6 100644 --- a/scouter/workflow/sub_workflows/core_scouter.py +++ b/scouter/workflow/sub_workflows/core_scouter.py @@ -88,11 +88,14 @@ class CoreScouter: start_to_close_timeout=timedelta(seconds=60) ) - metrics.LABORIOUS_DATA_WRITTEN_COUNT.labels( - pod_id=metadata['pod_id'], - model_name=input_data['model_name'], - pipeline_name=input_data['workflow_name'] - ).inc() + await workflow.execute_activity_method( + Activities.write_metrics, + { + **metadata, + }, + retry_policy=retry_policy, + start_to_close_timeout=timedelta(seconds=60) + ) if input_data.get('debug_data_package', False): await workflow.execute_activity_method(