diff --git a/docs/my-website/docs/proxy/logging.md b/docs/my-website/docs/proxy/logging.md index bfbc280db20..5aa78f73d72 100644 --- a/docs/my-website/docs/proxy/logging.md +++ b/docs/my-website/docs/proxy/logging.md @@ -3,9 +3,9 @@ import Tabs from '@theme/Tabs'; import TabItem from '@theme/TabItem'; -# Logging - Custom Callbacks, Langfuse, OpenTelemetry, Sentry +# 🔎 Logging - Custom Callbacks, Langfuse, s3 Bucket, Sentry, OpenTelemetry -Log Proxy Input, Output, Exceptions using Custom Callbacks, Langfuse, OpenTelemetry, LangFuse, DynamoDB +Log Proxy Input, Output, Exceptions using Custom Callbacks, Langfuse, OpenTelemetry, LangFuse, DynamoDB, s3 Bucket ## Custom Callback Class [Async] Use this when you want to run custom callbacks in `python` @@ -597,6 +597,62 @@ Here's the log view on Elastic Search. You can see the request `input`, `output` --> + +## Logging Proxy Input/Output - s3 Buckets + +We will use the `--config` to set +- `litellm.success_callback = ["s3"]` + +This will log all successfull LLM calls to s3 Bucket + +**Step 1** Set AWS Credentials in .env + +```shell +AWS_ACCESS_KEY_ID = "" +AWS_SECRET_ACCESS_KEY = "" +AWS_REGION_NAME = "" +``` + +**Step 2**: Create a `config.yaml` file and set `litellm_settings`: `success_callback` +```yaml +model_list: + - model_name: gpt-3.5-turbo + litellm_params: + model: gpt-3.5-turbo +litellm_settings: + success_callback: ["s3"] + s3_callback_params: + s3_bucket_name: logs-bucket-litellm # AWS Bucket Name for S3 + s3_region_name: us-west-2 # AWS Region Name for S3 + s3_aws_access_key_id: os.environ/AWS_ACCESS_KEY_ID # us os.environ/ to pass environment variables. This is AWS Access Key ID for S3 + s3_aws_secret_access_key: os.environ/AWS_SECRET_ACCESS_KEY # AWS Secret Access Key for S3 + s3_endpoint_url: https://s3.amazonaws.com # [OPTIONAL] S3 endpoint URL, if you want to use Backblaze/cloudflare s3 buckets +``` + +**Step 3**: Start the proxy, make a test request + +Start proxy +```shell +litellm --config config.yaml --debug +``` + +Test Request +```shell +curl --location 'http://0.0.0.0:8000/chat/completions' \ + --header 'Content-Type: application/json' \ + --data ' { + "model": "Azure OpenAI GPT-4 East", + "messages": [ + { + "role": "user", + "content": "what llm are you" + } + ] + }' +``` + +Your logs should be available on the specified s3 Bucket + ## Logging Proxy Input/Output - DynamoDB We will use the `--config` to set diff --git a/docs/my-website/sidebars.js b/docs/my-website/sidebars.js index 3c6e0011dbb..facda0e98df 100644 --- a/docs/my-website/sidebars.js +++ b/docs/my-website/sidebars.js @@ -111,12 +111,12 @@ const sidebars = { "proxy/users", "proxy/model_management", "proxy/reliability", + "proxy/caching", + "proxy/logging", "proxy/health", "proxy/call_hooks", "proxy/rules", - "proxy/caching", "proxy/alerting", - "proxy/logging", "proxy/streaming_logging", "proxy/deploy", "proxy/cli", diff --git a/litellm/__init__.py b/litellm/__init__.py index 670c19f4c13..837e5443420 100644 --- a/litellm/__init__.py +++ b/litellm/__init__.py @@ -135,6 +135,7 @@ model_fallbacks: Optional[List] = None # Deprecated for 'litellm.fallbacks' model_cost_map_url: str = "https://raw.githubusercontent.com/BerriAI/litellm/main/model_prices_and_context_window.json" suppress_debug_info = False dynamodb_table_name: Optional[str] = None +s3_callback_params: Optional[Dict] = None #### RELIABILITY #### request_timeout: Optional[float] = 6000 num_retries: Optional[int] = None # per model endpoint diff --git a/litellm/integrations/s3.py b/litellm/integrations/s3.py new file mode 100644 index 00000000000..e7f607b417e --- /dev/null +++ b/litellm/integrations/s3.py @@ -0,0 +1,145 @@ +#### 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 + + +class S3Logger: + # Class variables or attributes + def __init__( + self, + s3_bucket_name=None, + s3_region_name=None, + s3_api_version=None, + s3_use_ssl=True, + s3_verify=None, + s3_endpoint_url=None, + s3_aws_access_key_id=None, + s3_aws_secret_access_key=None, + s3_aws_session_token=None, + s3_config=None, + **kwargs, + ): + import boto3 + + try: + print_verbose("in init s3 logger") + + if litellm.s3_callback_params is not None: + # read in .env variables - example os.environ/AWS_BUCKET_NAME + for key, value in litellm.s3_callback_params.items(): + if type(value) is str and value.startswith("os.environ/"): + litellm.s3_callback_params[key] = litellm.get_secret(value) + # now set s3 params from litellm.s3_logger_params + s3_bucket_name = litellm.s3_callback_params.get("s3_bucket_name") + s3_region_name = litellm.s3_callback_params.get("s3_region_name") + s3_api_version = litellm.s3_callback_params.get("s3_api_version") + s3_use_ssl = litellm.s3_callback_params.get("s3_use_ssl") + s3_verify = litellm.s3_callback_params.get("s3_verify") + s3_endpoint_url = litellm.s3_callback_params.get("s3_endpoint_url") + s3_aws_access_key_id = litellm.s3_callback_params.get( + "s3_aws_access_key_id" + ) + s3_aws_secret_access_key = litellm.s3_callback_params.get( + "s3_aws_secret_access_key" + ) + s3_aws_session_token = litellm.s3_callback_params.get( + "s3_aws_session_token" + ) + s3_config = litellm.s3_callback_params.get("s3_config") + # done reading litellm.s3_callback_params + + self.bucket_name = s3_bucket_name + # Create an S3 client with custom endpoint URL + self.s3_client = boto3.client( + "s3", + region_name=s3_region_name, + endpoint_url=s3_endpoint_url, + api_version=s3_api_version, + use_ssl=s3_use_ssl, + verify=s3_verify, + aws_access_key_id=s3_aws_access_key_id, + aws_secret_access_key=s3_aws_secret_access_key, + aws_session_token=s3_aws_session_token, + config=s3_config, + **kwargs, + ) + 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 + ): + self.log_event(kwargs, response_obj, start_time, end_time, print_verbose) + + def log_event(self, kwargs, response_obj, start_time, end_time, print_verbose): + try: + print_verbose(f"s3 Logging - Enters logging function for model {kwargs}") + + # construct payload to send to s3 + # 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") + optional_params = kwargs.get("optional_params", {}) + call_type = kwargs.get("call_type", "litellm.completion") + usage = response_obj["usage"] + id = response_obj.get("id", str(uuid.uuid4())) + + # Build the initial payload + payload = { + "id": id, + "call_type": call_type, + "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, + } + + # 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 + s3_object_key = payload["id"] + + import json + + payload = json.dumps(payload) + + print_verbose(f"\ns3 Logger - Logging payload = {payload}") + + response = self.s3_client.put_object( + Bucket=self.bucket_name, + Key=s3_object_key, + Body=payload, + ContentType="application/json", + ContentLanguage="en", + ContentDisposition=f'inline; filename="{key}.json"', + ) + + print_verbose(f"Response from s3:{str(response)}") + + print_verbose(f"s3 Layer Logging - final response object: {response_obj}") + return response + except Exception as e: + traceback.print_exc() + print_verbose(f"s3 Layer Error - {str(e)}\n{traceback.format_exc()}") + pass diff --git a/litellm/tests/test_s3_logs.py b/litellm/tests/test_s3_logs.py new file mode 100644 index 00000000000..2a919d1272d --- /dev/null +++ b/litellm/tests/test_s3_logs.py @@ -0,0 +1,101 @@ +import sys +import os +import io, asyncio + +# import logging +# logging.basicConfig(level=logging.DEBUG) +sys.path.insert(0, os.path.abspath("../..")) + +from litellm import completion +import litellm + +litellm.num_retries = 3 + +import time, random +import pytest + + +def test_s3_logging(): + # all s3 requests need to be in one test function + # since we are modifying stdout, and pytests runs tests in parallel + # on circle ci - we only test litellm.acompletion() + try: + # pre + # redirect stdout to log_file + + litellm.success_callback = ["s3"] + litellm.s3_callback_params = { + "s3_bucket_name": "litellm-logs", + "s3_aws_secret_access_key": "os.environ/AWS_SECRET_ACCESS_KEY", + "s3_aws_access_key_id": "os.environ/AWS_ACCESS_KEY_ID", + } + litellm.set_verbose = True + + print("Testing async s3 logging") + + expected_keys = [] + + async def _test(): + return await litellm.acompletion( + model="gpt-3.5-turbo", + messages=[{"role": "user", "content": "This is a test"}], + max_tokens=10, + temperature=0.7, + user="ishaan-2", + ) + + response = asyncio.run(_test()) + print(f"response: {response}") + expected_keys.append(response.id) + + # # streaming + async + # async def _test2(): + # response = await litellm.acompletion( + # model="gpt-3.5-turbo", + # messages=[{"role": "user", "content": "what llm are u"}], + # max_tokens=10, + # temperature=0.7, + # user="ishaan-2", + # stream=True, + # ) + # async for chunk in response: + # pass + + # asyncio.run(_test2()) + + # aembedding() + # async def _test3(): + # return await litellm.aembedding( + # model="text-embedding-ada-002", input=["hi"], user="ishaan-2" + # ) + + # response = asyncio.run(_test3()) + # expected_keys.append(response.id) + # time.sleep(1) + + import boto3 + + s3 = boto3.client("s3") + bucket_name = "litellm-logs" + # List objects in the bucket + response = s3.list_objects(Bucket=bucket_name) + + # Sort the objects based on the LastModified timestamp + objects = sorted( + response["Contents"], key=lambda x: x["LastModified"], reverse=True + ) + # Get the keys of the most recent objects + most_recent_keys = [obj["Key"] for obj in objects] + print("\n most recent keys", most_recent_keys) + print("\n Expected keys: ", expected_keys) + for key in expected_keys: + assert key in most_recent_keys + 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 77e4de8e913..f3e743ec4cb 100644 --- a/litellm/utils.py +++ b/litellm/utils.py @@ -47,6 +47,7 @@ from .integrations.weights_biases import WeightsBiasesLogger from .integrations.custom_logger import CustomLogger from .integrations.langfuse import LangFuseLogger from .integrations.dynamodb import DyanmoDBLogger +from .integrations.s3 import S3Logger from .integrations.litedebugger import LiteDebugger from .proxy._types import KeyManagementSystem from openai import OpenAIError as OriginalError @@ -90,6 +91,7 @@ weightsBiasesLogger = None customLogger = None langFuseLogger = None dynamoLogger = None +s3Logger = None llmonitorLogger = None aispendLogger = None berrispendLogger = None @@ -1459,6 +1461,36 @@ class Logging: end_time=end_time, print_verbose=print_verbose, ) + if callback == "s3": + global s3Logger + if s3Logger is None: + s3Logger = S3Logger() + if self.stream: + if "complete_streaming_response" in self.model_call_details: + print_verbose( + "S3Logger Logger: Got Stream Event - Completed Stream Response" + ) + await s3Logger._async_log_event( + kwargs=self.model_call_details, + response_obj=self.model_call_details[ + "complete_streaming_response" + ], + start_time=start_time, + end_time=end_time, + print_verbose=print_verbose, + ) + else: + print_verbose( + "S3Logger Logger: Got Stream Event - No complete stream response as yet" + ) + else: + await s3Logger._async_log_event( + kwargs=self.model_call_details, + response_obj=result, + start_time=start_time, + end_time=end_time, + print_verbose=print_verbose, + ) if callback == "langfuse": global langFuseLogger print_verbose("reaches Async langfuse for logging!") @@ -1806,6 +1838,11 @@ def client(original_function): # we only support async dynamo db logging for acompletion/aembedding since that's used on proxy litellm._async_success_callback.append(callback) removed_async_items.append(index) + elif callback == "s3": + # s3 is an async callback, it's used for the proxy and needs to be async + # we only support async s3 logging for acompletion/aembedding since that's used on proxy + litellm._async_success_callback.append(callback) + removed_async_items.append(index) elif callback == "langfuse" and inspect.iscoroutinefunction( original_function ): @@ -4678,7 +4715,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, heliconeLogger, aispendLogger, berrispendLogger, supabaseClient, liteDebuggerClient, llmonitorLogger, promptLayerLogger, langFuseLogger, customLogger, weightsBiasesLogger, langsmithLogger, dynamoLogger + global sentry_sdk_instance, capture_exception, add_breadcrumb, posthog, slack_app, alerts_channel, traceloopLogger, heliconeLogger, aispendLogger, berrispendLogger, supabaseClient, liteDebuggerClient, llmonitorLogger, promptLayerLogger, langFuseLogger, customLogger, weightsBiasesLogger, langsmithLogger, dynamoLogger, s3Logger try: for callback in callback_list: print_verbose(f"callback: {callback}") @@ -4743,6 +4780,8 @@ def set_callbacks(callback_list, function_id=None): langFuseLogger = LangFuseLogger() elif callback == "dynamodb": dynamoLogger = DyanmoDBLogger() + elif callback == "s3": + s3Logger = S3Logger() elif callback == "wandb": weightsBiasesLogger = WeightsBiasesLogger() elif callback == "langsmith":