From 73bbd0a4460e808ba1e385293ec23bad16114fe8 Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Wed, 2 Apr 2025 17:40:25 -0700 Subject: [PATCH] emit lock acquired and released events --- .../db_transaction_queue/pod_lock_manager.py | 36 +++++++++++++++++++ 1 file changed, 36 insertions(+) diff --git a/litellm/proxy/db/db_transaction_queue/pod_lock_manager.py b/litellm/proxy/db/db_transaction_queue/pod_lock_manager.py index 5b640033a0a..3f63afe62a8 100644 --- a/litellm/proxy/db/db_transaction_queue/pod_lock_manager.py +++ b/litellm/proxy/db/db_transaction_queue/pod_lock_manager.py @@ -1,15 +1,20 @@ +import asyncio import uuid from typing import TYPE_CHECKING, Any, Optional from litellm._logging import verbose_proxy_logger +from litellm._service_logger import ServiceLogging from litellm.caching.redis_cache import RedisCache from litellm.constants import DEFAULT_CRON_JOB_LOCK_TTL_SECONDS +from litellm.types.services import ServiceTypes if TYPE_CHECKING: ProxyLogging = Any else: ProxyLogging = Any +service_logger_obj = ServiceLogging() # used for tracking current pod lock status + class PodLockManager: """ @@ -57,6 +62,7 @@ class PodLockManager: self.pod_id, self.cronjob_id, ) + return True else: # Check if the current pod already holds the lock @@ -70,6 +76,7 @@ class PodLockManager: self.pod_id, self.cronjob_id, ) + self._emit_acquired_lock_event(self.cronjob_id, self.pod_id) return True return False except Exception as e: @@ -104,6 +111,7 @@ class PodLockManager: self.pod_id, self.cronjob_id, ) + self._emit_released_lock_event(self.cronjob_id, self.pod_id) else: verbose_proxy_logger.debug( "Pod %s failed to release Redis lock for cronjob_id=%s", @@ -127,3 +135,31 @@ class PodLockManager: verbose_proxy_logger.error( f"Error releasing Redis lock for {self.cronjob_id}: {e}" ) + + @staticmethod + def _emit_acquired_lock_event(cronjob_id: str, pod_id: str): + asyncio.create_task( + service_logger_obj.async_service_success_hook( + service=ServiceTypes.POD_LOCK_MANAGER, + duration=DEFAULT_CRON_JOB_LOCK_TTL_SECONDS, + call_type="_emit_acquired_lock_event", + event_metadata={ + "gauge_labels": f"{cronjob_id}:{pod_id}", + "gauge_value": 1, + }, + ) + ) + + @staticmethod + def _emit_released_lock_event(cronjob_id: str, pod_id: str): + asyncio.create_task( + service_logger_obj.async_service_success_hook( + service=ServiceTypes.POD_LOCK_MANAGER, + duration=DEFAULT_CRON_JOB_LOCK_TTL_SECONDS, + call_type="_emit_released_lock_event", + event_metadata={ + "gauge_labels": f"{cronjob_id}:{pod_id}", + "gauge_value": 0, + }, + ) + )