From f3fdae6e30d13ae310de866aa98482ebf37f3611 Mon Sep 17 00:00:00 2001 From: onukura <26293997+onukura@users.noreply.github.com> Date: Thu, 4 Apr 2024 14:30:52 +0900 Subject: [PATCH 001/197] Update README.md --- README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/README.md b/README.md index 5a1440f970e..b5f63287c82 100644 --- a/README.md +++ b/README.md @@ -220,7 +220,7 @@ curl 'http://0.0.0.0:4000/key/generate' \ | [nlp_cloud](https://docs.litellm.ai/docs/providers/nlp_cloud) | ✅ | ✅ | ✅ | ✅ | | [aleph alpha](https://docs.litellm.ai/docs/providers/aleph_alpha) | ✅ | ✅ | ✅ | ✅ | | [petals](https://docs.litellm.ai/docs/providers/petals) | ✅ | ✅ | ✅ | ✅ | -| [ollama](https://docs.litellm.ai/docs/providers/ollama) | ✅ | ✅ | ✅ | ✅ | +| [ollama](https://docs.litellm.ai/docs/providers/ollama) | ✅ | ✅ | ✅ | ✅ | ✅ | | [deepinfra](https://docs.litellm.ai/docs/providers/deepinfra) | ✅ | ✅ | ✅ | ✅ | | [perplexity-ai](https://docs.litellm.ai/docs/providers/perplexity) | ✅ | ✅ | ✅ | ✅ | | [Groq AI](https://docs.litellm.ai/docs/providers/groq) | ✅ | ✅ | ✅ | ✅ | From 58c4b024479d7a7549ffc102c68d0906c3d0afc0 Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Sat, 6 Apr 2024 15:50:13 -0700 Subject: [PATCH 002/197] feat - make anthropic async --- litellm/llms/anthropic.py | 365 ++++++++++++++-------- litellm/llms/custom_httpx/http_handler.py | 4 +- litellm/main.py | 2 + 3 files changed, 231 insertions(+), 140 deletions(-) diff --git a/litellm/llms/anthropic.py b/litellm/llms/anthropic.py index b7b078b9b12..db41ae6e358 100644 --- a/litellm/llms/anthropic.py +++ b/litellm/llms/anthropic.py @@ -7,6 +7,9 @@ from typing import Callable, Optional, List from litellm.utils import ModelResponse, Usage, map_finish_reason, CustomStreamWrapper import litellm from .prompt_templates.factory import prompt_factory, custom_prompt +from litellm.llms.custom_httpx.http_handler import AsyncHTTPHandler + +async_handler = AsyncHTTPHandler() import httpx @@ -36,7 +39,9 @@ class AnthropicConfig: to pass metadata to anthropic, it's {"user_id": "any-relevant-information"} """ - max_tokens: Optional[int] = 4096 # anthropic requires a default value (Opus, Sonnet, and Haiku have the same default) + max_tokens: Optional[int] = ( + 4096 # anthropic requires a default value (Opus, Sonnet, and Haiku have the same default) + ) stop_sequences: Optional[list] = None temperature: Optional[int] = None top_p: Optional[int] = None @@ -46,7 +51,9 @@ class AnthropicConfig: def __init__( self, - max_tokens: Optional[int] = 4096, # You can pass in a value yourself or use the default value 4096 + max_tokens: Optional[ + int + ] = 4096, # You can pass in a value yourself or use the default value 4096 stop_sequences: Optional[list] = None, temperature: Optional[int] = None, top_p: Optional[int] = None, @@ -95,6 +102,169 @@ def validate_environment(api_key, user_headers): return headers +def process_response( + model, + response, + model_response, + _is_function_call, + stream, + logging_obj, + api_key, + data, + messages, + print_verbose, +): + ## LOGGING + logging_obj.post_call( + input=messages, + api_key=api_key, + original_response=response.text, + additional_args={"complete_input_dict": data}, + ) + print_verbose(f"raw model_response: {response.text}") + ## RESPONSE OBJECT + try: + completion_response = response.json() + except: + raise AnthropicError(message=response.text, status_code=response.status_code) + if "error" in completion_response: + raise AnthropicError( + message=str(completion_response["error"]), + status_code=response.status_code, + ) + elif len(completion_response["content"]) == 0: + raise AnthropicError( + message="No content in response", + status_code=response.status_code, + ) + else: + text_content = "" + tool_calls = [] + for content in completion_response["content"]: + if content["type"] == "text": + text_content += content["text"] + ## TOOL CALLING + elif content["type"] == "tool_use": + tool_calls.append( + { + "id": content["id"], + "type": "function", + "function": { + "name": content["name"], + "arguments": json.dumps(content["input"]), + }, + } + ) + + _message = litellm.Message( + tool_calls=tool_calls, + content=text_content or None, + ) + model_response.choices[0].message = _message # type: ignore + model_response._hidden_params["original_response"] = completion_response[ + "content" + ] # allow user to access raw anthropic tool calling response + + model_response.choices[0].finish_reason = map_finish_reason( + completion_response["stop_reason"] + ) + + print_verbose(f"_is_function_call: {_is_function_call}; stream: {stream}") + if _is_function_call and stream: + print_verbose("INSIDE ANTHROPIC STREAMING TOOL CALLING CONDITION BLOCK") + # return an iterator + streaming_model_response = ModelResponse(stream=True) + streaming_model_response.choices[0].finish_reason = model_response.choices[ + 0 + ].finish_reason + # streaming_model_response.choices = [litellm.utils.StreamingChoices()] + streaming_choice = litellm.utils.StreamingChoices() + streaming_choice.index = model_response.choices[0].index + _tool_calls = [] + print_verbose( + f"type of model_response.choices[0]: {type(model_response.choices[0])}" + ) + print_verbose(f"type of streaming_choice: {type(streaming_choice)}") + if isinstance(model_response.choices[0], litellm.Choices): + if getattr( + model_response.choices[0].message, "tool_calls", None + ) is not None and isinstance( + model_response.choices[0].message.tool_calls, list + ): + for tool_call in model_response.choices[0].message.tool_calls: + _tool_call = {**tool_call.dict(), "index": 0} + _tool_calls.append(_tool_call) + delta_obj = litellm.utils.Delta( + content=getattr(model_response.choices[0].message, "content", None), + role=model_response.choices[0].message.role, + tool_calls=_tool_calls, + ) + streaming_choice.delta = delta_obj + streaming_model_response.choices = [streaming_choice] + completion_stream = ModelResponseIterator( + model_response=streaming_model_response + ) + print_verbose( + "Returns anthropic CustomStreamWrapper with 'cached_response' streaming object" + ) + return CustomStreamWrapper( + completion_stream=completion_stream, + model=model, + custom_llm_provider="cached_response", + logging_obj=logging_obj, + ) + + ## CALCULATING USAGE + prompt_tokens = completion_response["usage"]["input_tokens"] + completion_tokens = completion_response["usage"]["output_tokens"] + total_tokens = prompt_tokens + completion_tokens + + model_response["created"] = int(time.time()) + model_response["model"] = model + usage = Usage( + prompt_tokens=prompt_tokens, + completion_tokens=completion_tokens, + total_tokens=total_tokens, + ) + model_response.usage = usage + return model_response + + +async def acompletion_function( + model: str, + messages: list, + api_base: str, + custom_prompt_dict: dict, + model_response: ModelResponse, + print_verbose: Callable, + encoding, + api_key, + logging_obj, + stream, + _is_function_call, + data=None, + optional_params=None, + litellm_params=None, + logger_fn=None, + headers={}, +): + response = await async_handler.post( + api_base, headers=headers, data=json.dumps(data) + ) + return process_response( + model=model, + response=response, + model_response=model_response, + _is_function_call=_is_function_call, + stream=stream, + logging_obj=logging_obj, + api_key=api_key, + data=data, + messages=messages, + print_verbose=print_verbose, + ) + + def completion( model: str, messages: list, @@ -106,6 +276,7 @@ def completion( api_key, logging_obj, optional_params=None, + acompletion=None, litellm_params=None, logger_fn=None, headers={}, @@ -184,148 +355,66 @@ def completion( }, ) print_verbose(f"_is_function_call: {_is_function_call}") - ## COMPLETION CALL - if ( - stream and not _is_function_call - ): # if function call - fake the streaming (need complete blocks for output parsing in openai format) - print_verbose("makes anthropic streaming POST request") - data["stream"] = stream - response = requests.post( - api_base, - headers=headers, - data=json.dumps(data), - stream=stream, - ) - - if response.status_code != 200: - raise AnthropicError( - status_code=response.status_code, message=response.text - ) - - return response.iter_lines() - else: - response = requests.post(api_base, headers=headers, data=json.dumps(data)) - if response.status_code != 200: - raise AnthropicError( - status_code=response.status_code, message=response.text - ) - - ## LOGGING - logging_obj.post_call( - input=messages, - api_key=api_key, - original_response=response.text, - additional_args={"complete_input_dict": data}, - ) - print_verbose(f"raw model_response: {response.text}") - ## RESPONSE OBJECT - try: - completion_response = response.json() - except: - raise AnthropicError( - message=response.text, status_code=response.status_code - ) - if "error" in completion_response: - raise AnthropicError( - message=str(completion_response["error"]), - status_code=response.status_code, - ) - elif len(completion_response["content"]) == 0: - raise AnthropicError( - message="No content in response", - status_code=response.status_code, - ) + if acompletion == True: + if optional_params.get("stream", False): + pass else: - text_content = "" - tool_calls = [] - for content in completion_response["content"]: - if content["type"] == "text": - text_content += content["text"] - ## TOOL CALLING - elif content["type"] == "tool_use": - tool_calls.append( - { - "id": content["id"], - "type": "function", - "function": { - "name": content["name"], - "arguments": json.dumps(content["input"]), - }, - } - ) - - _message = litellm.Message( - tool_calls=tool_calls, - content=text_content or None, + return acompletion_function( + model=model, + messages=messages, + data=data, + api_base=api_base, + custom_prompt_dict=custom_prompt_dict, + model_response=model_response, + print_verbose=print_verbose, + encoding=encoding, + api_key=api_key, + logging_obj=logging_obj, + optional_params=optional_params, + stream=stream, + _is_function_call=_is_function_call, + litellm_params=litellm_params, + logger_fn=logger_fn, + headers=headers, ) - model_response.choices[0].message = _message # type: ignore - model_response._hidden_params["original_response"] = completion_response[ - "content" - ] # allow user to access raw anthropic tool calling response - - model_response.choices[0].finish_reason = map_finish_reason( - completion_response["stop_reason"] + else: + ## COMPLETION CALL + if ( + stream and not _is_function_call + ): # if function call - fake the streaming (need complete blocks for output parsing in openai format) + print_verbose("makes anthropic streaming POST request") + data["stream"] = stream + response = requests.post( + api_base, + headers=headers, + data=json.dumps(data), + stream=stream, ) - print_verbose(f"_is_function_call: {_is_function_call}; stream: {stream}") - if _is_function_call and stream: - print_verbose("INSIDE ANTHROPIC STREAMING TOOL CALLING CONDITION BLOCK") - # return an iterator - streaming_model_response = ModelResponse(stream=True) - streaming_model_response.choices[0].finish_reason = model_response.choices[ - 0 - ].finish_reason - # streaming_model_response.choices = [litellm.utils.StreamingChoices()] - streaming_choice = litellm.utils.StreamingChoices() - streaming_choice.index = model_response.choices[0].index - _tool_calls = [] - print_verbose( - f"type of model_response.choices[0]: {type(model_response.choices[0])}" - ) - print_verbose(f"type of streaming_choice: {type(streaming_choice)}") - if isinstance(model_response.choices[0], litellm.Choices): - if getattr( - model_response.choices[0].message, "tool_calls", None - ) is not None and isinstance( - model_response.choices[0].message.tool_calls, list - ): - for tool_call in model_response.choices[0].message.tool_calls: - _tool_call = {**tool_call.dict(), "index": 0} - _tool_calls.append(_tool_call) - delta_obj = litellm.utils.Delta( - content=getattr(model_response.choices[0].message, "content", None), - role=model_response.choices[0].message.role, - tool_calls=_tool_calls, - ) - streaming_choice.delta = delta_obj - streaming_model_response.choices = [streaming_choice] - completion_stream = ModelResponseIterator( - model_response=streaming_model_response - ) - print_verbose( - "Returns anthropic CustomStreamWrapper with 'cached_response' streaming object" - ) - return CustomStreamWrapper( - completion_stream=completion_stream, - model=model, - custom_llm_provider="cached_response", - logging_obj=logging_obj, + if response.status_code != 200: + raise AnthropicError( + status_code=response.status_code, message=response.text ) - ## CALCULATING USAGE - prompt_tokens = completion_response["usage"]["input_tokens"] - completion_tokens = completion_response["usage"]["output_tokens"] - total_tokens = prompt_tokens + completion_tokens - - model_response["created"] = int(time.time()) - model_response["model"] = model - usage = Usage( - prompt_tokens=prompt_tokens, - completion_tokens=completion_tokens, - total_tokens=total_tokens, - ) - model_response.usage = usage - return model_response + return response.iter_lines() + else: + response = requests.post(api_base, headers=headers, data=json.dumps(data)) + if response.status_code != 200: + raise AnthropicError( + status_code=response.status_code, message=response.text + ) + return process_response( + model=model, + response=response, + model_response=model_response, + _is_function_call=_is_function_call, + stream=stream, + logging_obj=logging_obj, + api_key=api_key, + data=data, + messages=messages, + print_verbose=print_verbose, + ) class ModelResponseIterator: diff --git a/litellm/llms/custom_httpx/http_handler.py b/litellm/llms/custom_httpx/http_handler.py index 10314d831c5..51723a2f995 100644 --- a/litellm/llms/custom_httpx/http_handler.py +++ b/litellm/llms/custom_httpx/http_handler.py @@ -1,5 +1,5 @@ import httpx, asyncio -from typing import Optional +from typing import Optional, Union class AsyncHTTPHandler: @@ -25,7 +25,7 @@ class AsyncHTTPHandler: async def post( self, url: str, - data: Optional[dict] = None, + data: Optional[Union[dict, str]] = None, params: Optional[dict] = None, headers: Optional[dict] = None, ): diff --git a/litellm/main.py b/litellm/main.py index 5a9eb6e451d..b7e5a3ba921 100644 --- a/litellm/main.py +++ b/litellm/main.py @@ -304,6 +304,7 @@ async def acompletion( or custom_llm_provider == "vertex_ai" or custom_llm_provider == "gemini" or custom_llm_provider == "sagemaker" + or custom_llm_provider == "anthropic" or custom_llm_provider in litellm.openai_compatible_providers ): # currently implemented aiohttp calls for just azure, openai, hf, ollama, vertex ai soon all. init_response = await loop.run_in_executor(None, func_with_context) @@ -1184,6 +1185,7 @@ def completion( model=model, messages=messages, api_base=api_base, + acompletion=acompletion, custom_prompt_dict=litellm.custom_prompt_dict, model_response=model_response, print_verbose=print_verbose, From 8e5e99533b77ae66390885a1e0c40fba5f7b651a Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Sat, 6 Apr 2024 17:34:23 -0700 Subject: [PATCH 003/197] async streaming for anthropic --- litellm/main.py | 13 ------------- 1 file changed, 13 deletions(-) diff --git a/litellm/main.py b/litellm/main.py index b7e5a3ba921..75382b8ce3c 100644 --- a/litellm/main.py +++ b/litellm/main.py @@ -1197,19 +1197,6 @@ def completion( logging_obj=logging, headers=headers, ) - if ( - "stream" in optional_params - and optional_params["stream"] == True - and not isinstance(response, CustomStreamWrapper) - ): - # don't try to access stream object, - response = CustomStreamWrapper( - response, - model, - custom_llm_provider="anthropic", - logging_obj=logging, - ) - if optional_params.get("stream", False) or acompletion == True: ## LOGGING logging.post_call( From 7849c29f70a12660e0798237888c6656cd6213cc Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Sat, 6 Apr 2024 17:36:56 -0700 Subject: [PATCH 004/197] async anthropic streaming --- litellm/utils.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/litellm/utils.py b/litellm/utils.py index 5153a414bcc..1e67d63e540 100644 --- a/litellm/utils.py +++ b/litellm/utils.py @@ -8710,7 +8710,9 @@ class CustomStreamWrapper: return hold, curr_chunk def handle_anthropic_chunk(self, chunk): - str_line = chunk.decode("utf-8") # Convert bytes to string + str_line = chunk + if isinstance(chunk, bytes): # Handle binary data + str_line = chunk.decode("utf-8") # Convert bytes to string text = "" is_finished = False finish_reason = None @@ -9970,6 +9972,7 @@ class CustomStreamWrapper: or self.custom_llm_provider == "custom_openai" or self.custom_llm_provider == "text-completion-openai" or self.custom_llm_provider == "azure_text" + or self.custom_llm_provider == "anthropic" or self.custom_llm_provider == "huggingface" or self.custom_llm_provider == "ollama" or self.custom_llm_provider == "ollama_chat" From 5c796b436512d1f9addb303c634db1842d9116a3 Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Sat, 6 Apr 2024 17:53:06 -0700 Subject: [PATCH 005/197] async streaming anthropic --- litellm/llms/anthropic.py | 78 +++++++++++++++++++++-- litellm/llms/custom_httpx/http_handler.py | 19 ++++-- 2 files changed, 87 insertions(+), 10 deletions(-) diff --git a/litellm/llms/anthropic.py b/litellm/llms/anthropic.py index db41ae6e358..47e485ecb86 100644 --- a/litellm/llms/anthropic.py +++ b/litellm/llms/anthropic.py @@ -9,8 +9,6 @@ import litellm from .prompt_templates.factory import prompt_factory, custom_prompt from litellm.llms.custom_httpx.http_handler import AsyncHTTPHandler -async_handler = AsyncHTTPHandler() - import httpx @@ -18,6 +16,11 @@ class AnthropicConstants(Enum): HUMAN_PROMPT = "\n\nHuman: " AI_PROMPT = "\n\nAssistant: " + # constants from https://github.com/anthropics/anthropic-sdk-python/blob/main/src/anthropic/_constants.py + + +async_handler = AsyncHTTPHandler(timeout=httpx.Timeout(timeout=600.0, connect=5.0)) + class AnthropicError(Exception): def __init__(self, status_code, message): @@ -230,6 +233,42 @@ def process_response( return model_response +async def acompletion_stream_function( + model: str, + messages: list, + api_base: str, + custom_prompt_dict: dict, + model_response: ModelResponse, + print_verbose: Callable, + encoding, + api_key, + logging_obj, + stream, + _is_function_call, + data=None, + optional_params=None, + litellm_params=None, + logger_fn=None, + headers={}, +): + response = await async_handler.post( + api_base, headers=headers, data=json.dumps(data) + ) + + if response.status_code != 200: + raise AnthropicError(status_code=response.status_code, message=response.text) + + completion_stream = response.aiter_lines() + + streamwrapper = CustomStreamWrapper( + completion_stream=completion_stream, + model=model, + custom_llm_provider="anthropic", + logging_obj=logging_obj, + ) + return streamwrapper + + async def acompletion_function( model: str, messages: list, @@ -356,8 +395,29 @@ def completion( ) print_verbose(f"_is_function_call: {_is_function_call}") if acompletion == True: - if optional_params.get("stream", False): - pass + if ( + stream and not _is_function_call + ): # if function call - fake the streaming (need complete blocks for output parsing in openai format) + print_verbose("makes async anthropic streaming POST request") + data["stream"] = stream + return acompletion_stream_function( + model=model, + messages=messages, + data=data, + api_base=api_base, + custom_prompt_dict=custom_prompt_dict, + model_response=model_response, + print_verbose=print_verbose, + encoding=encoding, + api_key=api_key, + logging_obj=logging_obj, + optional_params=optional_params, + stream=stream, + _is_function_call=_is_function_call, + litellm_params=litellm_params, + logger_fn=logger_fn, + headers=headers, + ) else: return acompletion_function( model=model, @@ -396,7 +456,15 @@ def completion( status_code=response.status_code, message=response.text ) - return response.iter_lines() + completion_stream = response.iter_lines() + streaming_response = CustomStreamWrapper( + completion_stream=completion_stream, + model=model, + custom_llm_provider="anthropic", + logging_obj=logging_obj, + ) + return streaming_response + else: response = requests.post(api_base, headers=headers, data=json.dumps(data)) if response.status_code != 200: diff --git a/litellm/llms/custom_httpx/http_handler.py b/litellm/llms/custom_httpx/http_handler.py index 51723a2f995..c008b059330 100644 --- a/litellm/llms/custom_httpx/http_handler.py +++ b/litellm/llms/custom_httpx/http_handler.py @@ -1,15 +1,21 @@ import httpx, asyncio -from typing import Optional, Union +from typing import Optional, Union, Mapping, Any + +# https://www.python-httpx.org/advanced/timeouts +_DEFAULT_TIMEOUT = httpx.Timeout(timeout=5.0, connect=5.0) class AsyncHTTPHandler: - def __init__(self, concurrent_limit=1000): + def __init__( + self, timeout: httpx.Timeout = _DEFAULT_TIMEOUT, concurrent_limit=1000 + ): # Create a client with a connection pool self.client = httpx.AsyncClient( + timeout=timeout, limits=httpx.Limits( max_connections=concurrent_limit, max_keepalive_connections=concurrent_limit, - ) + ), ) async def close(self): @@ -25,12 +31,15 @@ class AsyncHTTPHandler: async def post( self, url: str, - data: Optional[Union[dict, str]] = None, + data: Optional[Union[dict, str]] = None, # type: ignore params: Optional[dict] = None, headers: Optional[dict] = None, ): response = await self.client.post( - url, data=data, params=params, headers=headers + url, + data=data, # type: ignore + params=params, + headers=headers, ) return response From 2cf41d3d9f4e5af65f517ef7a842ba2f07574a2c Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Sat, 6 Apr 2024 17:54:19 -0700 Subject: [PATCH 006/197] async ahtropic streaming --- litellm/llms/anthropic_text.py | 11 +++++++++-- 1 file changed, 9 insertions(+), 2 deletions(-) diff --git a/litellm/llms/anthropic_text.py b/litellm/llms/anthropic_text.py index bccc8c769cf..c9a9adfc26d 100644 --- a/litellm/llms/anthropic_text.py +++ b/litellm/llms/anthropic_text.py @@ -4,7 +4,7 @@ from enum import Enum import requests import time from typing import Callable, Optional -from litellm.utils import ModelResponse, Usage +from litellm.utils import ModelResponse, Usage, CustomStreamWrapper import litellm from .prompt_templates.factory import prompt_factory, custom_prompt import httpx @@ -162,8 +162,15 @@ def completion( raise AnthropicError( status_code=response.status_code, message=response.text ) + completion_stream = response.iter_lines() + stream_response = CustomStreamWrapper( + completion_stream=completion_stream, + model=model, + custom_llm_provider="anthropic", + logging_obj=logging_obj, + ) + return stream_response - return response.iter_lines() else: response = requests.post(api_base, headers=headers, data=json.dumps(data)) if response.status_code != 200: From 548b2b686118306dd751cd6a1815982ad675d70f Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Sat, 6 Apr 2024 17:55:26 -0700 Subject: [PATCH 007/197] test - async claude streaming --- litellm/tests/test_streaming.py | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) diff --git a/litellm/tests/test_streaming.py b/litellm/tests/test_streaming.py index ed629875283..a8e3f768076 100644 --- a/litellm/tests/test_streaming.py +++ b/litellm/tests/test_streaming.py @@ -831,22 +831,25 @@ def test_bedrock_claude_3_streaming(): pytest.fail(f"Error occurred: {e}") -def test_claude_3_streaming_finish_reason(): +@pytest.mark.asyncio +async def test_claude_3_streaming_finish_reason(): try: litellm.set_verbose = True messages = [ {"role": "system", "content": "Be helpful"}, {"role": "user", "content": "What do you know?"}, ] - response: ModelResponse = completion( # type: ignore + response: ModelResponse = await litellm.acompletion( # type: ignore model="claude-3-opus-20240229", messages=messages, stream=True, + max_tokens=10, ) complete_response = "" # Add any assertions here to check the response num_finish_reason = 0 - for idx, chunk in enumerate(response): + async for chunk in response: + print(f"chunk: {chunk}") if isinstance(chunk, ModelResponse): if chunk.choices[0].finish_reason is not None: num_finish_reason += 1 From 9be6b7ec7c450de262555f8f5a29231528b9efd2 Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Sat, 6 Apr 2024 18:07:41 -0700 Subject: [PATCH 008/197] ci/cd run again --- litellm/tests/test_streaming.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/litellm/tests/test_streaming.py b/litellm/tests/test_streaming.py index a8e3f768076..5e7609db9bd 100644 --- a/litellm/tests/test_streaming.py +++ b/litellm/tests/test_streaming.py @@ -846,7 +846,7 @@ async def test_claude_3_streaming_finish_reason(): max_tokens=10, ) complete_response = "" - # Add any assertions here to check the response + # Add any assertions here to-check the response num_finish_reason = 0 async for chunk in response: print(f"chunk: {chunk}") From fcf5aa278b1563fc9479e50705fcba21e258b7e8 Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Sat, 6 Apr 2024 18:19:28 -0700 Subject: [PATCH 009/197] fix - use anthropic class for clients --- litellm/llms/anthropic.py | 762 +++++++++++++++++++------------------- litellm/main.py | 5 +- 2 files changed, 389 insertions(+), 378 deletions(-) diff --git a/litellm/llms/anthropic.py b/litellm/llms/anthropic.py index 47e485ecb86..3eca11beff0 100644 --- a/litellm/llms/anthropic.py +++ b/litellm/llms/anthropic.py @@ -8,7 +8,7 @@ from litellm.utils import ModelResponse, Usage, map_finish_reason, CustomStreamW import litellm from .prompt_templates.factory import prompt_factory, custom_prompt from litellm.llms.custom_httpx.http_handler import AsyncHTTPHandler - +from .base import BaseLLM import httpx @@ -19,9 +19,6 @@ class AnthropicConstants(Enum): # constants from https://github.com/anthropics/anthropic-sdk-python/blob/main/src/anthropic/_constants.py -async_handler = AsyncHTTPHandler(timeout=httpx.Timeout(timeout=600.0, connect=5.0)) - - class AnthropicError(Exception): def __init__(self, status_code, message): self.status_code = status_code @@ -105,384 +102,402 @@ def validate_environment(api_key, user_headers): return headers -def process_response( - model, - response, - model_response, - _is_function_call, - stream, - logging_obj, - api_key, - data, - messages, - print_verbose, -): - ## LOGGING - logging_obj.post_call( - input=messages, - api_key=api_key, - original_response=response.text, - additional_args={"complete_input_dict": data}, - ) - print_verbose(f"raw model_response: {response.text}") - ## RESPONSE OBJECT - try: - completion_response = response.json() - except: - raise AnthropicError(message=response.text, status_code=response.status_code) - if "error" in completion_response: - raise AnthropicError( - message=str(completion_response["error"]), - status_code=response.status_code, - ) - elif len(completion_response["content"]) == 0: - raise AnthropicError( - message="No content in response", - status_code=response.status_code, - ) - else: - text_content = "" - tool_calls = [] - for content in completion_response["content"]: - if content["type"] == "text": - text_content += content["text"] - ## TOOL CALLING - elif content["type"] == "tool_use": - tool_calls.append( - { - "id": content["id"], - "type": "function", - "function": { - "name": content["name"], - "arguments": json.dumps(content["input"]), - }, - } - ) - - _message = litellm.Message( - tool_calls=tool_calls, - content=text_content or None, - ) - model_response.choices[0].message = _message # type: ignore - model_response._hidden_params["original_response"] = completion_response[ - "content" - ] # allow user to access raw anthropic tool calling response - - model_response.choices[0].finish_reason = map_finish_reason( - completion_response["stop_reason"] +class AnthropicChatCompletion(BaseLLM): + def __init__(self) -> None: + super().__init__() + self.async_handler = AsyncHTTPHandler( + timeout=httpx.Timeout(timeout=600.0, connect=5.0) ) - print_verbose(f"_is_function_call: {_is_function_call}; stream: {stream}") - if _is_function_call and stream: - print_verbose("INSIDE ANTHROPIC STREAMING TOOL CALLING CONDITION BLOCK") - # return an iterator - streaming_model_response = ModelResponse(stream=True) - streaming_model_response.choices[0].finish_reason = model_response.choices[ - 0 - ].finish_reason - # streaming_model_response.choices = [litellm.utils.StreamingChoices()] - streaming_choice = litellm.utils.StreamingChoices() - streaming_choice.index = model_response.choices[0].index - _tool_calls = [] - print_verbose( - f"type of model_response.choices[0]: {type(model_response.choices[0])}" + def process_response( + self, + model, + response, + model_response, + _is_function_call, + stream, + logging_obj, + api_key, + data, + messages, + print_verbose, + ): + ## LOGGING + logging_obj.post_call( + input=messages, + api_key=api_key, + original_response=response.text, + additional_args={"complete_input_dict": data}, ) - print_verbose(f"type of streaming_choice: {type(streaming_choice)}") - if isinstance(model_response.choices[0], litellm.Choices): - if getattr( - model_response.choices[0].message, "tool_calls", None - ) is not None and isinstance( - model_response.choices[0].message.tool_calls, list - ): - for tool_call in model_response.choices[0].message.tool_calls: - _tool_call = {**tool_call.dict(), "index": 0} - _tool_calls.append(_tool_call) - delta_obj = litellm.utils.Delta( - content=getattr(model_response.choices[0].message, "content", None), - role=model_response.choices[0].message.role, - tool_calls=_tool_calls, - ) - streaming_choice.delta = delta_obj - streaming_model_response.choices = [streaming_choice] - completion_stream = ModelResponseIterator( - model_response=streaming_model_response - ) - print_verbose( - "Returns anthropic CustomStreamWrapper with 'cached_response' streaming object" - ) - return CustomStreamWrapper( - completion_stream=completion_stream, - model=model, - custom_llm_provider="cached_response", - logging_obj=logging_obj, - ) - - ## CALCULATING USAGE - prompt_tokens = completion_response["usage"]["input_tokens"] - completion_tokens = completion_response["usage"]["output_tokens"] - total_tokens = prompt_tokens + completion_tokens - - model_response["created"] = int(time.time()) - model_response["model"] = model - usage = Usage( - prompt_tokens=prompt_tokens, - completion_tokens=completion_tokens, - total_tokens=total_tokens, - ) - model_response.usage = usage - return model_response - - -async def acompletion_stream_function( - model: str, - messages: list, - api_base: str, - custom_prompt_dict: dict, - model_response: ModelResponse, - print_verbose: Callable, - encoding, - api_key, - logging_obj, - stream, - _is_function_call, - data=None, - optional_params=None, - litellm_params=None, - logger_fn=None, - headers={}, -): - response = await async_handler.post( - api_base, headers=headers, data=json.dumps(data) - ) - - if response.status_code != 200: - raise AnthropicError(status_code=response.status_code, message=response.text) - - completion_stream = response.aiter_lines() - - streamwrapper = CustomStreamWrapper( - completion_stream=completion_stream, - model=model, - custom_llm_provider="anthropic", - logging_obj=logging_obj, - ) - return streamwrapper - - -async def acompletion_function( - model: str, - messages: list, - api_base: str, - custom_prompt_dict: dict, - model_response: ModelResponse, - print_verbose: Callable, - encoding, - api_key, - logging_obj, - stream, - _is_function_call, - data=None, - optional_params=None, - litellm_params=None, - logger_fn=None, - headers={}, -): - response = await async_handler.post( - api_base, headers=headers, data=json.dumps(data) - ) - return process_response( - model=model, - response=response, - model_response=model_response, - _is_function_call=_is_function_call, - stream=stream, - logging_obj=logging_obj, - api_key=api_key, - data=data, - messages=messages, - print_verbose=print_verbose, - ) - - -def completion( - model: str, - messages: list, - api_base: str, - custom_prompt_dict: dict, - model_response: ModelResponse, - print_verbose: Callable, - encoding, - api_key, - logging_obj, - optional_params=None, - acompletion=None, - litellm_params=None, - logger_fn=None, - headers={}, -): - headers = validate_environment(api_key, headers) - _is_function_call = False - messages = copy.deepcopy(messages) - optional_params = copy.deepcopy(optional_params) - if model in custom_prompt_dict: - # check if the model has a registered custom prompt - model_prompt_details = custom_prompt_dict[model] - prompt = custom_prompt( - role_dict=model_prompt_details["roles"], - initial_prompt_value=model_prompt_details["initial_prompt_value"], - final_prompt_value=model_prompt_details["final_prompt_value"], - messages=messages, - ) - else: - # Separate system prompt from rest of message - system_prompt_indices = [] - system_prompt = "" - for idx, message in enumerate(messages): - if message["role"] == "system": - system_prompt += message["content"] - system_prompt_indices.append(idx) - if len(system_prompt_indices) > 0: - for idx in reversed(system_prompt_indices): - messages.pop(idx) - if len(system_prompt) > 0: - optional_params["system"] = system_prompt - # Format rest of message according to anthropic guidelines + print_verbose(f"raw model_response: {response.text}") + ## RESPONSE OBJECT try: - messages = prompt_factory( - model=model, messages=messages, custom_llm_provider="anthropic" + completion_response = response.json() + except: + raise AnthropicError( + message=response.text, status_code=response.status_code ) - except Exception as e: - raise AnthropicError(status_code=400, message=str(e)) - - ## Load Config - config = litellm.AnthropicConfig.get_config() - for k, v in config.items(): - if ( - k not in optional_params - ): # completion(top_k=3) > anthropic_config(top_k=3) <- allows for dynamic variables to be passed in - optional_params[k] = v - - ## Handle Tool Calling - if "tools" in optional_params: - _is_function_call = True - headers["anthropic-beta"] = "tools-2024-04-04" - - anthropic_tools = [] - for tool in optional_params["tools"]: - new_tool = tool["function"] - new_tool["input_schema"] = new_tool.pop("parameters") # rename key - anthropic_tools.append(new_tool) - - optional_params["tools"] = anthropic_tools - - stream = optional_params.pop("stream", None) - - data = { - "model": model, - "messages": messages, - **optional_params, - } - - ## LOGGING - logging_obj.pre_call( - input=messages, - api_key=api_key, - additional_args={ - "complete_input_dict": data, - "api_base": api_base, - "headers": headers, - }, - ) - print_verbose(f"_is_function_call: {_is_function_call}") - if acompletion == True: - if ( - stream and not _is_function_call - ): # if function call - fake the streaming (need complete blocks for output parsing in openai format) - print_verbose("makes async anthropic streaming POST request") - data["stream"] = stream - return acompletion_stream_function( - model=model, - messages=messages, - data=data, - api_base=api_base, - custom_prompt_dict=custom_prompt_dict, - model_response=model_response, - print_verbose=print_verbose, - encoding=encoding, - api_key=api_key, - logging_obj=logging_obj, - optional_params=optional_params, - stream=stream, - _is_function_call=_is_function_call, - litellm_params=litellm_params, - logger_fn=logger_fn, - headers=headers, + if "error" in completion_response: + raise AnthropicError( + message=str(completion_response["error"]), + status_code=response.status_code, + ) + elif len(completion_response["content"]) == 0: + raise AnthropicError( + message="No content in response", + status_code=response.status_code, ) else: - return acompletion_function( - model=model, + text_content = "" + tool_calls = [] + for content in completion_response["content"]: + if content["type"] == "text": + text_content += content["text"] + ## TOOL CALLING + elif content["type"] == "tool_use": + tool_calls.append( + { + "id": content["id"], + "type": "function", + "function": { + "name": content["name"], + "arguments": json.dumps(content["input"]), + }, + } + ) + + _message = litellm.Message( + tool_calls=tool_calls, + content=text_content or None, + ) + model_response.choices[0].message = _message # type: ignore + model_response._hidden_params["original_response"] = completion_response[ + "content" + ] # allow user to access raw anthropic tool calling response + + model_response.choices[0].finish_reason = map_finish_reason( + completion_response["stop_reason"] + ) + + print_verbose(f"_is_function_call: {_is_function_call}; stream: {stream}") + if _is_function_call and stream: + print_verbose("INSIDE ANTHROPIC STREAMING TOOL CALLING CONDITION BLOCK") + # return an iterator + streaming_model_response = ModelResponse(stream=True) + streaming_model_response.choices[0].finish_reason = model_response.choices[ + 0 + ].finish_reason + # streaming_model_response.choices = [litellm.utils.StreamingChoices()] + streaming_choice = litellm.utils.StreamingChoices() + streaming_choice.index = model_response.choices[0].index + _tool_calls = [] + print_verbose( + f"type of model_response.choices[0]: {type(model_response.choices[0])}" + ) + print_verbose(f"type of streaming_choice: {type(streaming_choice)}") + if isinstance(model_response.choices[0], litellm.Choices): + if getattr( + model_response.choices[0].message, "tool_calls", None + ) is not None and isinstance( + model_response.choices[0].message.tool_calls, list + ): + for tool_call in model_response.choices[0].message.tool_calls: + _tool_call = {**tool_call.dict(), "index": 0} + _tool_calls.append(_tool_call) + delta_obj = litellm.utils.Delta( + content=getattr(model_response.choices[0].message, "content", None), + role=model_response.choices[0].message.role, + tool_calls=_tool_calls, + ) + streaming_choice.delta = delta_obj + streaming_model_response.choices = [streaming_choice] + completion_stream = ModelResponseIterator( + model_response=streaming_model_response + ) + print_verbose( + "Returns anthropic CustomStreamWrapper with 'cached_response' streaming object" + ) + return CustomStreamWrapper( + completion_stream=completion_stream, + model=model, + custom_llm_provider="cached_response", + logging_obj=logging_obj, + ) + + ## CALCULATING USAGE + prompt_tokens = completion_response["usage"]["input_tokens"] + completion_tokens = completion_response["usage"]["output_tokens"] + total_tokens = prompt_tokens + completion_tokens + + model_response["created"] = int(time.time()) + model_response["model"] = model + usage = Usage( + prompt_tokens=prompt_tokens, + completion_tokens=completion_tokens, + total_tokens=total_tokens, + ) + model_response.usage = usage + return model_response + + async def acompletion_stream_function( + self, + model: str, + messages: list, + api_base: str, + custom_prompt_dict: dict, + model_response: ModelResponse, + print_verbose: Callable, + encoding, + api_key, + logging_obj, + stream, + _is_function_call, + data=None, + optional_params=None, + litellm_params=None, + logger_fn=None, + headers={}, + ): + response = await self.async_handler.post( + api_base, headers=headers, data=json.dumps(data) + ) + + if response.status_code != 200: + raise AnthropicError( + status_code=response.status_code, message=response.text + ) + + completion_stream = response.aiter_lines() + + streamwrapper = CustomStreamWrapper( + completion_stream=completion_stream, + model=model, + custom_llm_provider="anthropic", + logging_obj=logging_obj, + ) + return streamwrapper + + async def acompletion_function( + self, + model: str, + messages: list, + api_base: str, + custom_prompt_dict: dict, + model_response: ModelResponse, + print_verbose: Callable, + encoding, + api_key, + logging_obj, + stream, + _is_function_call, + data=None, + optional_params=None, + litellm_params=None, + logger_fn=None, + headers={}, + ): + response = await self.async_handler.post( + api_base, headers=headers, data=json.dumps(data) + ) + return self.process_response( + model=model, + response=response, + model_response=model_response, + _is_function_call=_is_function_call, + stream=stream, + logging_obj=logging_obj, + api_key=api_key, + data=data, + messages=messages, + print_verbose=print_verbose, + ) + + def completion( + self, + model: str, + messages: list, + api_base: str, + custom_prompt_dict: dict, + model_response: ModelResponse, + print_verbose: Callable, + encoding, + api_key, + logging_obj, + optional_params=None, + acompletion=None, + litellm_params=None, + logger_fn=None, + headers={}, + ): + headers = validate_environment(api_key, headers) + _is_function_call = False + messages = copy.deepcopy(messages) + optional_params = copy.deepcopy(optional_params) + if model in custom_prompt_dict: + # check if the model has a registered custom prompt + model_prompt_details = custom_prompt_dict[model] + prompt = custom_prompt( + role_dict=model_prompt_details["roles"], + initial_prompt_value=model_prompt_details["initial_prompt_value"], + final_prompt_value=model_prompt_details["final_prompt_value"], messages=messages, - data=data, - api_base=api_base, - custom_prompt_dict=custom_prompt_dict, - model_response=model_response, - print_verbose=print_verbose, - encoding=encoding, - api_key=api_key, - logging_obj=logging_obj, - optional_params=optional_params, - stream=stream, - _is_function_call=_is_function_call, - litellm_params=litellm_params, - logger_fn=logger_fn, - headers=headers, ) - else: - ## COMPLETION CALL - if ( - stream and not _is_function_call - ): # if function call - fake the streaming (need complete blocks for output parsing in openai format) - print_verbose("makes anthropic streaming POST request") - data["stream"] = stream - response = requests.post( - api_base, - headers=headers, - data=json.dumps(data), - stream=stream, - ) - - if response.status_code != 200: - raise AnthropicError( - status_code=response.status_code, message=response.text - ) - - completion_stream = response.iter_lines() - streaming_response = CustomStreamWrapper( - completion_stream=completion_stream, - model=model, - custom_llm_provider="anthropic", - logging_obj=logging_obj, - ) - return streaming_response - else: - response = requests.post(api_base, headers=headers, data=json.dumps(data)) - if response.status_code != 200: - raise AnthropicError( - status_code=response.status_code, message=response.text + # Separate system prompt from rest of message + system_prompt_indices = [] + system_prompt = "" + for idx, message in enumerate(messages): + if message["role"] == "system": + system_prompt += message["content"] + system_prompt_indices.append(idx) + if len(system_prompt_indices) > 0: + for idx in reversed(system_prompt_indices): + messages.pop(idx) + if len(system_prompt) > 0: + optional_params["system"] = system_prompt + # Format rest of message according to anthropic guidelines + try: + messages = prompt_factory( + model=model, messages=messages, custom_llm_provider="anthropic" ) - return process_response( - model=model, - response=response, - model_response=model_response, - _is_function_call=_is_function_call, - stream=stream, - logging_obj=logging_obj, - api_key=api_key, - data=data, - messages=messages, - print_verbose=print_verbose, - ) + except Exception as e: + raise AnthropicError(status_code=400, message=str(e)) + + ## Load Config + config = litellm.AnthropicConfig.get_config() + for k, v in config.items(): + if ( + k not in optional_params + ): # completion(top_k=3) > anthropic_config(top_k=3) <- allows for dynamic variables to be passed in + optional_params[k] = v + + ## Handle Tool Calling + if "tools" in optional_params: + _is_function_call = True + headers["anthropic-beta"] = "tools-2024-04-04" + + anthropic_tools = [] + for tool in optional_params["tools"]: + new_tool = tool["function"] + new_tool["input_schema"] = new_tool.pop("parameters") # rename key + anthropic_tools.append(new_tool) + + optional_params["tools"] = anthropic_tools + + stream = optional_params.pop("stream", None) + + data = { + "model": model, + "messages": messages, + **optional_params, + } + + ## LOGGING + logging_obj.pre_call( + input=messages, + api_key=api_key, + additional_args={ + "complete_input_dict": data, + "api_base": api_base, + "headers": headers, + }, + ) + print_verbose(f"_is_function_call: {_is_function_call}") + if acompletion == True: + if ( + stream and not _is_function_call + ): # if function call - fake the streaming (need complete blocks for output parsing in openai format) + print_verbose("makes async anthropic streaming POST request") + data["stream"] = stream + return self.acompletion_stream_function( + model=model, + messages=messages, + data=data, + api_base=api_base, + custom_prompt_dict=custom_prompt_dict, + model_response=model_response, + print_verbose=print_verbose, + encoding=encoding, + api_key=api_key, + logging_obj=logging_obj, + optional_params=optional_params, + stream=stream, + _is_function_call=_is_function_call, + litellm_params=litellm_params, + logger_fn=logger_fn, + headers=headers, + ) + else: + return self.acompletion_function( + model=model, + messages=messages, + data=data, + api_base=api_base, + custom_prompt_dict=custom_prompt_dict, + model_response=model_response, + print_verbose=print_verbose, + encoding=encoding, + api_key=api_key, + logging_obj=logging_obj, + optional_params=optional_params, + stream=stream, + _is_function_call=_is_function_call, + litellm_params=litellm_params, + logger_fn=logger_fn, + headers=headers, + ) + else: + ## COMPLETION CALL + if ( + stream and not _is_function_call + ): # if function call - fake the streaming (need complete blocks for output parsing in openai format) + print_verbose("makes anthropic streaming POST request") + data["stream"] = stream + response = requests.post( + api_base, + headers=headers, + data=json.dumps(data), + stream=stream, + ) + + if response.status_code != 200: + raise AnthropicError( + status_code=response.status_code, message=response.text + ) + + completion_stream = response.iter_lines() + streaming_response = CustomStreamWrapper( + completion_stream=completion_stream, + model=model, + custom_llm_provider="anthropic", + logging_obj=logging_obj, + ) + return streaming_response + + else: + response = requests.post( + api_base, headers=headers, data=json.dumps(data) + ) + if response.status_code != 200: + raise AnthropicError( + status_code=response.status_code, message=response.text + ) + return self.process_response( + model=model, + response=response, + model_response=model_response, + _is_function_call=_is_function_call, + stream=stream, + logging_obj=logging_obj, + api_key=api_key, + data=data, + messages=messages, + print_verbose=print_verbose, + ) + + def embedding(self): + # logic for parsing in - calling - parsing out model embedding calls + pass class ModelResponseIterator: @@ -509,8 +524,3 @@ class ModelResponseIterator: raise StopAsyncIteration self.is_done = True return self.model_response - - -def embedding(): - # logic for parsing in - calling - parsing out model embedding calls - pass diff --git a/litellm/main.py b/litellm/main.py index 75382b8ce3c..a387c91475d 100644 --- a/litellm/main.py +++ b/litellm/main.py @@ -39,7 +39,6 @@ from litellm.utils import ( get_optional_params_image_gen, ) from .llms import ( - anthropic, anthropic_text, together_ai, ai21, @@ -68,6 +67,7 @@ from .llms import ( from .llms.openai import OpenAIChatCompletion, OpenAITextCompletion from .llms.azure import AzureChatCompletion from .llms.azure_text import AzureTextCompletion +from .llms.anthropic import AnthropicChatCompletion from .llms.huggingface_restapi import Huggingface from .llms.prompt_templates.factory import ( prompt_factory, @@ -99,6 +99,7 @@ from litellm.utils import ( dotenv.load_dotenv() # Loading env variables using dotenv openai_chat_completions = OpenAIChatCompletion() openai_text_completions = OpenAITextCompletion() +anthropic_chat_completions = AnthropicChatCompletion() azure_chat_completions = AzureChatCompletion() azure_text_completions = AzureTextCompletion() huggingface = Huggingface() @@ -1181,7 +1182,7 @@ def completion( or get_secret("ANTHROPIC_API_BASE") or "https://api.anthropic.com/v1/messages" ) - response = anthropic.completion( + response = anthropic_chat_completions.completion( model=model, messages=messages, api_base=api_base, From 2622f0351bc79b0021126725967822a60c4a1408 Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Sat, 6 Apr 2024 18:26:52 -0700 Subject: [PATCH 010/197] (ci/cd) run again --- litellm/tests/test_streaming.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/litellm/tests/test_streaming.py b/litellm/tests/test_streaming.py index 5e7609db9bd..6041788baa6 100644 --- a/litellm/tests/test_streaming.py +++ b/litellm/tests/test_streaming.py @@ -2288,7 +2288,7 @@ async def test_acompletion_claude_3_function_call_with_streaming(): elif chunk.choices[0].finish_reason is not None: # last chunk validate_final_streaming_function_calling_chunk(chunk=chunk) idx += 1 - # raise Exception("it worked!") + # raise Exception("it worked! ") except Exception as e: pytest.fail(f"Error occurred: {e}") From f08486448c0543d3d91df5123a547dfddba82e87 Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Sat, 6 Apr 2024 18:28:07 -0700 Subject: [PATCH 011/197] fix - test streaming --- .circleci/config.yml | 1 + litellm/tests/test_streaming.py | 3 +++ 2 files changed, 4 insertions(+) diff --git a/.circleci/config.yml b/.circleci/config.yml index 92892d3ff84..2bad708cb36 100644 --- a/.circleci/config.yml +++ b/.circleci/config.yml @@ -42,6 +42,7 @@ jobs: pip install lunary==0.2.5 pip install "langfuse==2.7.3" pip install numpydoc + pip install nest-asyncio==1.6.0 pip install traceloop-sdk==0.0.69 pip install openai pip install prisma diff --git a/litellm/tests/test_streaming.py b/litellm/tests/test_streaming.py index 6041788baa6..4ca61303fd4 100644 --- a/litellm/tests/test_streaming.py +++ b/litellm/tests/test_streaming.py @@ -26,6 +26,9 @@ litellm.logging = False litellm.set_verbose = True litellm.num_retries = 3 litellm.cache = None +import nest_asyncio + +nest_asyncio.apply() score = 0 From a38d3b17c5bb499f073ce162b879dfb8fdce3327 Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Sat, 6 Apr 2024 19:16:27 -0700 Subject: [PATCH 012/197] ci/cd run async handler --- litellm/llms/anthropic.py | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) diff --git a/litellm/llms/anthropic.py b/litellm/llms/anthropic.py index 3eca11beff0..d836ed8db54 100644 --- a/litellm/llms/anthropic.py +++ b/litellm/llms/anthropic.py @@ -105,9 +105,6 @@ def validate_environment(api_key, user_headers): class AnthropicChatCompletion(BaseLLM): def __init__(self) -> None: super().__init__() - self.async_handler = AsyncHTTPHandler( - timeout=httpx.Timeout(timeout=600.0, connect=5.0) - ) def process_response( self, @@ -258,6 +255,9 @@ class AnthropicChatCompletion(BaseLLM): logger_fn=None, headers={}, ): + self.async_handler = AsyncHTTPHandler( + timeout=httpx.Timeout(timeout=600.0, connect=5.0) + ) response = await self.async_handler.post( api_base, headers=headers, data=json.dumps(data) ) @@ -296,6 +296,9 @@ class AnthropicChatCompletion(BaseLLM): logger_fn=None, headers={}, ): + self.async_handler = AsyncHTTPHandler( + timeout=httpx.Timeout(timeout=600.0, connect=5.0) + ) response = await self.async_handler.post( api_base, headers=headers, data=json.dumps(data) ) From 9be250c0f0852c4a9851ed2612b97f97b06d54e4 Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Sat, 6 Apr 2024 19:27:26 -0700 Subject: [PATCH 013/197] add exit and aenter --- litellm/llms/custom_httpx/http_handler.py | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/litellm/llms/custom_httpx/http_handler.py b/litellm/llms/custom_httpx/http_handler.py index c008b059330..67e6c80da6f 100644 --- a/litellm/llms/custom_httpx/http_handler.py +++ b/litellm/llms/custom_httpx/http_handler.py @@ -22,6 +22,13 @@ class AsyncHTTPHandler: # Close the client when you're done with it await self.client.aclose() + async def __aenter__(self): + return self.client + + async def __aexit__(self): + # close the client when exiting + await self.client.aclose() + async def get( self, url: str, params: Optional[dict] = None, headers: Optional[dict] = None ): From d51e853b609f36c6e4d1a223df6f701b3bca5204 Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Sat, 6 Apr 2024 19:28:51 -0700 Subject: [PATCH 014/197] undo adding next-asyncio --- .circleci/config.yml | 1 - litellm/tests/test_streaming.py | 3 --- 2 files changed, 4 deletions(-) diff --git a/.circleci/config.yml b/.circleci/config.yml index 2bad708cb36..92892d3ff84 100644 --- a/.circleci/config.yml +++ b/.circleci/config.yml @@ -42,7 +42,6 @@ jobs: pip install lunary==0.2.5 pip install "langfuse==2.7.3" pip install numpydoc - pip install nest-asyncio==1.6.0 pip install traceloop-sdk==0.0.69 pip install openai pip install prisma diff --git a/litellm/tests/test_streaming.py b/litellm/tests/test_streaming.py index 4ca61303fd4..6041788baa6 100644 --- a/litellm/tests/test_streaming.py +++ b/litellm/tests/test_streaming.py @@ -26,9 +26,6 @@ litellm.logging = False litellm.set_verbose = True litellm.num_retries = 3 litellm.cache = None -import nest_asyncio - -nest_asyncio.apply() score = 0 From b4b882c5d63078cd6b51737e14bbf5e87173cd15 Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Sun, 7 Apr 2024 09:57:27 -0700 Subject: [PATCH 015/197] =?UTF-8?q?bump:=20version=201.34.33=20=E2=86=92?= =?UTF-8?q?=201.34.34?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pyproject.toml | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index 60258d00f66..85a58728eed 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [tool.poetry] name = "litellm" -version = "1.34.33" +version = "1.34.34" description = "Library to easily interface with LLM API providers" authors = ["BerriAI"] license = "MIT" @@ -80,7 +80,7 @@ requires = ["poetry-core", "wheel"] build-backend = "poetry.core.masonry.api" [tool.commitizen] -version = "1.34.33" +version = "1.34.34" version_files = [ "pyproject.toml:^version" ] From 559a4cde23863af48b11c0bb1605dec26dab4b80 Mon Sep 17 00:00:00 2001 From: Gregory Nwosu <193151+gregnwosu@users.noreply.github.com> Date: Mon, 8 Apr 2024 02:03:54 +0100 Subject: [PATCH 016/197] created defaults for response["eval_count"] there is no way in litellm to disable the cache in ollama that is removing the eval_count response keys from the json. This PR allows the code to create sensible defaults for when the response is empty see - https://github.com/ollama/ollama/issues/1573 - https://github.com/ollama/ollama/issues/2023 --- litellm/llms/ollama.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/litellm/llms/ollama.py b/litellm/llms/ollama.py index 779896abfdd..a14c3cb5031 100644 --- a/litellm/llms/ollama.py +++ b/litellm/llms/ollama.py @@ -229,7 +229,7 @@ def get_ollama_response( model_response["created"] = int(time.time()) model_response["model"] = "ollama/" + model prompt_tokens = response_json.get("prompt_eval_count", len(encoding.encode(prompt))) # type: ignore - completion_tokens = response_json["eval_count"] + completion_tokens = response_json.get("eval_count", len(response_json.get("message",dict()).get("content", ""))) model_response["usage"] = litellm.Usage( prompt_tokens=prompt_tokens, completion_tokens=completion_tokens, @@ -331,7 +331,7 @@ async def ollama_acompletion(url, data, model_response, encoding, logging_obj): model_response["created"] = int(time.time()) model_response["model"] = "ollama/" + data["model"] prompt_tokens = response_json.get("prompt_eval_count", len(encoding.encode(data["prompt"]))) # type: ignore - completion_tokens = response_json["eval_count"] + completion_tokens = response_json.get("eval_count", len(response_json.get("message",dict()).get("content", ""))) model_response["usage"] = litellm.Usage( prompt_tokens=prompt_tokens, completion_tokens=completion_tokens, From 1ace1921553e733563f9cb9d90a716a789da80ba Mon Sep 17 00:00:00 2001 From: unclecode Date: Mon, 8 Apr 2024 12:42:24 +0800 Subject: [PATCH 017/197] Fix issue #2832: Add protected_namespaces to Config class within utils.py, router.py and completion.py to avoid the warning message. --- litellm/types/completion.py | 2 +- litellm/types/router.py | 5 ++++- litellm/utils.py | 1 + 3 files changed, 6 insertions(+), 2 deletions(-) diff --git a/litellm/types/completion.py b/litellm/types/completion.py index 5eac9057560..5302d16b9ab 100644 --- a/litellm/types/completion.py +++ b/litellm/types/completion.py @@ -32,5 +32,5 @@ class CompletionRequest(BaseModel): model_list: Optional[List[str]] = None class Config: - # allow kwargs extra = "allow" + protected_namespaces = () \ No newline at end of file diff --git a/litellm/types/router.py b/litellm/types/router.py index dc29bb94919..16f7413d223 100644 --- a/litellm/types/router.py +++ b/litellm/types/router.py @@ -41,6 +41,8 @@ class RouterConfig(BaseModel): "latency-based-routing", ] = "simple-shuffle" + class Config: + protected_namespaces = () class ModelInfo(BaseModel): id: Optional[ @@ -141,7 +143,8 @@ class Deployment(BaseModel): class Config: extra = "allow" - + protected_namespaces = () + def __contains__(self, key): # Define custom behavior for the 'in' operator return hasattr(self, key) diff --git a/litellm/utils.py b/litellm/utils.py index 6a58d56db11..4bf877b0dd5 100644 --- a/litellm/utils.py +++ b/litellm/utils.py @@ -236,6 +236,7 @@ class HiddenParams(OpenAIObject): class Config: extra = "allow" + protected_namespaces = () def get(self, key, default=None): # Custom .get() method to access attributes with a default value if the attribute doesn't exist From 5554e2c3599b628df5a7a04598af9901119ec3b8 Mon Sep 17 00:00:00 2001 From: unclecode Date: Mon, 8 Apr 2024 12:49:40 +0800 Subject: [PATCH 018/197] Continue fixing the issue #2832: Add protected_namespaces to another to Config class within the router.py --- litellm/types/router.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/litellm/types/router.py b/litellm/types/router.py index 16f7413d223..920725131a0 100644 --- a/litellm/types/router.py +++ b/litellm/types/router.py @@ -12,6 +12,9 @@ class ModelConfig(BaseModel): tpm: int rpm: int + class Config: + protected_namespaces = () + class RouterConfig(BaseModel): model_list: List[ModelConfig] @@ -144,7 +147,7 @@ class Deployment(BaseModel): class Config: extra = "allow" protected_namespaces = () - + def __contains__(self, key): # Define custom behavior for the 'in' operator return hasattr(self, key) From d099591a09097ec3eb3dd75803091f34492e10cc Mon Sep 17 00:00:00 2001 From: Krrish Dholakia Date: Mon, 8 Apr 2024 07:29:58 -0700 Subject: [PATCH 019/197] docs(sidebars.js): refactor ordering --- docs/my-website/sidebars.js | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/my-website/sidebars.js b/docs/my-website/sidebars.js index 5d5e24371ee..7ee4b9b4de4 100644 --- a/docs/my-website/sidebars.js +++ b/docs/my-website/sidebars.js @@ -163,7 +163,6 @@ const sidebars = { "debugging/local_debugging", "observability/callbacks", "observability/custom_callback", - "observability/lunary_integration", "observability/langfuse_integration", "observability/sentry", "observability/promptlayer_integration", @@ -171,6 +170,7 @@ const sidebars = { "observability/langsmith_integration", "observability/slack_integration", "observability/traceloop_integration", + "observability/lunary_integration", "observability/athina_integration", "observability/helicone_integration", "observability/supabase_integration", From 28e4706bfdcb41b13589b09b1a5d4bd955599679 Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Mon, 8 Apr 2024 12:02:40 -0700 Subject: [PATCH 020/197] test - re-order embedding responses --- litellm/proxy/tests/test_openai_embedding.py | 126 +++++++++++++++++++ 1 file changed, 126 insertions(+) create mode 100644 litellm/proxy/tests/test_openai_embedding.py diff --git a/litellm/proxy/tests/test_openai_embedding.py b/litellm/proxy/tests/test_openai_embedding.py new file mode 100644 index 00000000000..3763f4edd75 --- /dev/null +++ b/litellm/proxy/tests/test_openai_embedding.py @@ -0,0 +1,126 @@ +import openai +import asyncio + + +async def async_request(client, model, input_data): + response = await client.embeddings.create(model=model, input=input_data) + response = response.dict() + data_list = response["data"] + for i, embedding in enumerate(data_list): + embedding["embedding"] = [] + current_index = embedding["index"] + assert i == current_index + return response + + +async def main(): + client = openai.AsyncOpenAI(api_key="sk-1234", base_url="http://0.0.0.0:4000") + models = [ + "text-embedding-ada-002", + "text-embedding-ada-002", + "text-embedding-ada-002", + ] + inputs = [ + [ + "5", + "6", + "7", + "8", + "9", + "10", + "11", + "12", + "13", + "14", + "15", + "16", + "17", + "18", + "19", + "20", + ], + ["1", "2", "3", "4", "5", "6"], + [ + "1", + "2", + "3", + "4", + "5", + "6", + "7", + "8", + "9", + "10", + "11", + "12", + "13", + "14", + "15", + "16", + "17", + "18", + "19", + "20", + ], + [ + "1", + "2", + "3", + "4", + "5", + "6", + "7", + "8", + "9", + "10", + "11", + "12", + "13", + "14", + "15", + "16", + "17", + "18", + "19", + "20", + ], + [ + "1", + "2", + "3", + "4", + "5", + "6", + "7", + "8", + "9", + "10", + "11", + "12", + "13", + "14", + "15", + "16", + "17", + "18", + "19", + "20", + ], + ["1", "2", "3"], + ] + + tasks = [] + for model, input_data in zip(models, inputs): + task = async_request(client, model, input_data) + tasks.append(task) + + responses = await asyncio.gather(*tasks) + print(responses) + for response in responses: + data_list = response["data"] + for embedding in data_list: + embedding["embedding"] = [] + print(response) + + +asyncio.run(main()) From 48bfc45cb0d27dc2ade60f3797ec234c1d947758 Mon Sep 17 00:00:00 2001 From: Krrish Dholakia Date: Mon, 8 Apr 2024 12:17:57 -0700 Subject: [PATCH 021/197] fix(utils.py): fix reordering of items for cached embeddings ensures cached embedding item is returned in correct order --- litellm/proxy/_new_secret_config.yaml | 4 +++ litellm/tests/test_caching.py | 42 +++++++++++++++++++++++++++ litellm/utils.py | 8 ++++- 3 files changed, 53 insertions(+), 1 deletion(-) diff --git a/litellm/proxy/_new_secret_config.yaml b/litellm/proxy/_new_secret_config.yaml index 90209a9e629..8d0afd067f7 100644 --- a/litellm/proxy/_new_secret_config.yaml +++ b/litellm/proxy/_new_secret_config.yaml @@ -12,6 +12,10 @@ model_list: api_version: "2023-07-01-preview" stream_timeout: 0.001 model_name: azure-gpt-3.5 +- model_name: text-embedding-ada-002 + litellm_params: + model: text-embedding-ada-002 + api_key: os.environ/OPENAI_API_KEY - model_name: gpt-instruct litellm_params: model: gpt-3.5-turbo-instruct diff --git a/litellm/tests/test_caching.py b/litellm/tests/test_caching.py index da46359d967..835d3611b83 100644 --- a/litellm/tests/test_caching.py +++ b/litellm/tests/test_caching.py @@ -345,6 +345,48 @@ async def test_embedding_caching_azure_individual_items(): assert embedding_val_2._hidden_params["cache_hit"] == True +@pytest.mark.asyncio +async def test_embedding_caching_azure_individual_items_reordered(): + """ + Tests caching for individual items in an embedding list + + - Cache an item + - call aembedding(..) with the item + 1 unique item + - compare to a 2nd aembedding (...) with 2 unique items + ``` + embedding_1 = ["hey how's it going", "I'm doing well"] + embedding_val_1 = embedding(...) + + embedding_2 = ["hey how's it going", "I'm fine"] + embedding_val_2 = embedding(...) + + assert embedding_val_1[0]["id"] == embedding_val_2[0]["id"] + ``` + """ + litellm.cache = Cache() + common_msg = f"{uuid.uuid4()}" + common_msg_2 = f"hey how's it going {uuid.uuid4()}" + embedding_1 = [common_msg_2, common_msg] + embedding_2 = [ + common_msg, + f"I'm fine {uuid.uuid4()}", + ] + + embedding_val_1 = await aembedding( + model="azure/azure-embedding-model", input=embedding_1, caching=True + ) + embedding_val_2 = await aembedding( + model="azure/azure-embedding-model", input=embedding_2, caching=True + ) + print(f"embedding_val_2._hidden_params: {embedding_val_2._hidden_params}") + assert embedding_val_2._hidden_params["cache_hit"] == True + + assert embedding_val_2.data[0]["embedding"] == embedding_val_1.data[1]["embedding"] + assert embedding_val_2.data[0]["index"] != embedding_val_1.data[1]["index"] + assert embedding_val_2.data[0]["index"] == 0 + assert embedding_val_1.data[1]["index"] == 1 + + @pytest.mark.asyncio async def test_redis_cache_basic(): """ diff --git a/litellm/utils.py b/litellm/utils.py index 6a58d56db11..affa811f331 100644 --- a/litellm/utils.py +++ b/litellm/utils.py @@ -3174,7 +3174,13 @@ def client(original_function): for val in non_null_list: idx, cr = val # (idx, cr) tuple if cr is not None: - final_embedding_cached_response.data[idx] = cr + final_embedding_cached_response.data[idx] = ( + Embedding( + embedding=cr["embedding"], + index=idx, + object="embedding", + ) + ) if len(remaining_list) == 0: # LOG SUCCESS cache_hit = True From 2fc169e6a0b0de4259ee1ec6fd60b61303d1659e Mon Sep 17 00:00:00 2001 From: Krrish Dholakia Date: Mon, 8 Apr 2024 12:19:11 -0700 Subject: [PATCH 022/197] refactor(main.py): trigger new build --- litellm/main.py | 1 - 1 file changed, 1 deletion(-) diff --git a/litellm/main.py b/litellm/main.py index 1ee16f36ff2..aace6565b2c 100644 --- a/litellm/main.py +++ b/litellm/main.py @@ -12,7 +12,6 @@ from typing import Any, Literal, Union, BinaryIO from functools import partial import dotenv, traceback, random, asyncio, time, contextvars from copy import deepcopy - import httpx import litellm from ._logging import verbose_logger From 75d2eb61b422dee4c9dfc16d8bf0e7f6a415a3e5 Mon Sep 17 00:00:00 2001 From: Krrish Dholakia Date: Mon, 8 Apr 2024 12:19:46 -0700 Subject: [PATCH 023/197] =?UTF-8?q?bump:=20version=201.34.34=20=E2=86=92?= =?UTF-8?q?=201.34.35?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pyproject.toml | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index 85a58728eed..6715e1a676d 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [tool.poetry] name = "litellm" -version = "1.34.34" +version = "1.34.35" description = "Library to easily interface with LLM API providers" authors = ["BerriAI"] license = "MIT" @@ -80,7 +80,7 @@ requires = ["poetry-core", "wheel"] build-backend = "poetry.core.masonry.api" [tool.commitizen] -version = "1.34.34" +version = "1.34.35" version_files = [ "pyproject.toml:^version" ] From 1bb07e54d49d14dd12863a9bb37f19d15a1d0be4 Mon Sep 17 00:00:00 2001 From: Ishaan Jaff Date: Mon, 8 Apr 2024 13:20:28 -0700 Subject: [PATCH 024/197] ui - pass extra litellm params --- .../src/components/model_dashboard.tsx | 33 +++++++++++++++++++ 1 file changed, 33 insertions(+) diff --git a/ui/litellm-dashboard/src/components/model_dashboard.tsx b/ui/litellm-dashboard/src/components/model_dashboard.tsx index 61d2a5277c0..4194b56cc51 100644 --- a/ui/litellm-dashboard/src/components/model_dashboard.tsx +++ b/ui/litellm-dashboard/src/components/model_dashboard.tsx @@ -33,6 +33,7 @@ import { import { Badge, BadgeDelta, Button } from "@tremor/react"; import RequestAccess from "./request_model_access"; import { Typography } from "antd"; +import TextArea from "antd/es/input/TextArea"; const { Title: Title2, Link } = Typography; @@ -183,6 +184,22 @@ const ModelDashboard: React.FC = ({ // Add key-value pair to model_info dictionary modelInfoObj[key] = value; } + + + if (key == "litellm_extra_params") { + console.log("litellm_extra_params:", value); + let litellmExtraParams = {}; + try { + litellmExtraParams = JSON.parse(value); + } + catch (error) { + message.error("Failed to parse LiteLLM Extra Paras: " + error); + throw new Error("Failed to parse litellm_extra_params: " + error); + } + for (const [key, value] of Object.entries(litellmExtraParams)) { + litellmParamsObj[key] = value; + } + } } const new_model: Model = { @@ -420,6 +437,22 @@ const ModelDashboard: React.FC = ({ } + +