From bb04a340a56051bd02bad4a8b46f0f281d3e74cc Mon Sep 17 00:00:00 2001 From: Krrish Dholakia Date: Wed, 10 Jan 2024 20:52:01 +0530 Subject: [PATCH 1/5] fix(lowest_latency.py): add back tpm/rpm checks, configurable time window --- litellm/router.py | 25 ++- litellm/router_strategy/lowest_latency.py | 205 ++++++++++++++++--- litellm/tests/test_lowest_latency_routing.py | 133 +++++++++++- 3 files changed, 312 insertions(+), 51 deletions(-) diff --git a/litellm/router.py b/litellm/router.py index 311afeb446d..f6355550980 100644 --- a/litellm/router.py +++ b/litellm/router.py @@ -105,6 +105,7 @@ class Router: "usage-based-routing", "latency-based-routing", ] = "simple-shuffle", + routing_strategy_args: dict = {}, # just for latency-based routing ) -> None: self.set_verbose = set_verbose self.deployment_names: List = ( @@ -217,7 +218,9 @@ class Router: litellm.callbacks.append(self.lowesttpm_logger) # type: ignore elif routing_strategy == "latency-based-routing": self.lowestlatency_logger = LowestLatencyLoggingHandler( - router_cache=self.cache, model_list=self.model_list + router_cache=self.cache, + model_list=self.model_list, + routing_args=routing_strategy_args, ) if isinstance(litellm.callbacks, list): litellm.callbacks.append(self.lowestlatency_logger) # type: ignore @@ -1427,9 +1430,8 @@ class Router: http_client=httpx.AsyncClient( transport=AsyncCustomHTTPTransport(), limits=httpx.Limits( - max_connections=1000, - max_keepalive_connections=100 - ) + max_connections=1000, max_keepalive_connections=100 + ), ), # type: ignore ) self.cache.set_cache( @@ -1449,9 +1451,8 @@ class Router: http_client=httpx.Client( transport=CustomHTTPTransport(), limits=httpx.Limits( - max_connections=1000, - max_keepalive_connections=100 - ) + max_connections=1000, max_keepalive_connections=100 + ), ), # type: ignore ) self.cache.set_cache( @@ -1471,10 +1472,9 @@ class Router: max_retries=max_retries, http_client=httpx.AsyncClient( limits=httpx.Limits( - max_connections=1000, - max_keepalive_connections=100 + max_connections=1000, max_keepalive_connections=100 ) - ) + ), ) self.cache.set_cache( key=cache_key, @@ -1492,10 +1492,9 @@ class Router: max_retries=max_retries, http_client=httpx.Client( limits=httpx.Limits( - max_connections=1000, - max_keepalive_connections=100 + max_connections=1000, max_keepalive_connections=100 ) - ) + ), ) self.cache.set_cache( key=cache_key, diff --git a/litellm/router_strategy/lowest_latency.py b/litellm/router_strategy/lowest_latency.py index 43e28a8b3aa..3f8cb513b47 100644 --- a/litellm/router_strategy/lowest_latency.py +++ b/litellm/router_strategy/lowest_latency.py @@ -1,25 +1,46 @@ #### What this does #### # picks based on response time (for streaming, this is time to first token) - +from pydantic import BaseModel, Extra, Field, root_validator import dotenv, os, requests, random -from typing import Optional +from typing import Optional, Union, List, Dict from datetime import datetime, timedelta dotenv.load_dotenv() # Loading env variables using dotenv import traceback from litellm.caching import DualCache from litellm.integrations.custom_logger import CustomLogger +from litellm import ModelResponse +from litellm import token_counter + + +class LiteLLMBase(BaseModel): + """ + Implements default functions, all pydantic objects should have. + """ + + def json(self, **kwargs): + try: + return self.model_dump() # noqa + except: + # if using pydantic v1 + return self.dict() + + +class RoutingArgs(LiteLLMBase): + ttl: int = 1 * 60 * 60 # 1 hour class LowestLatencyLoggingHandler(CustomLogger): test_flag: bool = False logged_success: int = 0 logged_failure: int = 0 - default_cache_time_seconds: int = 1 * 60 * 60 # 1 hour - def __init__(self, router_cache: DualCache, model_list: list): + def __init__( + self, router_cache: DualCache, model_list: list, routing_args: dict = {} + ): self.router_cache = router_cache self.model_list = model_list + self.routing_args = RoutingArgs(**routing_args) def log_success_event(self, kwargs, response_obj, start_time, end_time): try: @@ -37,25 +58,64 @@ class LowestLatencyLoggingHandler(CustomLogger): if model_group is None or id is None: return - response_ms = end_time - start_time - # ------------ # Setup values # ------------ - latency_key = f"{model_group}_latency_map" + """ + { + {model_group}_map: { + id: { + "latency": [..] + f"{date:hour:minute}" : {"tpm": 34, "rpm": 3} + } + } + } + """ + latency_key = f"{model_group}_map" + + current_date = datetime.now().strftime("%Y-%m-%d") + current_hour = datetime.now().strftime("%H") + current_minute = datetime.now().strftime("%M") + precise_minute = f"{current_date}-{current_hour}-{current_minute}" + + response_ms: timedelta = end_time - start_time + + final_value = response_ms + total_tokens = 0 + + if isinstance(response_obj, ModelResponse): + completion_tokens = response_obj.usage.completion_tokens + total_tokens = response_obj.usage.total_tokens + final_value = float(completion_tokens / response_ms.total_seconds()) # ------------ # Update usage # ------------ - ## Latency request_count_dict = self.router_cache.get_cache(key=latency_key) or {} - if id in request_count_dict and isinstance(request_count_dict[id], list): - request_count_dict[id] = request_count_dict[id].append(response_ms) - else: - request_count_dict[id] = [response_ms] - self.router_cache.set_cache(key=latency_key, value=request_count_dict, ttl=self.default_cache_time_seconds) # reset map within window + if id not in request_count_dict: + request_count_dict[id] = {} + + ## Latency + request_count_dict[id].setdefault("latency", []).append(final_value) + + if precise_minute not in request_count_dict[id]: + request_count_dict[id][precise_minute] = {} + + ## TPM + request_count_dict[id][precise_minute]["tpm"] = ( + request_count_dict[id][precise_minute].get("tpm", 0) + total_tokens + ) + + ## RPM + request_count_dict[id][precise_minute]["rpm"] = ( + request_count_dict[id][precise_minute].get("rpm", 0) + 1 + ) + + self.router_cache.set_cache( + key=latency_key, value=request_count_dict, ttl=self.routing_args.ttl + ) # reset map within window ### TESTING ### if self.test_flag: @@ -80,25 +140,64 @@ class LowestLatencyLoggingHandler(CustomLogger): if model_group is None or id is None: return - response_ms = end_time - start_time - # ------------ # Setup values # ------------ - latency_key = f"{model_group}_latency_map" + """ + { + {model_group}_map: { + id: { + "latency": [..] + f"{date:hour:minute}" : {"tpm": 34, "rpm": 3} + } + } + } + """ + latency_key = f"{model_group}_map" + + current_date = datetime.now().strftime("%Y-%m-%d") + current_hour = datetime.now().strftime("%H") + current_minute = datetime.now().strftime("%M") + precise_minute = f"{current_date}-{current_hour}-{current_minute}" + + response_ms: timedelta = end_time - start_time + + final_value = response_ms + total_tokens = 0 + + if isinstance(response_obj, ModelResponse): + completion_tokens = response_obj.usage.completion_tokens + total_tokens = response_obj.usage.total_tokens + final_value = float(completion_tokens / response_ms.total_seconds()) # ------------ # Update usage # ------------ - ## Latency request_count_dict = self.router_cache.get_cache(key=latency_key) or {} - if id in request_count_dict and isinstance(request_count_dict[id], list): - request_count_dict[id] = request_count_dict[id] + [response_ms] - else: - request_count_dict[id] = [response_ms] - self.router_cache.set_cache(key=latency_key, value=request_count_dict, ttl=self.default_cache_time_seconds) # reset map within window + if id not in request_count_dict: + request_count_dict[id] = {} + + ## Latency + request_count_dict[id].setdefault("latency", []).append(final_value) + + if precise_minute not in request_count_dict[id]: + request_count_dict[id][precise_minute] = {} + + ## TPM + request_count_dict[id][precise_minute]["tpm"] = ( + request_count_dict[id][precise_minute].get("tpm", 0) + total_tokens + ) + + ## RPM + request_count_dict[id][precise_minute]["rpm"] = ( + request_count_dict[id][precise_minute].get("rpm", 0) + 1 + ) + + self.router_cache.set_cache( + key=latency_key, value=request_count_dict, ttl=self.routing_args.ttl + ) # reset map within window ### TESTING ### if self.test_flag: @@ -107,12 +206,18 @@ class LowestLatencyLoggingHandler(CustomLogger): traceback.print_exc() pass - def get_available_deployments(self, model_group: str, healthy_deployments: list): + def get_available_deployments( + self, + model_group: str, + healthy_deployments: list, + messages: Optional[List[Dict[str, str]]] = None, + input: Optional[Union[str, List]] = None, + ): """ Returns a deployment with the lowest latency """ # get list of potential deployments - latency_key = f"{model_group}_latency_map" + latency_key = f"{model_group}_map" request_count_dict = self.router_cache.get_cache(key=latency_key) or {} @@ -120,6 +225,12 @@ class LowestLatencyLoggingHandler(CustomLogger): # Find lowest used model # ---------------------- lowest_latency = float("inf") + + current_date = datetime.now().strftime("%Y-%m-%d") + current_hour = datetime.now().strftime("%H") + current_minute = datetime.now().strftime("%M") + precise_minute = f"{current_date}-{current_hour}-{current_minute}" + deployment = None if request_count_dict is None: # base case @@ -129,9 +240,17 @@ class LowestLatencyLoggingHandler(CustomLogger): for d in healthy_deployments: ## if healthy deployment not yet used if d["model_info"]["id"] not in all_deployments: - all_deployments[d["model_info"]["id"]] = [0] + all_deployments[d["model_info"]["id"]] = { + "latency": [0], + precise_minute: {"tpm": 0, "rpm": 0}, + } - for item, item_latency in all_deployments.items(): + try: + input_tokens = token_counter(messages=messages, text=input) + except: + input_tokens = 0 + + for item, item_map in all_deployments.items(): ## get the item from model list _deployment = None for m in healthy_deployments: @@ -140,18 +259,38 @@ class LowestLatencyLoggingHandler(CustomLogger): if _deployment is None: continue # skip to next one - - # get average latency - total = 0.0 + + _deployment_tpm = ( + _deployment.get("tpm", None) + or _deployment.get("litellm_params", {}).get("tpm", None) + or _deployment.get("model_info", {}).get("tpm", None) + or float("inf") + ) + + _deployment_rpm = ( + _deployment.get("rpm", None) + or _deployment.get("litellm_params", {}).get("rpm", None) + or _deployment.get("model_info", {}).get("rpm", None) + or float("inf") + ) + item_latency = item_map.get("latency", []) + item_rpm = item_map.get(precise_minute, {}).get("rpm", 0) + item_tpm = item_map.get(precise_minute, {}).get("tpm", 0) + + # get average latency + total: float = 0.0 for _call_latency in item_latency: - if isinstance(_call_latency, timedelta): - total += float(_call_latency.total_seconds()) - elif isinstance(_call_latency, float): + if isinstance(_call_latency, float): total += _call_latency - item_latency = total/len(item_latency) + item_latency = total / len(item_latency) if item_latency == 0: deployment = _deployment break + elif ( + item_tpm + input_tokens > _deployment_tpm + or item_rpm + 1 > _deployment_rpm + ): # if user passed in tpm / rpm in the model_list + continue elif item_latency < lowest_latency: lowest_latency = item_latency deployment = _deployment diff --git a/litellm/tests/test_lowest_latency_routing.py b/litellm/tests/test_lowest_latency_routing.py index 6805bfa5882..c9b1e7972c0 100644 --- a/litellm/tests/test_lowest_latency_routing.py +++ b/litellm/tests/test_lowest_latency_routing.py @@ -48,13 +48,55 @@ def test_latency_updated(): start_time=start_time, end_time=end_time, ) - latency_key = f"{model_group}_latency_map" - assert end_time - start_time == test_cache.get_cache(key=latency_key)[deployment_id][0] + latency_key = f"{model_group}_map" + assert ( + end_time - start_time + == test_cache.get_cache(key=latency_key)[deployment_id]["latency"][0] + ) # test_tpm_rpm_updated() +def test_latency_updated_custom_ttl(): + """ + Invalidate the cached request. + + Test that the cache is empty + """ + test_cache = DualCache() + model_list = [] + cache_time = 3 + lowest_latency_logger = LowestLatencyLoggingHandler( + router_cache=test_cache, model_list=model_list, routing_args={"ttl": cache_time} + ) + model_group = "gpt-3.5-turbo" + deployment_id = "1234" + kwargs = { + "litellm_params": { + "metadata": { + "model_group": "gpt-3.5-turbo", + "deployment": "azure/chatgpt-v-2", + }, + "model_info": {"id": deployment_id}, + } + } + start_time = time.time() + response_obj = {"usage": {"total_tokens": 50}} + time.sleep(5) + end_time = time.time() + lowest_latency_logger.log_success_event( + response_obj=response_obj, + kwargs=kwargs, + start_time=start_time, + end_time=end_time, + ) + latency_key = f"{model_group}_map" + assert isinstance(test_cache.get_cache(key=latency_key), dict) + time.sleep(cache_time) + assert test_cache.get_cache(key=latency_key) is None + + def test_get_available_deployments(): test_cache = DualCache() model_list = [ @@ -133,6 +175,90 @@ def test_get_available_deployments(): # test_get_available_deployments() +def test_get_available_endpoints_tpm_rpm_check(): + """ + Pass in list of 2 valid models + + Update cache with 1 model clearly being at tpm/rpm limit + + assert that only the valid model is returned + """ + test_cache = DualCache() + model_list = [ + { + "model_name": "gpt-3.5-turbo", + "litellm_params": {"model": "azure/chatgpt-v-2"}, + "model_info": {"id": "1234", "rpm": 10}, + }, + { + "model_name": "gpt-3.5-turbo", + "litellm_params": {"model": "azure/chatgpt-v-2"}, + "model_info": {"id": "5678", "rpm": 3}, + }, + ] + lowest_latency_logger = LowestLatencyLoggingHandler( + router_cache=test_cache, model_list=model_list + ) + model_group = "gpt-3.5-turbo" + ## DEPLOYMENT 1 ## + deployment_id = "1234" + kwargs = { + "litellm_params": { + "metadata": { + "model_group": "gpt-3.5-turbo", + "deployment": "azure/chatgpt-v-2", + }, + "model_info": {"id": deployment_id}, + } + } + for _ in range(3): + start_time = time.time() + response_obj = {"usage": {"total_tokens": 50}} + time.sleep(0.05) + end_time = time.time() + lowest_latency_logger.log_success_event( + response_obj=response_obj, + kwargs=kwargs, + start_time=start_time, + end_time=end_time, + ) + ## DEPLOYMENT 2 ## + deployment_id = "5678" + kwargs = { + "litellm_params": { + "metadata": { + "model_group": "gpt-3.5-turbo", + "deployment": "azure/chatgpt-v-2", + }, + "model_info": {"id": deployment_id}, + } + } + for _ in range(3): + start_time = time.time() + response_obj = {"usage": {"total_tokens": 20}} + time.sleep(2) + end_time = time.time() + lowest_latency_logger.log_success_event( + response_obj=response_obj, + kwargs=kwargs, + start_time=start_time, + end_time=end_time, + ) + + ## CHECK WHAT'S SELECTED ## + print( + lowest_latency_logger.get_available_deployments( + model_group=model_group, healthy_deployments=model_list + ) + ) + assert ( + lowest_latency_logger.get_available_deployments( + model_group=model_group, healthy_deployments=model_list + )["model_info"]["id"] + == "1234" + ) + + def test_router_get_available_deployments(): """ Test if routers 'get_available_deployments' returns the fastest deployment @@ -213,9 +339,6 @@ def test_router_get_available_deployments(): assert router.get_available_deployment(model="azure-model")["model_info"]["id"] == 2 -# test_get_available_deployments() - - # test_router_get_available_deployments() From 162f6f1ed3c6b4ce13370ac92d0ad7cd93d6cd79 Mon Sep 17 00:00:00 2001 From: Krrish Dholakia Date: Wed, 10 Jan 2024 20:58:29 +0530 Subject: [PATCH 2/5] refactor: refactor key tests --- litellm/tests/test_deployed_proxy_keygen.py | 98 ++--- litellm/tests/test_proxy_server_keys.py | 379 ++++++++++---------- 2 files changed, 248 insertions(+), 229 deletions(-) diff --git a/litellm/tests/test_deployed_proxy_keygen.py b/litellm/tests/test_deployed_proxy_keygen.py index e627609436f..e0acee083c0 100644 --- a/litellm/tests/test_deployed_proxy_keygen.py +++ b/litellm/tests/test_deployed_proxy_keygen.py @@ -1,63 +1,63 @@ -import sys, os, time -import traceback -from dotenv import load_dotenv +# import sys, os, time +# import traceback +# from dotenv import load_dotenv -load_dotenv() -import os, io +# load_dotenv() +# import os, io -# this file is to test litellm/proxy +# # this file is to test litellm/proxy -sys.path.insert( - 0, os.path.abspath("../..") -) # Adds the parent directory to the system path -import pytest, logging, requests -import litellm -from litellm import embedding, completion, completion_cost, Timeout -from litellm import RateLimitError +# sys.path.insert( +# 0, os.path.abspath("../..") +# ) # Adds the parent directory to the system path +# import pytest, logging, requests +# import litellm +# from litellm import embedding, completion, completion_cost, Timeout +# from litellm import RateLimitError -def test_add_new_key(): - max_retries = 3 - retry_delay = 1 # seconds +# def test_add_new_key(): +# max_retries = 3 +# retry_delay = 1 # seconds - for retry in range(max_retries + 1): - try: - # Your test data - test_data = { - "models": ["gpt-3.5-turbo", "gpt-4", "claude-2", "azure-model"], - "aliases": {"mistral-7b": "gpt-3.5-turbo"}, - "duration": "20m", - } - print("testing proxy server") +# for retry in range(max_retries + 1): +# try: +# # Your test data +# test_data = { +# "models": ["gpt-3.5-turbo", "gpt-4", "claude-2", "azure-model"], +# "aliases": {"mistral-7b": "gpt-3.5-turbo"}, +# "duration": "20m", +# } +# print("testing proxy server") - # Your bearer token - token = os.getenv("PROXY_MASTER_KEY") - headers = {"Authorization": f"Bearer {token}"} +# # Your bearer token +# token = os.getenv("PROXY_MASTER_KEY") +# headers = {"Authorization": f"Bearer {token}"} - staging_endpoint = "https://litellm-litellm-pr-1376.up.railway.app" - main_endpoint = "https://litellm-staging.up.railway.app" +# staging_endpoint = "https://litellm-litellm-pr-1376.up.railway.app" +# main_endpoint = "https://litellm-staging.up.railway.app" - # Make a request to the staging endpoint - response = requests.post( - main_endpoint + "/key/generate", json=test_data, headers=headers - ) +# # Make a request to the staging endpoint +# response = requests.post( +# main_endpoint + "/key/generate", json=test_data, headers=headers +# ) - print(f"response: {response.text}") +# print(f"response: {response.text}") - if response.status_code == 200: - result = response.json() - break # Successful response, exit the loop - elif response.status_code == 503 and retry < max_retries: - print( - f"Retrying in {retry_delay} seconds... (Retry {retry + 1}/{max_retries})" - ) - time.sleep(retry_delay) - else: - assert False, f"Unexpected response status code: {response.status_code}" +# if response.status_code == 200: +# result = response.json() +# break # Successful response, exit the loop +# elif response.status_code == 503 and retry < max_retries: +# print( +# f"Retrying in {retry_delay} seconds... (Retry {retry + 1}/{max_retries})" +# ) +# time.sleep(retry_delay) +# else: +# assert False, f"Unexpected response status code: {response.status_code}" - except Exception as e: - print(traceback.format_exc()) - pytest.fail(f"An error occurred {e}") +# except Exception as e: +# print(traceback.format_exc()) +# pytest.fail(f"An error occurred {e}") -test_add_new_key() +# test_add_new_key() diff --git a/litellm/tests/test_proxy_server_keys.py b/litellm/tests/test_proxy_server_keys.py index 4aa2c2e26e5..763c546021e 100644 --- a/litellm/tests/test_proxy_server_keys.py +++ b/litellm/tests/test_proxy_server_keys.py @@ -15,111 +15,105 @@ import litellm from litellm import embedding, completion, completion_cost, Timeout from litellm import RateLimitError -# Configure logging -logging.basicConfig( - level=logging.DEBUG, # Set the desired logging level - format="%(asctime)s - %(levelname)s - %(message)s", -) + +import sys, os, time +import traceback +from dotenv import load_dotenv + +load_dotenv() +import os, io + +# this file is to test litellm/proxy from concurrent.futures import ThreadPoolExecutor -# test /chat/completion request to the proxy -from fastapi.testclient import TestClient -from fastapi import FastAPI -from litellm.proxy.proxy_server import ( - router, - save_worker_config, - startup_event, - shutdown_event, -) # Replace with the actual module where your FastAPI router is defined +sys.path.insert( + 0, os.path.abspath("../..") +) # Adds the parent directory to the system path -filepath = os.path.dirname(os.path.abspath(__file__)) -config_fp = f"{filepath}/test_configs/test_config.yaml" -save_worker_config( - config=config_fp, - model=None, - alias=None, - api_base=None, - api_version=None, - debug=True, - temperature=None, - max_tokens=None, - request_timeout=600, - max_budget=None, - telemetry=False, - drop_params=True, - add_function_to_prompt=False, - headers=None, - save=False, - use_queue=False, -) -app = FastAPI() -app.include_router(router) # Include your router in the test app +import pytest, logging, requests +import litellm +from litellm import embedding, completion, completion_cost, Timeout +from litellm import RateLimitError +from github import Github +import subprocess -@app.on_event("startup") -async def wrapper_startup_event(): - await startup_event() +# Function to execute a command and return the output +def run_command(command): + process = subprocess.Popen(command, stdout=subprocess.PIPE, shell=True) + output, _ = process.communicate() + return output.decode().strip() -@app.on_event("shutdown") -async def wrapper_shutdown_event(): - await shutdown_event() +# Retrieve the current branch name +branch_name = run_command("git rev-parse --abbrev-ref HEAD") + +# GitHub personal access token (with repo scope) or use username and password +access_token = os.getenv("GITHUB_ACCESS_TOKEN") +# Instantiate the PyGithub library's Github object +g = Github(access_token) + +# Provide the owner and name of the repository where the pull request is located +repository_owner = "BerriAI" +repository_name = "litellm" + +# Get the repository object +repo = g.get_repo(f"{repository_owner}/{repository_name}") + +# Iterate through the pull requests to find the one related to your branch +for pr in repo.get_pulls(): + print(f"in here! {pr.head.ref}") + if pr.head.ref == branch_name: + pr_number = pr.number + break + +print(f"The pull request number for branch {branch_name} is: {pr_number}") -# Here you create a fixture that will be used by your tests -# Make sure the fixture returns TestClient(app) -@pytest.fixture(autouse=True) -def client(): - from litellm.proxy.proxy_server import cleanup_router_config_variables +def test_add_new_key(): + max_retries = 3 + retry_delay = 1 # seconds - cleanup_router_config_variables() - with TestClient(app) as client: - yield client - - -def test_add_new_key(client): - try: - # Your test data - test_data = { - "models": ["gpt-3.5-turbo", "gpt-4", "claude-2", "azure-model"], - "aliases": {"mistral-7b": "gpt-3.5-turbo"}, - "duration": "20m", - } - print("testing proxy server") - # Your bearer token - token = os.getenv("PROXY_MASTER_KEY") - - headers = {"Authorization": f"Bearer {token}"} - response = client.post("/key/generate", json=test_data, headers=headers) - print(f"response: {response.text}") - assert response.status_code == 200 - result = response.json() - assert result["key"].startswith("sk-") - - def _post_data(): - json_data = { - "model": "azure-model", - "messages": [ - { - "role": "user", - "content": f"this is a test request, write a short poem {time.time()}", - } - ], + for retry in range(max_retries + 1): + try: + # Your test data + test_data = { + "models": ["gpt-3.5-turbo", "gpt-4", "claude-2", "azure-model"], + "aliases": {"mistral-7b": "gpt-3.5-turbo"}, + "duration": "20m", } - response = client.post( - "/chat/completions", - json=json_data, - headers={"Authorization": f"Bearer {result['key']}"}, + print("testing proxy server") + + # Your bearer token + token = os.getenv("PROXY_MASTER_KEY") + headers = {"Authorization": f"Bearer {token}"} + + endpoint = f"https://litellm-litellm-pr-{pr_number}.up.railway.app" + + # Make a request to the staging endpoint + response = requests.post( + endpoint + "/key/generate", json=test_data, headers=headers ) - return response - _post_data() - print(f"Received response: {result}") - except Exception as e: - pytest.fail(f"LiteLLM Proxy test failed. Exception: {str(e)}") + print(f"response: {response.text}") + + if response.status_code == 200: + result = response.json() + break # Successful response, exit the loop + elif response.status_code == 503 and retry < max_retries: + print( + f"Retrying in {retry_delay} seconds... (Retry {retry + 1}/{max_retries})" + ) + time.sleep(retry_delay) + else: + assert False, f"Unexpected response status code: {response.status_code}" + + except Exception as e: + print(traceback.format_exc()) + pytest.fail(f"An error occurred {e}") -def test_update_new_key(client): +def test_update_new_key(): try: # Your test data test_data = { @@ -130,17 +124,23 @@ def test_update_new_key(client): print("testing proxy server") # Your bearer token token = os.getenv("PROXY_MASTER_KEY") - headers = {"Authorization": f"Bearer {token}"} - response = client.post("/key/generate", json=test_data, headers=headers) - print(f"response: {response.text}") + + endpoint = f"https://litellm-litellm-pr-{pr_number}.up.railway.app" + + # Make a request to the staging endpoint + response = requests.post( + endpoint + "/key/generate", json=test_data, headers=headers + ) assert response.status_code == 200 result = response.json() assert result["key"].startswith("sk-") def _post_data(): json_data = {"models": ["bedrock-models"], "key": result["key"]} - response = client.post("/key/update", json=json_data, headers=headers) + response = requests.post( + endpoint + "/key/generate", json=json_data, headers=headers + ) print(f"response text: {response.text}") assert response.status_code == 200 return response @@ -151,101 +151,120 @@ def test_update_new_key(client): pytest.fail(f"LiteLLM Proxy test failed. Exception: {str(e)}") -# # Run the test - only runs via pytest +# def test_add_new_key_max_parallel_limit(): +# try: +# # Your test data +# test_data = {"duration": "20m", "max_parallel_requests": 1} +# # Your bearer token +# token = os.getenv("PROXY_MASTER_KEY") +# headers = {"Authorization": f"Bearer {token}"} + +# endpoint = f"https://litellm-litellm-pr-{pr_number}.up.railway.app" +# print(f"endpoint: {endpoint}") +# # Make a request to the staging endpoint +# response = requests.post( +# endpoint + "/key/generate", json=test_data, headers=headers +# ) +# assert response.status_code == 200 +# result = response.json() + +# # load endpoint with model +# model_data = { +# "model_name": "azure-model", +# "litellm_params": { +# "model": "azure/chatgpt-v-2", +# "api_key": os.getenv("AZURE_API_KEY"), +# "api_base": os.getenv("AZURE_API_BASE"), +# "api_version": os.getenv("AZURE_API_VERSION") +# } +# } +# response = requests.post(endpoint + "/model/new", json=model_data, headers=headers) +# assert response.status_code == 200 +# print(f"response text: {response.text}") -def test_add_new_key_max_parallel_limit(client): - try: - # Your test data - test_data = {"duration": "20m", "max_parallel_requests": 1} - # Your bearer token - token = os.getenv("PROXY_MASTER_KEY") +# def _post_data(): +# json_data = { +# "model": "azure-model", +# "messages": [ +# { +# "role": "user", +# "content": f"this is a test request, write a short poem {time.time()}", +# } +# ], +# } +# # Your bearer token +# response = requests.post( +# endpoint + "/chat/completions", json=json_data, headers={"Authorization": f"Bearer {result['key']}"} +# ) +# return response - headers = {"Authorization": f"Bearer {token}"} - response = client.post("/key/generate", json=test_data, headers=headers) - print(f"response: {response.text}") - assert response.status_code == 200 - result = response.json() +# def _run_in_parallel(): +# with ThreadPoolExecutor(max_workers=2) as executor: +# future1 = executor.submit(_post_data) +# future2 = executor.submit(_post_data) - def _post_data(): - json_data = { - "model": "azure-model", - "messages": [ - { - "role": "user", - "content": f"this is a test request, write a short poem {time.time()}", - } - ], - } - response = client.post( - "/chat/completions", - json=json_data, - headers={"Authorization": f"Bearer {result['key']}"}, - ) - return response +# # Obtain the results from the futures +# response1 = future1.result() +# print(f"response1 text: {response1.text}") +# response2 = future2.result() +# print(f"response2 text: {response2.text}") +# if response1.status_code == 429 or response2.status_code == 429: +# pass +# else: +# raise Exception() - def _run_in_parallel(): - with ThreadPoolExecutor(max_workers=2) as executor: - future1 = executor.submit(_post_data) - future2 = executor.submit(_post_data) +# _run_in_parallel() +# except Exception as e: +# pytest.fail(f"LiteLLM Proxy test failed. Exception: {str(e)}") - # Obtain the results from the futures - response1 = future1.result() - response2 = future2.result() - if response1.status_code == 429 or response2.status_code == 429: - pass - else: - raise Exception() +# def test_add_new_key_max_parallel_limit_streaming(): +# try: +# # Your test data +# test_data = {"duration": "20m", "max_parallel_requests": 1} +# # Your bearer token +# token = os.getenv("PROXY_MASTER_KEY") +# headers = {"Authorization": f"Bearer {token}"} - _run_in_parallel() - except Exception as e: - pytest.fail(f"LiteLLM Proxy test failed. Exception: {str(e)}") +# endpoint = f"https://litellm-litellm-pr-{pr_number}.up.railway.app" +# # Make a request to the staging endpoint +# response = requests.post( +# endpoint + "/key/generate", json=test_data, headers=headers +# ) +# print(f"response: {response.text}") +# assert response.status_code == 200 +# result = response.json() -def test_add_new_key_max_parallel_limit_streaming(client): - try: - # Your test data - test_data = {"duration": "20m", "max_parallel_requests": 1} - # Your bearer token - token = os.getenv("PROXY_MASTER_KEY") +# def _post_data(): +# json_data = { +# "model": "azure-model", +# "messages": [ +# { +# "role": "user", +# "content": f"this is a test request, write a short poem {time.time()}", +# } +# ], +# "stream": True, +# } +# response = requests.post( +# endpoint + "/chat/completions", json=json_data, headers={"Authorization": f"Bearer {result['key']}"} +# ) +# return response - headers = {"Authorization": f"Bearer {token}"} - response = client.post("/key/generate", json=test_data, headers=headers) - print(f"response: {response.text}") - assert response.status_code == 200 - result = response.json() +# def _run_in_parallel(): +# with ThreadPoolExecutor(max_workers=2) as executor: +# future1 = executor.submit(_post_data) +# future2 = executor.submit(_post_data) - def _post_data(): - json_data = { - "model": "azure-model", - "messages": [ - { - "role": "user", - "content": f"this is a test request, write a short poem {time.time()}", - } - ], - "stream": True, - } - response = client.post( - "/chat/completions", - json=json_data, - headers={"Authorization": f"Bearer {result['key']}"}, - ) - return response +# # Obtain the results from the futures +# response1 = future1.result() +# response2 = future2.result() +# if response1.status_code == 429 or response2.status_code == 429: +# pass +# else: +# raise Exception() - def _run_in_parallel(): - with ThreadPoolExecutor(max_workers=2) as executor: - future1 = executor.submit(_post_data) - future2 = executor.submit(_post_data) - - # Obtain the results from the futures - response1 = future1.result() - response2 = future2.result() - if response1.status_code == 429 or response2.status_code == 429: - pass - else: - raise Exception() - - _run_in_parallel() - except Exception as e: - pytest.fail(f"LiteLLM Proxy test failed. Exception: {str(e)}") +# _run_in_parallel() +# except Exception as e: +# pytest.fail(f"LiteLLM Proxy test failed. Exception: {str(e)}") From 14a65eb730327296da14db9945afdc725df5bb55 Mon Sep 17 00:00:00 2001 From: Krrish Dholakia Date: Wed, 10 Jan 2024 21:05:12 +0530 Subject: [PATCH 3/5] test(test_proxy_server_keys.py): removing as this is now tested via the docker build job --- litellm/tests/test_proxy_server_keys.py | 132 ++++++++++++------------ 1 file changed, 66 insertions(+), 66 deletions(-) diff --git a/litellm/tests/test_proxy_server_keys.py b/litellm/tests/test_proxy_server_keys.py index 763c546021e..398fb68742a 100644 --- a/litellm/tests/test_proxy_server_keys.py +++ b/litellm/tests/test_proxy_server_keys.py @@ -70,85 +70,85 @@ for pr in repo.get_pulls(): print(f"The pull request number for branch {branch_name} is: {pr_number}") -def test_add_new_key(): - max_retries = 3 - retry_delay = 1 # seconds +# def test_add_new_key(): +# max_retries = 3 +# retry_delay = 10 # seconds - for retry in range(max_retries + 1): - try: - # Your test data - test_data = { - "models": ["gpt-3.5-turbo", "gpt-4", "claude-2", "azure-model"], - "aliases": {"mistral-7b": "gpt-3.5-turbo"}, - "duration": "20m", - } - print("testing proxy server") +# for retry in range(max_retries + 1): +# try: +# # Your test data +# test_data = { +# "models": ["gpt-3.5-turbo", "gpt-4", "claude-2", "azure-model"], +# "aliases": {"mistral-7b": "gpt-3.5-turbo"}, +# "duration": "20m", +# } +# print("testing proxy server") - # Your bearer token - token = os.getenv("PROXY_MASTER_KEY") - headers = {"Authorization": f"Bearer {token}"} +# # Your bearer token +# token = os.getenv("PROXY_MASTER_KEY") +# headers = {"Authorization": f"Bearer {token}"} - endpoint = f"https://litellm-litellm-pr-{pr_number}.up.railway.app" +# endpoint = f"https://litellm-litellm-pr-{pr_number}.up.railway.app" - # Make a request to the staging endpoint - response = requests.post( - endpoint + "/key/generate", json=test_data, headers=headers - ) +# # Make a request to the staging endpoint +# response = requests.post( +# endpoint + "/key/generate", json=test_data, headers=headers +# ) - print(f"response: {response.text}") +# print(f"response: {response.text}") - if response.status_code == 200: - result = response.json() - break # Successful response, exit the loop - elif response.status_code == 503 and retry < max_retries: - print( - f"Retrying in {retry_delay} seconds... (Retry {retry + 1}/{max_retries})" - ) - time.sleep(retry_delay) - else: - assert False, f"Unexpected response status code: {response.status_code}" +# if response.status_code == 200: +# result = response.json() +# break # Successful response, exit the loop +# elif response.status_code == 503 and retry < max_retries: +# print( +# f"Retrying in {retry_delay} seconds... (Retry {retry + 1}/{max_retries})" +# ) +# time.sleep(retry_delay) +# else: +# assert False, f"Unexpected response status code: {response.status_code}" - except Exception as e: - print(traceback.format_exc()) - pytest.fail(f"An error occurred {e}") +# except Exception as e: +# print(traceback.format_exc()) +# pytest.fail(f"An error occurred {e}") -def test_update_new_key(): - try: - # Your test data - test_data = { - "models": ["gpt-3.5-turbo", "gpt-4", "claude-2", "azure-model"], - "aliases": {"mistral-7b": "gpt-3.5-turbo"}, - "duration": "20m", - } - print("testing proxy server") - # Your bearer token - token = os.getenv("PROXY_MASTER_KEY") - headers = {"Authorization": f"Bearer {token}"} +# def test_update_new_key(): +# try: +# # Your test data +# test_data = { +# "models": ["gpt-3.5-turbo", "gpt-4", "claude-2", "azure-model"], +# "aliases": {"mistral-7b": "gpt-3.5-turbo"}, +# "duration": "20m", +# } +# print("testing proxy server") +# # Your bearer token +# token = os.getenv("PROXY_MASTER_KEY") +# headers = {"Authorization": f"Bearer {token}"} - endpoint = f"https://litellm-litellm-pr-{pr_number}.up.railway.app" +# endpoint = f"https://litellm-litellm-pr-{pr_number}.up.railway.app" - # Make a request to the staging endpoint - response = requests.post( - endpoint + "/key/generate", json=test_data, headers=headers - ) - assert response.status_code == 200 - result = response.json() - assert result["key"].startswith("sk-") +# # Make a request to the staging endpoint +# response = requests.post( +# endpoint + "/key/generate", json=test_data, headers=headers +# ) +# assert response.status_code == 200 +# result = response.json() +# assert result["key"].startswith("sk-") - def _post_data(): - json_data = {"models": ["bedrock-models"], "key": result["key"]} - response = requests.post( - endpoint + "/key/generate", json=json_data, headers=headers - ) - print(f"response text: {response.text}") - assert response.status_code == 200 - return response +# def _post_data(): +# json_data = {"models": ["bedrock-models"], "key": result["key"]} +# response = requests.post( +# endpoint + "/key/generate", json=json_data, headers=headers +# ) +# print(f"response text: {response.text}") +# assert response.status_code == 200 +# return response - _post_data() - print(f"Received response: {result}") - except Exception as e: - pytest.fail(f"LiteLLM Proxy test failed. Exception: {str(e)}") +# _post_data() +# print(f"Received response: {result}") +# except Exception as e: +# pytest.fail(f"LiteLLM Proxy test failed. Exception: {str(e)}") # def test_add_new_key_max_parallel_limit(): From 7f269e92c578d36326b33930c9319040d30a7d45 Mon Sep 17 00:00:00 2001 From: Krrish Dholakia Date: Wed, 10 Jan 2024 21:17:30 +0530 Subject: [PATCH 4/5] test(test_completion_with_retries.py): remove duplicate test --- .circleci/config.yml | 7 ------- litellm/tests/test_completion_with_retries.py | 18 ------------------ 2 files changed, 25 deletions(-) diff --git a/.circleci/config.yml b/.circleci/config.yml index 3ea6b7fca9a..2f2b581985a 100644 --- a/.circleci/config.yml +++ b/.circleci/config.yml @@ -45,13 +45,6 @@ jobs: paths: - ./venv key: v1-dependencies-{{ checksum ".circleci/requirements.txt" }} - - run: - name: Run prisma ./entrypoint.sh - command: | - set +e - chmod +x entrypoint.sh - ./entrypoint.sh - set -e - run: name: Black Formatting command: | diff --git a/litellm/tests/test_completion_with_retries.py b/litellm/tests/test_completion_with_retries.py index 42279453137..e59d1d6e13a 100644 --- a/litellm/tests/test_completion_with_retries.py +++ b/litellm/tests/test_completion_with_retries.py @@ -29,20 +29,6 @@ def logger_fn(user_model_dict): pass -# normal call -def test_completion_custom_provider_model_name(): - try: - response = completion_with_retries( - model="together_ai/togethercomputer/llama-2-70b-chat", - messages=messages, - logger_fn=logger_fn, - ) - # Add any assertions here to check the response - print(response) - except Exception as e: - pytest.fail(f"Error occurred: {e}") - - # completion with num retries + impact on exception mapping def test_completion_with_num_retries(): try: @@ -75,7 +61,3 @@ def test_completion_with_0_num_retries(): except Exception as e: print("exception", e) pass - - -# Call the test function -test_completion_with_0_num_retries() From e44d3e51aa0ac15690c13985d039493a118faf4a Mon Sep 17 00:00:00 2001 From: Krrish Dholakia Date: Wed, 10 Jan 2024 21:26:38 +0530 Subject: [PATCH 5/5] ci(config.yml): run prisma generate before testing --- .circleci/config.yml | 7 ++ docs/my-website/docs/routing.md | 113 +++++++++++++++++--------------- 2 files changed, 67 insertions(+), 53 deletions(-) diff --git a/.circleci/config.yml b/.circleci/config.yml index 2f2b581985a..3ea6b7fca9a 100644 --- a/.circleci/config.yml +++ b/.circleci/config.yml @@ -45,6 +45,13 @@ jobs: paths: - ./venv key: v1-dependencies-{{ checksum ".circleci/requirements.txt" }} + - run: + name: Run prisma ./entrypoint.sh + command: | + set +e + chmod +x entrypoint.sh + ./entrypoint.sh + set -e - run: name: Black Formatting command: | diff --git a/docs/my-website/docs/routing.md b/docs/my-website/docs/routing.md index ce2f491b219..2115e280277 100644 --- a/docs/my-website/docs/routing.md +++ b/docs/my-website/docs/routing.md @@ -77,7 +77,65 @@ print(response) Router provides 4 strategies for routing your calls across multiple deployments: - + + + +Picks the deployment with the lowest response time. + +It caches, and updates the response times for deployments based on when a request was sent and received from a deployment. + +[**How to test**](https://github.com/BerriAI/litellm/blob/main/litellm/tests/test_lowest_latency_routing.py) + +```python +from litellm import Router +import asyncio + +model_list = [{ ... }] + +# init router +router = Router(model_list=model_list, routing_strategy="latency-based-routing") # 👈 set routing strategy + +## CALL 1+2 +tasks = [] +response = None +final_response = None +for _ in range(2): + tasks.append(router.acompletion(model=model, messages=messages)) +response = await asyncio.gather(*tasks) + +if response is not None: + ## CALL 3 + await asyncio.sleep(1) # let the cache update happen + picked_deployment = router.lowestlatency_logger.get_available_deployments( + model_group=model, healthy_deployments=router.healthy_deployments + ) + final_response = await router.acompletion(model=model, messages=messages) + print(f"min deployment id: {picked_deployment}") + print(f"model id: {final_response._hidden_params['model_id']}") + assert ( + final_response._hidden_params["model_id"] + == picked_deployment["model_info"]["id"] + ) +``` + +### Set Time Window + +Set time window for how far back to consider when averaging latency for a deployment. + +**In Router** +```python +router = Router(..., routing_strategy_args={"ttl": 10}) +``` + +**In Proxy** + +```yaml +router_settings: + routing_strategy_args: {"ttl": 10} +``` + + + **Default** Picks a deployment based on the provided **Requests per minute (rpm) or Tokens per minute (tpm)** @@ -235,58 +293,7 @@ asyncio.run(router_acompletion()) ``` - - -Picks the deployment with the lowest response time. - -It caches, and updates the response times for deployments based on when a request was sent and received from a deployment. - -[**How to test**](https://github.com/BerriAI/litellm/blob/main/litellm/tests/test_lowest_latency_routing.py) - -```python -from litellm import Router -import asyncio - -model_list = [{ # list of model deployments - "model_name": "gpt-3.5-turbo", # model alias - "litellm_params": { # params for litellm completion/embedding call - "model": "azure/chatgpt-v-2", # actual model name - "api_key": os.getenv("AZURE_API_KEY"), - "api_version": os.getenv("AZURE_API_VERSION"), - "api_base": os.getenv("AZURE_API_BASE"), - } -}, { - "model_name": "gpt-3.5-turbo", - "litellm_params": { # params for litellm completion/embedding call - "model": "azure/chatgpt-functioncalling", - "api_key": os.getenv("AZURE_API_KEY"), - "api_version": os.getenv("AZURE_API_VERSION"), - "api_base": os.getenv("AZURE_API_BASE"), - } -}, { - "model_name": "gpt-3.5-turbo", - "litellm_params": { # params for litellm completion/embedding call - "model": "gpt-3.5-turbo", - "api_key": os.getenv("OPENAI_API_KEY"), - } -}] - -# init router -router = Router(model_list=model_list, routing_strategy="latency-based-routing") -async def router_acompletion(): - response = await router.acompletion( - model="gpt-3.5-turbo", - messages=[{"role": "user", "content": "Hey, how's it going?"}] - ) - print(response) - return response - -asyncio.run(router_acompletion()) -``` - - - ## Basic Reliability @@ -608,4 +615,4 @@ def __init__( "latency-based-routing", ] = "simple-shuffle", ): -``` +``` \ No newline at end of file