from typing import Sequence, Type, Any from temporalio.worker import Worker, PollerBehaviorAutoscaling from temporalio.client import Client import os 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'), ('ACTIVITY_POLLER_BEHAVIUR_MINIMUM', '10'), ('ACTIVITY_POLLER_BEHAVIUR_INITIAL', '100'), ('ACTIVITY_POLLER_BEHAVIUR_MAXIMUM', '200'), ] def prepare_worker( main_workflow: Type, other_workflows: Sequence[Type], activities: Sequence[Any], temporal_client: Client, ): 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])) return [ Worker( temporal_client, task_queue='scouter-queue', workflows=[main_workflow, *other_workflows], 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_cached_workflows=local_workflow_parameters['MAX_CACHED_WORKFLOWS'], workflow_task_poller_behavior=PollerBehaviorAutoscaling( minimum=local_workflow_parameters['WORKFLOW_POLLER_BEHAVIUR_MINIMUM'], initial=local_workflow_parameters['WORKFLOW_POLLER_BEHAVIUR_INITIAL'], maximum=local_workflow_parameters['WORKFLOW_POLLER_BEHAVIUR_MAXIMUM'], ), activity_task_poller_behavior=PollerBehaviorAutoscaling( minimum=local_workflow_parameters['ACTIVITY_POLLER_BEHAVIUR_MINIMUM'], initial=local_workflow_parameters['ACTIVITY_POLLER_BEHAVIUR_INITIAL'], maximum=local_workflow_parameters['ACTIVITY_POLLER_BEHAVIUR_MAXIMUM'], ), ) ]