From 666a0a6e13d76f2fca016c3c8795fb7872cdef56 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Tue, 16 Dec 2025 16:01:47 -0300 Subject: [PATCH] SIENTIAPDE-1441 Refactor prepare_worker.py to improve type annotations and code clarity. Update import statements and enhance formatting for better readability. Adjust comments in worker.py for consistency. --- scouter/worker/prepare_worker.py | 26 +++++++++++++------------- scouter/worker/worker.py | 3 +-- 2 files changed, 14 insertions(+), 15 deletions(-) diff --git a/scouter/worker/prepare_worker.py b/scouter/worker/prepare_worker.py index b261f68..09fced1 100644 --- a/scouter/worker/prepare_worker.py +++ b/scouter/worker/prepare_worker.py @@ -1,17 +1,15 @@ -from typing import Sequence, Type, Any - -from temporalio.worker import Worker, PollerBehaviorAutoscaling -from temporalio.client import Client - import os +from collections.abc import Sequence +from typing import Any +from temporalio.client import Client +from temporalio.worker import PollerBehaviorAutoscaling, Worker parameters = [ ('MAX_CONCURRENT_WORKFLOW_TASKS', '200'), ('MAX_CONCURRENT_ACTIVITIES', '200'), ('MAX_CONCURRENT_LOCAL_ACTIVITIES', '200'), ('MAX_CACHED_WORKFLOWS', '200'), - ('WORKFLOW_POLLER_BEHAVIUR_MINIMUM', '10'), ('WORKFLOW_POLLER_BEHAVIUR_INITIAL', '100'), ('WORKFLOW_POLLER_BEHAVIUR_MAXIMUM', '200'), @@ -20,21 +18,21 @@ parameters = [ ('ACTIVITY_POLLER_BEHAVIUR_MAXIMUM', '200'), ] + def prepare_worker( - main_workflow: Type, - other_workflows: Sequence[Type], + main_workflow: type, + other_workflows: Sequence[type], activities: Sequence[Any], temporal_client: Client, ) -> Worker: - main_workflow_name = main_workflow.__name__.upper() local_workflow_parameters = {} for parameter in parameters: local_workflow_parameters[parameter[0]] = int( - os.getenv(main_workflow_name + '_' + parameter[0], parameter[1])) - + os.getenv(main_workflow_name + '_' + parameter[0], parameter[1]) + ) return Worker( temporal_client, @@ -43,7 +41,9 @@ def prepare_worker( activities=[*activities], max_concurrent_workflow_tasks=local_workflow_parameters['MAX_CONCURRENT_WORKFLOW_TASKS'], max_concurrent_activities=local_workflow_parameters['MAX_CONCURRENT_ACTIVITIES'], - max_concurrent_local_activities=local_workflow_parameters['MAX_CONCURRENT_LOCAL_ACTIVITIES'], + max_concurrent_local_activities=local_workflow_parameters[ + 'MAX_CONCURRENT_LOCAL_ACTIVITIES' + ], max_cached_workflows=local_workflow_parameters['MAX_CACHED_WORKFLOWS'], workflow_task_poller_behavior=PollerBehaviorAutoscaling( minimum=local_workflow_parameters['WORKFLOW_POLLER_BEHAVIUR_MINIMUM'], @@ -55,4 +55,4 @@ def prepare_worker( initial=local_workflow_parameters['ACTIVITY_POLLER_BEHAVIUR_INITIAL'], maximum=local_workflow_parameters['ACTIVITY_POLLER_BEHAVIUR_MAXIMUM'], ), - ) \ No newline at end of file + ) diff --git a/scouter/worker/worker.py b/scouter/worker/worker.py index d5f6be8..cd333f4 100644 --- a/scouter/worker/worker.py +++ b/scouter/worker/worker.py @@ -1,6 +1,5 @@ from temporalio import client, workflow from temporalio.runtime import PrometheusConfig, Runtime, TelemetryConfig -from temporalio.worker import PollerBehaviorAutoscaling, Worker from scouter.worker.prepare_worker import prepare_worker @@ -29,7 +28,7 @@ SDK_METRICS_PORT = int(os.getenv('HTTP_SDK_METRICS_PORT', '9091')) # For optmized latency, Temporal docs recommends fixed slots, ensuring -# high concurency levels. +# high concurency levels. MAX_CONCURRENT_WORKFLOW_TASKS = int(os.getenv('MAX_CONCURRENT_WORKFLOW_TASKS', '200')) MAX_CONCURRENT_ACTIVITIES = int(os.getenv('MAX_CONCURRENT_ACTIVITIES', '200'))