From a1b004900b6484fd539be5d653c897b30f815572 Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Mon, 18 Mar 2024 10:55:29 -0700 Subject: [PATCH 1/9] (feat) add datadog integration --- litellm/proxy/proxy_server.py | 35 +++++++++++++++++++++++++++++++++++ 1 file changed, 35 insertions(+) diff --git a/litellm/proxy/proxy_server.py b/litellm/proxy/proxy_server.py index 827b07eef1..cd5ce42363 100644 --- a/litellm/proxy/proxy_server.py +++ b/litellm/proxy/proxy_server.py @@ -9,6 +9,14 @@ import warnings import importlib import warnings +from datadog import statsd, initialize as datadog_initialize + +# Define DataDog client + +options = {"statsd_host": "127.0.0.1", "statsd_port": 8125} + +datadog_initialize(**options) # type: ignore + def showwarning(message, category, filename, lineno, file=None, line=None): traceback_info = f"{filename}:{lineno}: {category.__name__}: {message}\n" @@ -228,6 +236,31 @@ try: app.mount("/ui", StaticFiles(directory=ui_path, html=True), name="ui") except: pass + +from starlette.middleware.base import BaseHTTPMiddleware + + +class DatadogMetricsMiddleware(BaseHTTPMiddleware): + async def dispatch(self, request: Request, call_next): + start_time = time.time() + + # Request Count + statsd.increment("litellm.proxy.requests") + + try: + response = await call_next(request) + except Exception as e: + # Error Count + statsd.increment("litellm.proxy.errors") + raise e + + # Calculate response time + response_time = (time.time() - start_time) * 1000 # in milliseconds + statsd.distribution("litellm.proxy.response_time", response_time) + + return response + + app.add_middleware( CORSMiddleware, allow_origins=origins, @@ -236,6 +269,8 @@ app.add_middleware( allow_headers=["*"], ) +app.add_middleware(DatadogMetricsMiddleware) + from typing import Dict From ac826851fa7371289f999bc4cc37cee249a32d4e Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Mon, 18 Mar 2024 16:01:47 -0700 Subject: [PATCH 2/9] (feat) v0 datadog logger --- litellm/utils.py | 33 ++++++++++++++++++++++++++++++++- 1 file changed, 32 insertions(+), 1 deletion(-) diff --git a/litellm/utils.py b/litellm/utils.py index d36ba4e1a4..24dbef51a7 100644 --- a/litellm/utils.py +++ b/litellm/utils.py @@ -65,6 +65,7 @@ from .integrations.langsmith import LangsmithLogger from .integrations.weights_biases import WeightsBiasesLogger from .integrations.custom_logger import CustomLogger from .integrations.langfuse import LangFuseLogger +from .integrations.datadog import DataDogLogger from .integrations.dynamodb import DyanmoDBLogger from .integrations.s3 import S3Logger from .integrations.clickhouse import ClickhouseLogger @@ -121,6 +122,7 @@ langsmithLogger = None weightsBiasesLogger = None customLogger = None langFuseLogger = None +dataDogLogger = None dynamoLogger = None s3Logger = None genericAPILogger = None @@ -1473,6 +1475,33 @@ class Logging: user_id=kwargs.get("user", None), print_verbose=print_verbose, ) + if callback == "datadog": + global dataDogLogger + verbose_logger.debug("reaches datadog for success logging!") + kwargs = {} + for k, v in self.model_call_details.items(): + if ( + k != "original_response" + ): # copy.deepcopy raises errors as this could be a coroutine + kwargs[k] = v + # this only logs streaming once, complete_streaming_response exists i.e when stream ends + if self.stream: + verbose_logger.debug( + f"datadog: is complete_streaming_response in kwargs: {kwargs.get('complete_streaming_response', None)}" + ) + if complete_streaming_response is None: + continue + else: + print_verbose("reaches datadog for streaming logging!") + result = kwargs["complete_streaming_response"] + dataDogLogger.log_event( + kwargs=kwargs, + response_obj=result, + start_time=start_time, + end_time=end_time, + user_id=kwargs.get("user", None), + print_verbose=print_verbose, + ) if callback == "generic": global genericAPILogger verbose_logger.debug("reaches langfuse for success logging!") @@ -6082,7 +6111,7 @@ def validate_environment(model: Optional[str] = None) -> dict: def set_callbacks(callback_list, function_id=None): - global sentry_sdk_instance, capture_exception, add_breadcrumb, posthog, slack_app, alerts_channel, traceloopLogger, athinaLogger, heliconeLogger, aispendLogger, berrispendLogger, supabaseClient, liteDebuggerClient, llmonitorLogger, promptLayerLogger, langFuseLogger, customLogger, weightsBiasesLogger, langsmithLogger, dynamoLogger, s3Logger + global sentry_sdk_instance, capture_exception, add_breadcrumb, posthog, slack_app, alerts_channel, traceloopLogger, athinaLogger, heliconeLogger, aispendLogger, berrispendLogger, supabaseClient, liteDebuggerClient, llmonitorLogger, promptLayerLogger, langFuseLogger, customLogger, weightsBiasesLogger, langsmithLogger, dynamoLogger, s3Logger, dataDogLogger try: for callback in callback_list: print_verbose(f"callback: {callback}") @@ -6148,6 +6177,8 @@ def set_callbacks(callback_list, function_id=None): promptLayerLogger = PromptLayerLogger() elif callback == "langfuse": langFuseLogger = LangFuseLogger() + elif callback == "datadog": + dataDogLogger = DataDogLogger() elif callback == "dynamodb": dynamoLogger = DyanmoDBLogger() elif callback == "s3": From 6a7c86fc58105975a6788f6e1bec4bab16469c05 Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Mon, 18 Mar 2024 16:22:59 -0700 Subject: [PATCH 3/9] (test) dd test --- litellm/tests/test_datadog.py | 27 +++++++++++++++++++++++++++ 1 file changed, 27 insertions(+) create mode 100644 litellm/tests/test_datadog.py diff --git a/litellm/tests/test_datadog.py b/litellm/tests/test_datadog.py new file mode 100644 index 0000000000..789b39fd9d --- /dev/null +++ b/litellm/tests/test_datadog.py @@ -0,0 +1,27 @@ +import sys +import os +import io + +sys.path.insert(0, os.path.abspath("../..")) + +from litellm import completion +import litellm +import pytest + +import time + + +@pytest.mark.skip(reason="beta test - this is a new feature") +def test_datadog_logging(): + try: + litellm.success_callback = ["datadog"] + litellm.set_verbose = True + response = completion( + model="gpt-3.5-turbo", + messages=[{"role": "user", "content": "what llm are u"}], + max_tokens=10, + temperature=0.2, + ) + print(response) + except Exception as e: + print(e) From 2afddf45364304f4c1edd8a02edbc0078443e08e Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Mon, 18 Mar 2024 16:27:01 -0700 Subject: [PATCH 4/9] (feat) init datadog logger --- litellm/integrations/datadog.py | 141 ++++++++++++++++++++++++++++++++ 1 file changed, 141 insertions(+) create mode 100644 litellm/integrations/datadog.py diff --git a/litellm/integrations/datadog.py b/litellm/integrations/datadog.py new file mode 100644 index 0000000000..d5adac1fca --- /dev/null +++ b/litellm/integrations/datadog.py @@ -0,0 +1,141 @@ +#### What this does #### +# On success + failure, log events to Supabase + +import dotenv, os +import requests + +dotenv.load_dotenv() # Loading env variables using dotenv +import traceback +import datetime, subprocess, sys +import litellm, uuid +from litellm._logging import print_verbose, verbose_logger +from datadog_api_client.v2 import ApiClient, Configuration +from datadog import statsd, api as datadog_api, initialize as datadog_initialize + +# Define DataDog client + +from datadog_api_client import ApiClient, Configuration +from datadog_api_client.v2.api.logs_api import LogsApi +from datadog_api_client.v2.model.log import Log +from datadog_api_client.v2.model import * +from datadog_api_client.v2.models import * + + +class DataDogLogger: + # Class variables or attributes + def __init__( + self, + **kwargs, + ): + + try: + verbose_logger.debug(f"in init datadog logger") + pass + + except Exception as e: + print_verbose(f"Got exception on init s3 client {str(e)}") + raise e + + async def _async_log_event( + self, kwargs, response_obj, start_time, end_time, print_verbose, user_id + ): + self.log_event(kwargs, response_obj, start_time, end_time, print_verbose) + + def log_event( + self, kwargs, response_obj, start_time, end_time, user_id, print_verbose + ): + try: + verbose_logger.debug( + f"datadog Logging - Enters logging function for model {kwargs}" + ) + litellm_params = kwargs.get("litellm_params", {}) + metadata = ( + litellm_params.get("metadata", {}) or {} + ) # if litellm_params['metadata'] == None + messages = kwargs.get("messages") + optional_params = kwargs.get("optional_params", {}) + call_type = kwargs.get("call_type", "litellm.completion") + cache_hit = kwargs.get("cache_hit", False) + usage = response_obj["usage"] + id = response_obj.get("id", str(uuid.uuid4())) + usage = dict(usage) + try: + response_time = (end_time - start_time).total_seconds() + except: + response_time = None + + try: + response_obj = dict(response_obj) + except: + response_obj = response_obj + + # Clean Metadata before logging - never log raw metadata + # the raw metadata can contain circular references which leads to infinite recursion + # we clean out all extra litellm metadata params before logging + clean_metadata = {} + if isinstance(metadata, dict): + for key, value in metadata.items(): + # clean litellm metadata before logging + if key in [ + "endpoint", + "caching_groups", + "previous_models", + ]: + continue + else: + clean_metadata[key] = value + + # Build the initial payload + payload = { + "id": id, + "call_type": call_type, + "cache_hit": cache_hit, + "startTime": start_time, + "endTime": end_time, + "responseTime (seconds)": response_time, + "model": kwargs.get("model", ""), + "user": kwargs.get("user", ""), + "modelParameters": optional_params, + "spend": kwargs.get("response_cost", 0), + "messages": messages, + "response": response_obj, + "usage": usage, + "metadata": clean_metadata, + } + + # Ensure everything in the payload is converted to str + for key, value in payload.items(): + try: + payload[key] = str(value) + except: + # non blocking if it can't cast to a str + pass + import json + + payload = json.dumps(payload) + + print_verbose(f"\ndd Logger - Logging payload = {payload}") + + configuration = Configuration() + with ApiClient(configuration) as api_client: + api_instance = LogsApi(api_client) + body = HTTPLog( + [ + HTTPLogItem( + ddsource="litellm", + message=payload, + service="litellm-server", + ), + ] + ) + response = api_instance.submit_log(body) + + print_verbose( + f"Datadog Layer Logging - final response object: {response_obj}" + ) + except Exception as e: + traceback.print_exc() + verbose_logger.debug( + f"Datadog Layer Error - {str(e)}\n{traceback.format_exc()}" + ) + pass From 2bce2ee7e3a2ba650757fcd0618fdb607465b4a9 Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Mon, 18 Mar 2024 16:27:55 -0700 Subject: [PATCH 5/9] (fix) log to datadog --- requirements.txt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/requirements.txt b/requirements.txt index eaff0fb711..5d6b2307d1 100644 --- a/requirements.txt +++ b/requirements.txt @@ -18,7 +18,7 @@ google-generativeai==0.3.2 # for vertex ai calls async_generator==1.10.0 # for async ollama calls traceloop-sdk==0.5.3 # for open telemetry logging langfuse>=2.6.3 # for langfuse self-hosted logging -clickhouse_connect==0.7.0 +datadog-api-client==2.23.0 # for datadog logging orjson==3.9.15 # fast /embedding responses apscheduler==3.10.4 # for resetting budget in background fastapi-sso==0.10.0 # admin UI, SSO From 325eb35b448b70421a3dca65f7811289b70d4d2d Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Mon, 18 Mar 2024 16:31:16 -0700 Subject: [PATCH 6/9] (cleanup) proxy --- litellm/proxy/proxy_server.py | 35 ----------------------------------- 1 file changed, 35 deletions(-) diff --git a/litellm/proxy/proxy_server.py b/litellm/proxy/proxy_server.py index cd5ce42363..ac59d5ac7a 100644 --- a/litellm/proxy/proxy_server.py +++ b/litellm/proxy/proxy_server.py @@ -9,14 +9,6 @@ import warnings import importlib import warnings -from datadog import statsd, initialize as datadog_initialize - -# Define DataDog client - -options = {"statsd_host": "127.0.0.1", "statsd_port": 8125} - -datadog_initialize(**options) # type: ignore - def showwarning(message, category, filename, lineno, file=None, line=None): traceback_info = f"{filename}:{lineno}: {category.__name__}: {message}\n" @@ -237,30 +229,6 @@ try: except: pass -from starlette.middleware.base import BaseHTTPMiddleware - - -class DatadogMetricsMiddleware(BaseHTTPMiddleware): - async def dispatch(self, request: Request, call_next): - start_time = time.time() - - # Request Count - statsd.increment("litellm.proxy.requests") - - try: - response = await call_next(request) - except Exception as e: - # Error Count - statsd.increment("litellm.proxy.errors") - raise e - - # Calculate response time - response_time = (time.time() - start_time) * 1000 # in milliseconds - statsd.distribution("litellm.proxy.response_time", response_time) - - return response - - app.add_middleware( CORSMiddleware, allow_origins=origins, @@ -269,9 +237,6 @@ app.add_middleware( allow_headers=["*"], ) -app.add_middleware(DatadogMetricsMiddleware) - - from typing import Dict api_key_header = APIKeyHeader( From 17295286d35fe9787fa59642f1bbf2fde03af721 Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Mon, 18 Mar 2024 16:40:01 -0700 Subject: [PATCH 7/9] (cleanup) extra line --- litellm/proxy/proxy_server.py | 2 -- 1 file changed, 2 deletions(-) diff --git a/litellm/proxy/proxy_server.py b/litellm/proxy/proxy_server.py index ac59d5ac7a..c320d0d1bb 100644 --- a/litellm/proxy/proxy_server.py +++ b/litellm/proxy/proxy_server.py @@ -228,7 +228,6 @@ try: app.mount("/ui", StaticFiles(directory=ui_path, html=True), name="ui") except: pass - app.add_middleware( CORSMiddleware, allow_origins=origins, @@ -236,7 +235,6 @@ app.add_middleware( allow_methods=["*"], allow_headers=["*"], ) - from typing import Dict api_key_header = APIKeyHeader( From fedb4e5703cbc1a4b1740283c6ad082d2f8c0ceb Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Mon, 18 Mar 2024 16:40:58 -0700 Subject: [PATCH 8/9] (fix) clean dd imports --- litellm/integrations/datadog.py | 26 ++++++++++++++------------ 1 file changed, 14 insertions(+), 12 deletions(-) diff --git a/litellm/integrations/datadog.py b/litellm/integrations/datadog.py index d5adac1fca..f5db5bf1f7 100644 --- a/litellm/integrations/datadog.py +++ b/litellm/integrations/datadog.py @@ -9,16 +9,6 @@ import traceback import datetime, subprocess, sys import litellm, uuid from litellm._logging import print_verbose, verbose_logger -from datadog_api_client.v2 import ApiClient, Configuration -from datadog import statsd, api as datadog_api, initialize as datadog_initialize - -# Define DataDog client - -from datadog_api_client import ApiClient, Configuration -from datadog_api_client.v2.api.logs_api import LogsApi -from datadog_api_client.v2.model.log import Log -from datadog_api_client.v2.model import * -from datadog_api_client.v2.models import * class DataDogLogger: @@ -27,6 +17,14 @@ class DataDogLogger: self, **kwargs, ): + from datadog_api_client import ApiClient, Configuration + + # check if the correct env variables are set + if os.getenv("DD_API_KEY", None) is None: + raise Exception("DD_API_KEY is not set, set 'DD_API_KEY=<>") + if os.getenv("DD_SITE", None) is None: + raise Exception("DD_SITE is not set in .env, set 'DD_SITE=<>") + self.configuration = Configuration() try: verbose_logger.debug(f"in init datadog logger") @@ -45,6 +43,11 @@ class DataDogLogger: self, kwargs, response_obj, start_time, end_time, user_id, print_verbose ): try: + # Define DataDog client + from datadog_api_client.v2.api.logs_api import LogsApi + from datadog_api_client.v2 import ApiClient + from datadog_api_client.v2.models import HTTPLogItem, HTTPLog + verbose_logger.debug( f"datadog Logging - Enters logging function for model {kwargs}" ) @@ -116,8 +119,7 @@ class DataDogLogger: print_verbose(f"\ndd Logger - Logging payload = {payload}") - configuration = Configuration() - with ApiClient(configuration) as api_client: + with ApiClient(self.configuration) as api_client: api_instance = LogsApi(api_client) body = HTTPLog( [ From 0dab5841cdab0d3a51ae3229b9a752577597f97a Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Mon, 18 Mar 2024 16:43:04 -0700 Subject: [PATCH 9/9] (fix) proxy extra spacing --- litellm/proxy/proxy_server.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/litellm/proxy/proxy_server.py b/litellm/proxy/proxy_server.py index c320d0d1bb..827b07eef1 100644 --- a/litellm/proxy/proxy_server.py +++ b/litellm/proxy/proxy_server.py @@ -235,6 +235,8 @@ app.add_middleware( allow_methods=["*"], allow_headers=["*"], ) + + from typing import Dict api_key_header = APIKeyHeader(