diff --git a/scouter/worker/worker.py b/scouter/worker/worker.py index 9bdfcc2..f636d59 100644 --- a/scouter/worker/worker.py +++ b/scouter/worker/worker.py @@ -30,12 +30,20 @@ async def main(): logger = get_logger(__name__) - logger.info(f"Starting Worker with pod_id: {POD_ID}") + metadata = { + 'pod_id': POD_ID, + 'model_name': '-', + 'model_id': '-', + 'workflow_name': '-', + 'schedule_name': '-', + } - logger.info("Starting prometheus client...") + logger.custom_info(f"Starting Worker with pod_id: {POD_ID}", metadata) + + logger.custom_info("Starting prometheus client...", metadata) start_prometheus_server() - logger.info('Starting Notification Handler...') + logger.custom_info('Starting Notification Handler...', metadata) mongo_config = build_mongodb_config() notification_handler = NotificationHandler( @@ -45,7 +53,7 @@ async def main(): project_name=os.getenv('PROJECT_NAME', 'scouter') ) - logger.info('Starting Activities...') + logger.custom_info('Starting Activities...', metadata) activities = Activities( logger=logger, @@ -55,7 +63,7 @@ async def main(): mongodb_config=build_mongodb_config() ) - logger.info('Starting Faker Activities...') + logger.custom_info('Starting Faker Activities...', metadata) faker_activities = Faker( logger=logger, @@ -64,7 +72,8 @@ async def main(): 'KAFKA_BOOTSTRAP_SERVERS', 'localhost:9092') ) - logger.info(f'Starting SDK Metrics Server on port {SDK_METRICS_PORT}...') + logger.custom_info( + f'Starting SDK Metrics Server on port {SDK_METRICS_PORT}...', metadata) new_runtime = Runtime( telemetry=TelemetryConfig( @@ -73,7 +82,7 @@ async def main(): ) ) - logger.info('Starting Temporal Client...') + logger.custom_info('Starting Temporal Client...', metadata) temporal_client = await client.Client.connect( target_host=host, @@ -81,7 +90,7 @@ async def main(): runtime=new_runtime ) - logger.info('Starting Workers...') + logger.custom_info('Starting Workers...', metadata) workers = [ Worker( @@ -120,13 +129,14 @@ async def main(): for w in workers: handlers.append(w.run()) - logger.info('Workers started successfully') + logger.custom_info('Workers started successfully', metadata) try: await asyncio.gather(*handlers) except BaseException as e: # NOSONAR - logger.error("An unhandled exception occurred: %s", e, exc_info=True) + logger.custom_error("An unhandled exception occurred: %s", + e, exc_info=True, metadata=metadata) finally: if notification_handler: notification_handler.shutdown()