diff --git a/docs/my-website/docs/proxy/logging.md b/docs/my-website/docs/proxy/logging.md index 0c5fd6bb91d..3f5596fc9e6 100644 --- a/docs/my-website/docs/proxy/logging.md +++ b/docs/my-website/docs/proxy/logging.md @@ -8,6 +8,7 @@ import TabItem from '@theme/TabItem'; Log Proxy Input, Output, Exceptions using Custom Callbacks, Langfuse, OpenTelemetry, LangFuse, DynamoDB, s3 Bucket - [Async Custom Callbacks](#custom-callback-class-async) +- [Async Custom Callback APIs](#custom-callback-apis-async) - [Logging to Langfuse](#logging-proxy-inputoutput---langfuse) - [Logging to s3 Buckets](#logging-proxy-inputoutput---s3-buckets) - [Logging to DynamoDB](#logging-proxy-inputoutput---dynamodb) @@ -297,6 +298,106 @@ ModelResponse( ``` +## Custom Callback APIs [Async] + +:::info + +This is an Enterprise only feature [Get Started with Enterprise here](https://github.com/BerriAI/litellm/tree/main/enterprise) + +::: + +Use this if you: +- Want to use custom callbacks written in a non Python programming language +- Want your callbacks to run on a different microservice + +#### Step 1. Create your generic logging API endpoint +Set up a generic API endpoint that can receive data in JSON format. The data will be included within a "data" field. + +Your server should support the following Request format: + +```shell +curl --location https://your-domain.com/log-event \ + --request POST \ + --header "Content-Type: application/json" \ + --data '{ + "data": { + "id": "chatcmpl-8sgE89cEQ4q9biRtxMvDfQU1O82PT", + "call_type": "acompletion", + "cache_hit": "None", + "startTime": "2024-02-15 16:18:44.336280", + "endTime": "2024-02-15 16:18:45.045539", + "model": "gpt-3.5-turbo", + "user": "ishaan-2", + "modelParameters": "{'temperature': 0.7, 'max_tokens': 10, 'user': 'ishaan-2', 'extra_body': {}}", + "messages": "[{'role': 'user', 'content': 'This is a test'}]", + "response": "ModelResponse(id='chatcmpl-8sgE89cEQ4q9biRtxMvDfQU1O82PT', choices=[Choices(finish_reason='length', index=0, message=Message(content='Great! How can I assist you with this test', role='assistant'))], created=1708042724, model='gpt-3.5-turbo-0613', object='chat.completion', system_fingerprint=None, usage=Usage(completion_tokens=10, prompt_tokens=11, total_tokens=21))", + "usage": "Usage(completion_tokens=10, prompt_tokens=11, total_tokens=21)", + "metadata": "{}", + "cost": "3.65e-05" + } + }' +``` + +Reference FastAPI Python Server + +Here's a reference FastAPI Server that is compatible with LiteLLM Proxy: + +```python +# this is an example endpoint to receive data from litellm +from fastapi import FastAPI, HTTPException, Request + +app = FastAPI() + + +@app.post("/log-event") +async def log_event(request: Request): + try: + print("Received /log-event request") + # Assuming the incoming request has JSON data + data = await request.json() + print("Received request data:") + print(data) + + # Your additional logic can go here + # For now, just printing the received data + + return {"message": "Request received successfully"} + except Exception as e: + print(f"Error processing request: {str(e)}") + import traceback + + traceback.print_exc() + raise HTTPException(status_code=500, detail="Internal Server Error") + + +if __name__ == "__main__": + import uvicorn + uvicorn.run(app, host="127.0.0.1", port=8000) + + +``` + + +#### Step 2. Set your `GENERIC_LOGGER_ENDPOINT` to the endpoint + route we should send callback logs to + +```shell +os.environ["GENERIC_LOGGER_ENDPOINT"] = "http://localhost:8000/log-event" +``` + +#### Step 3. Create a `config.yaml` file and set `litellm_settings`: `success_callback` = ["generic"] + +Example litellm proxy config.yaml +```yaml +model_list: + - model_name: gpt-3.5-turbo + litellm_params: + model: gpt-3.5-turbo +litellm_settings: + success_callback: ["generic"] +``` + +Start the LiteLLM Proxy and make a test request to verify the logs reached your callback API + ## Logging Proxy Input/Output - Langfuse We will use the `--config` to set `litellm.success_callback = ["langfuse"]` this will log all successfull LLM calls to langfuse diff --git a/enterprise/callbacks/example_logging_api.py b/enterprise/callbacks/example_logging_api.py new file mode 100644 index 00000000000..57ea99a6748 --- /dev/null +++ b/enterprise/callbacks/example_logging_api.py @@ -0,0 +1,31 @@ +# this is an example endpoint to receive data from litellm +from fastapi import FastAPI, HTTPException, Request + +app = FastAPI() + + +@app.post("/log-event") +async def log_event(request: Request): + try: + print("Received /log-event request") + # Assuming the incoming request has JSON data + data = await request.json() + print("Received request data:") + print(data) + + # Your additional logic can go here + # For now, just printing the received data + + return {"message": "Request received successfully"} + except Exception as e: + print(f"Error processing request: {str(e)}") + import traceback + + traceback.print_exc() + raise HTTPException(status_code=500, detail="Internal Server Error") + + +if __name__ == "__main__": + import uvicorn + + uvicorn.run(app, host="127.0.0.1", port=8000) diff --git a/enterprise/callbacks/generic_api_callback.py b/enterprise/callbacks/generic_api_callback.py new file mode 100644 index 00000000000..076c13d5eef --- /dev/null +++ b/enterprise/callbacks/generic_api_callback.py @@ -0,0 +1,128 @@ +# callback to make a request to an API endpoint + +#### What this does #### +# On success, logs events to Promptlayer +import dotenv, os +import requests + +from litellm.proxy._types import UserAPIKeyAuth +from litellm.caching import DualCache + +from typing import Literal, Union + +dotenv.load_dotenv() # Loading env variables using dotenv +import traceback + + +#### 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 + + +class GenericAPILogger: + # Class variables or attributes + def __init__(self, endpoint=None, headers=None): + try: + if endpoint == None: + # check env for "GENERIC_LOGGER_ENDPOINT" + if os.getenv("GENERIC_LOGGER_ENDPOINT"): + # Do something with the endpoint + endpoint = os.getenv("GENERIC_LOGGER_ENDPOINT") + else: + # Handle the case when the endpoint is not found in the environment variables + raise ValueError( + f"endpoint not set for GenericAPILogger, GENERIC_LOGGER_ENDPOINT not found in environment variables" + ) + headers = headers or litellm.generic_logger_headers + self.endpoint = endpoint + self.headers = headers + + verbose_logger.debug( + f"in init GenericAPILogger, endpoint {self.endpoint}, headers {self.headers}" + ) + + pass + + except Exception as e: + print_verbose(f"Got exception on init GenericAPILogger client {str(e)}") + raise e + + # This is sync, because we run this in a separate thread. Running in a sepearate thread ensures it will never block an LLM API call + # Experience with s3, Langfuse shows that async logging events are complicated and can block LLM calls + def log_event( + self, kwargs, response_obj, start_time, end_time, user_id, print_verbose + ): + try: + verbose_logger.debug( + f"GenericAPILogger Logging - Enters logging function for model {kwargs}" + ) + + # construct payload to send custom logger + # follows the same params as langfuse.py + litellm_params = kwargs.get("litellm_params", {}) + metadata = ( + litellm_params.get("metadata", {}) or {} + ) # if litellm_params['metadata'] == None + messages = kwargs.get("messages") + cost = kwargs.get("response_cost", 0.0) + 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())) + + # Build the initial payload + payload = { + "id": id, + "call_type": call_type, + "cache_hit": cache_hit, + "startTime": start_time, + "endTime": end_time, + "model": kwargs.get("model", ""), + "user": kwargs.get("user", ""), + "modelParameters": optional_params, + "messages": messages, + "response": response_obj, + "usage": usage, + "metadata": metadata, + "cost": cost, + } + + # 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 + + data = { + "data": payload, + } + data = json.dumps(data) + print_verbose(f"\nGeneric Logger - Logging payload = {data}") + + # make request to endpoint with payload + response = requests.post(self.endpoint, json=data, headers=self.headers) + + response_status = response.status_code + response_text = response.text + + print_verbose( + f"Generic Logger - final response status = {response_status}, response text = {response_text}" + ) + return response + except Exception as e: + traceback.print_exc() + verbose_logger.debug(f"Generic - {str(e)}\n{traceback.format_exc()}") + pass diff --git a/litellm/__init__.py b/litellm/__init__.py index 00d7488c94e..a7f232f7693 100644 --- a/litellm/__init__.py +++ b/litellm/__init__.py @@ -146,6 +146,7 @@ model_cost_map_url: str = "https://raw.githubusercontent.com/BerriAI/litellm/mai suppress_debug_info = False dynamodb_table_name: Optional[str] = None s3_callback_params: Optional[Dict] = None +generic_logger_headers: Optional[Dict] = None default_key_generate_params: Optional[Dict] = None upperbound_key_generate_params: Optional[Dict] = None default_team_settings: Optional[List] = None diff --git a/litellm/tests/test_custom_api_logger.py b/litellm/tests/test_custom_api_logger.py new file mode 100644 index 00000000000..bddce9a0878 --- /dev/null +++ b/litellm/tests/test_custom_api_logger.py @@ -0,0 +1,46 @@ +import sys +import os +import io, asyncio + +# import logging +# logging.basicConfig(level=logging.DEBUG) +sys.path.insert(0, os.path.abspath("../..")) +print("Modified sys.path:", sys.path) + + +from litellm import completion +import litellm + +litellm.num_retries = 3 + +import time, random +import pytest + + +@pytest.mark.asyncio +@pytest.mark.skip(reason="new beta feature, will be testing in our ci/cd soon") +async def test_custom_api_logging(): + try: + litellm.success_callback = ["generic"] + litellm.set_verbose = True + os.environ["GENERIC_LOGGER_ENDPOINT"] = "http://localhost:8000/log-event" + + print("Testing generic api logging") + + await litellm.acompletion( + model="gpt-3.5-turbo", + messages=[{"role": "user", "content": f"This is a test"}], + max_tokens=10, + temperature=0.7, + user="ishaan-2", + ) + + except Exception as e: + pytest.fail(f"An exception occurred - {e}") + finally: + # post, close log file and verify + # Reset stdout to the original value + print("Passed! Testing async s3 logging") + + +# test_s3_logging() diff --git a/litellm/utils.py b/litellm/utils.py index 35bf4640f8a..d3efcea738f 100644 --- a/litellm/utils.py +++ b/litellm/utils.py @@ -12,6 +12,7 @@ import litellm import dotenv, json, traceback, threading, base64, ast import subprocess, os +from os.path import abspath, join, dirname import litellm, openai import itertools import random, uuid, requests @@ -34,6 +35,7 @@ from dataclasses import ( from importlib import resources # filename = pkg_resources.resource_filename(__name__, "llms/tokenizers") + try: filename = str( resources.files().joinpath("llms/tokenizers") # type: ignore @@ -45,6 +47,7 @@ except: os.environ[ "TIKTOKEN_CACHE_DIR" ] = filename # use local copy of tiktoken b/c of - https://github.com/BerriAI/litellm/issues/1071 + encoding = tiktoken.get_encoding("cl100k_base") import importlib.metadata from ._logging import verbose_logger @@ -82,6 +85,20 @@ from .exceptions import ( BudgetExceededError, UnprocessableEntityError, ) + +# Import Enterprise features +project_path = abspath(join(dirname(__file__), "..", "..")) +# Add the "enterprise" directory to sys.path +verbose_logger.debug(f"current project_path: {project_path}") +enterprise_path = abspath(join(project_path, "enterprise")) +sys.path.append(enterprise_path) + +verbose_logger.debug(f"sys.path: {sys.path}") +try: + from enterprise.callbacks.generic_api_callback import GenericAPILogger +except Exception as e: + verbose_logger.debug(f"Exception import enterprise features {str(e)}") + from typing import cast, List, Dict, Union, Optional, Literal, Any from .caching import Cache from concurrent.futures import ThreadPoolExecutor @@ -107,6 +124,7 @@ customLogger = None langFuseLogger = None dynamoLogger = None s3Logger = None +genericAPILogger = None llmonitorLogger = None aispendLogger = None berrispendLogger = None @@ -1369,6 +1387,35 @@ class Logging: user_id=kwargs.get("user", None), print_verbose=print_verbose, ) + if callback == "generic": + global genericAPILogger + verbose_logger.debug("reaches langfuse 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"is complete_streaming_response in kwargs: {kwargs.get('complete_streaming_response', None)}" + ) + if complete_streaming_response is None: + break + else: + print_verbose("reaches langfuse for streaming logging!") + result = kwargs["complete_streaming_response"] + if genericAPILogger is None: + genericAPILogger = GenericAPILogger() + genericAPILogger.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 == "cache" and litellm.cache is not None: # this only logs streaming once, complete_streaming_response exists i.e when stream ends print_verbose("success_callback: reaches cache for logging!")