From 14e13666582c0ee3a67f5f53f0d1da2813342984 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Mon, 14 Jul 2025 16:24:03 -0300 Subject: [PATCH] SIENTIAPDE-1150 refactor: update orchestration workflow to conditionally start activity methods for pipeline timestamps based on report availability, enhancing efficiency and clarity --- orchestrator/workflows/orchestrator.py | 54 ++++++++++++++------------ 1 file changed, 30 insertions(+), 24 deletions(-) diff --git a/orchestrator/workflows/orchestrator.py b/orchestrator/workflows/orchestrator.py index 7053746..574f627 100644 --- a/orchestrator/workflows/orchestrator.py +++ b/orchestrator/workflows/orchestrator.py @@ -199,32 +199,38 @@ class Orchestrator: start_to_close_timeout=timedelta(seconds=60) ) - update_pipelines_timestamps_handler = workflow.start_activity_method( - Activities.update_pipelines_timestamps, - { - 'updated_pipelines': schedule_update_report - }, - retry_policy=retry_policy, - start_to_close_timeout=timedelta(seconds=60) - ) + if schedule_update_report: - create_pipelines_timestamps_handler = workflow.start_activity_method( - Activities.create_pipelines_timestamps, - { - 'created_pipelines': schedule_insertion_report - }, - retry_policy=retry_policy, - start_to_close_timeout=timedelta(seconds=60) - ) + update_pipelines_timestamps_handler = workflow.start_activity_method( + Activities.update_pipelines_timestamps, + { + 'updated_pipelines': schedule_update_report + }, + retry_policy=retry_policy, + start_to_close_timeout=timedelta(seconds=60) + ) - delete_pipelines_timestamps_handler = workflow.start_activity_method( - Activities.delete_pipelines_timestamps, - { - 'deleted_pipelines': schedule_deletion_report - }, - retry_policy=retry_policy, - start_to_close_timeout=timedelta(seconds=60) - ) + if schedule_insertion_report: + + create_pipelines_timestamps_handler = workflow.start_activity_method( + Activities.create_pipelines_timestamps, + { + 'created_pipelines': schedule_insertion_report + }, + retry_policy=retry_policy, + start_to_close_timeout=timedelta(seconds=60) + ) + + if schedule_deletion_report: + + delete_pipelines_timestamps_handler = workflow.start_activity_method( + Activities.delete_pipelines_timestamps, + { + 'deleted_pipelines': schedule_deletion_report + }, + retry_policy=retry_policy, + start_to_close_timeout=timedelta(seconds=60) + ) await schedule_report_handler await slot_report_handler