diff --git a/litellm/integrations/langfuse.py b/litellm/integrations/langfuse.py index 51ad4d19abe..c2612feb805 100644 --- a/litellm/integrations/langfuse.py +++ b/litellm/integrations/langfuse.py @@ -12,7 +12,9 @@ import litellm class LangFuseLogger: # Class variables or attributes - def __init__(self, langfuse_public_key=None, langfuse_secret=None): + def __init__( + self, langfuse_public_key=None, langfuse_secret=None, flush_interval=1 + ): try: from langfuse import Langfuse except Exception as e: @@ -31,7 +33,7 @@ class LangFuseLogger: host=self.langfuse_host, release=self.langfuse_release, debug=self.langfuse_debug, - flush_interval=1, # flush interval in seconds + flush_interval=flush_interval, # flush interval in seconds ) # set the current langfuse project id in the environ diff --git a/litellm/integrations/slack_alerting.py b/litellm/integrations/slack_alerting.py index adc56698c39..e3e5811ed5f 100644 --- a/litellm/integrations/slack_alerting.py +++ b/litellm/integrations/slack_alerting.py @@ -12,6 +12,7 @@ from litellm.caching import DualCache import asyncio import aiohttp from litellm.llms.custom_httpx.http_handler import AsyncHTTPHandler +import datetime class SlackAlerting: @@ -47,6 +48,18 @@ class SlackAlerting: self.internal_usage_cache = DualCache() self.async_http_handler = AsyncHTTPHandler() self.alert_to_webhook_url = alert_to_webhook_url + self.langfuse_logger = None + + try: + from litellm.integrations.langfuse import LangFuseLogger + + self.langfuse_logger = LangFuseLogger( + os.getenv("LANGFUSE_PUBLIC_KEY"), + os.getenv("LANGFUSE_SECRET_KEY"), + flush_interval=1, + ) + except: + pass pass @@ -93,39 +106,68 @@ class SlackAlerting: request_info: str, request_data: Optional[dict] = None, kwargs: Optional[dict] = None, + type: Literal["hanging_request", "slow_response"] = "hanging_request", + start_time: Optional[datetime.datetime] = None, + end_time: Optional[datetime.datetime] = None, ): import uuid # For now: do nothing as we're debugging why this is not working as expected + if request_data is not None: + trace_id = request_data.get("metadata", {}).get( + "trace_id", None + ) # get langfuse trace id + if trace_id is None: + trace_id = "litellm-alert-trace-" + str(uuid.uuid4()) + request_data["metadata"]["trace_id"] = trace_id + elif kwargs is not None: + _litellm_params = kwargs.get("litellm_params", {}) + trace_id = _litellm_params.get("metadata", {}).get( + "trace_id", None + ) # get langfuse trace id + if trace_id is None: + trace_id = "litellm-alert-trace-" + str(uuid.uuid4()) + _litellm_params["metadata"]["trace_id"] = trace_id + + # Log hanging request as an error on langfuse + if type == "hanging_request": + if self.langfuse_logger is not None: + _logging_kwargs = copy.deepcopy(request_data) + if _logging_kwargs is None: + _logging_kwargs = {} + _logging_kwargs["litellm_params"] = {} + request_data = request_data or {} + _logging_kwargs["litellm_params"]["metadata"] = request_data.get( + "metadata", {} + ) + # log to langfuse in a separate thread + import threading + + threading.Thread( + target=self.langfuse_logger.log_event, + args=( + _logging_kwargs, + None, + start_time, + end_time, + None, + print, + "ERROR", + "Requests is hanging", + ), + ).start() + + _langfuse_host = os.environ.get("LANGFUSE_HOST", "https://cloud.langfuse.com") + _langfuse_project_id = os.environ.get("LANGFUSE_PROJECT_ID") + + # langfuse urls look like: https://us.cloud.langfuse.com/project/************/traces/litellm-alert-trace-ididi9dk-09292-************ + + _langfuse_url = ( + f"{_langfuse_host}/project/{_langfuse_project_id}/traces/{trace_id}" + ) + request_info += f"\n🪢 Langfuse Trace: {_langfuse_url}" return request_info - # if request_data is not None: - # trace_id = request_data.get("metadata", {}).get( - # "trace_id", None - # ) # get langfuse trace id - # if trace_id is None: - # trace_id = "litellm-alert-trace-" + str(uuid.uuid4()) - # request_data["metadata"]["trace_id"] = trace_id - # elif kwargs is not None: - # _litellm_params = kwargs.get("litellm_params", {}) - # trace_id = _litellm_params.get("metadata", {}).get( - # "trace_id", None - # ) # get langfuse trace id - # if trace_id is None: - # trace_id = "litellm-alert-trace-" + str(uuid.uuid4()) - # _litellm_params["metadata"]["trace_id"] = trace_id - - # _langfuse_host = os.environ.get("LANGFUSE_HOST", "https://cloud.langfuse.com") - # _langfuse_project_id = os.environ.get("LANGFUSE_PROJECT_ID") - - # # langfuse urls look like: https://us.cloud.langfuse.com/project/************/traces/litellm-alert-trace-ididi9dk-09292-************ - - # _langfuse_url = ( - # f"{_langfuse_host}/project/{_langfuse_project_id}/traces/{trace_id}" - # ) - # request_info += f"\n🪢 Langfuse Trace: {_langfuse_url}" - # return request_info - def _response_taking_too_long_callback( self, kwargs, # kwargs to completion @@ -194,7 +236,7 @@ class SlackAlerting: if time_difference_float > self.alerting_threshold: if "langfuse" in litellm.success_callback: request_info = self._add_langfuse_trace_id_to_alert( - request_info=request_info, kwargs=kwargs + request_info=request_info, kwargs=kwargs, type="slow_response" ) # add deployment latencies to alert if ( @@ -222,8 +264,8 @@ class SlackAlerting: async def response_taking_too_long( self, - start_time: Optional[float] = None, - end_time: Optional[float] = None, + start_time: Optional[datetime.datetime] = None, + end_time: Optional[datetime.datetime] = None, type: Literal["hanging_request", "slow_response"] = "hanging_request", request_data: Optional[dict] = None, ): @@ -243,10 +285,6 @@ class SlackAlerting: except: messages = "" request_info = f"\nRequest Model: `{model}`\nMessages: `{messages}`" - if "langfuse" in litellm.success_callback: - request_info = self._add_langfuse_trace_id_to_alert( - request_info=request_info, request_data=request_data - ) else: request_info = "" @@ -288,6 +326,15 @@ class SlackAlerting: f"`Requests are hanging - {self.alerting_threshold}s+ request time`" ) + if "langfuse" in litellm.success_callback: + request_info = self._add_langfuse_trace_id_to_alert( + request_info=request_info, + request_data=request_data, + type="hanging_request", + start_time=start_time, + end_time=end_time, + ) + # add deployment latencies to alert _deployment_latency_map = self._get_deployment_latencies_to_alert( metadata=request_data.get("metadata", {})