litellm/tests/logging_callback_tests/test_amazing_s3_logs.py
yuneng-jiang 6a0d03914c
test: drop the cwd-relative sys.path.insert calls from the test suite (#37802)
* test: drop the cwd-relative sys.path.insert calls from the test suite

TQ003 stands at 1,077 across 1,058 files, and 1,015 of them are the same shape:
sys.path.insert(0, os.path.abspath("../..")) and its deeper siblings. The
argument resolves against the working directory rather than the file, so from
the repo root, where every job runs pytest, it inserts the directory two levels
above the checkout. It has never pointed at litellm. The package is installed
into the environment anyway, which is what actually makes the import work, and
what the rule's message has said all along.

Removing them leaves 1,634 imports of sys and os with no remaining reference,
and those go too, except where another test module imports the name back out of
the file. The rest of TQ003 is 62 call sites that resolve against __file__ or a
variable, which are a different question and are left alone.

Collection is identical either way: 45,871 tests and the same 51 pre-existing
collection errors before and after, and ruff reports no new undefined name.

* test: drop the duplicate imports the sys.path sweep exposed to F811

* test(pre-call-utils): restore the os import the new bedrock tests need
2026-08-22 09:25:58 -07:00

472 lines
15 KiB
Python

import io, asyncio
from collections import defaultdict
# import logging
# logging.basicConfig(level=logging.DEBUG)
from litellm import completion
import litellm
litellm.num_retries = 3
import time, random
import pytest
import boto3
from litellm._logging import verbose_logger
import logging
class _FakeS3Paginator:
def __init__(self, objects):
self.objects = objects
def paginate(self, Bucket):
keys = sorted(self.objects[Bucket])
if not keys:
return [{}]
return [{"Contents": [{"Key": key} for key in keys]}]
class _FakeS3Client:
def __init__(self):
self.objects = defaultdict(dict)
def clear(self):
self.objects.clear()
def put_object(self, Bucket, Key, Body, **_kwargs):
self.objects[Bucket][Key] = Body
return {"ResponseMetadata": {"HTTPStatusCode": 200}}
def delete_object(self, Bucket, Key):
self.objects[Bucket].pop(Key, None)
return {"ResponseMetadata": {"HTTPStatusCode": 204}}
def get_paginator(self, name):
assert name == "list_objects_v2"
return _FakeS3Paginator(self.objects)
def list_objects(self, Bucket):
keys = sorted(self.objects[Bucket])
return {"Contents": [{"Key": key, "LastModified": 0} for key in keys]}
_FAKE_S3_CLIENT = _FakeS3Client()
@pytest.fixture(autouse=True)
def fake_s3_client(monkeypatch):
_FAKE_S3_CLIENT.clear()
def fake_boto3_client(service_name, *args, **kwargs):
assert service_name == "s3"
return _FAKE_S3_CLIENT
monkeypatch.setattr(boto3, "client", fake_boto3_client)
litellm.success_callback = []
litellm.callbacks = []
yield _FAKE_S3_CLIENT
litellm.success_callback = []
litellm.callbacks = []
@pytest.mark.asyncio
@pytest.mark.parametrize(
"sync_mode,streaming", [(True, True), (True, False), (False, True), (False, False)]
)
@pytest.mark.flaky(retries=3, delay=1)
async def test_basic_s3_logging(sync_mode, streaming):
verbose_logger.setLevel(level=logging.DEBUG)
litellm.success_callback = ["s3"]
litellm.s3_callback_params = {
"s3_bucket_name": "load-testing-oct",
"s3_aws_secret_access_key": "os.environ/AWS_SECRET_ACCESS_KEY",
"s3_aws_access_key_id": "os.environ/AWS_ACCESS_KEY_ID",
"s3_region_name": "us-west-2",
}
litellm.set_verbose = True
response_id = None
if sync_mode is True:
response = litellm.completion(
model="gpt-5-mini",
messages=[{"role": "user", "content": "This is a test"}],
mock_response="It's simple to use and easy to get started",
stream=streaming,
)
if streaming:
for chunk in response:
print()
response_id = chunk.id
else:
response_id = response.id
time.sleep(2)
else:
response = await litellm.acompletion(
model="gpt-5-mini",
messages=[{"role": "user", "content": "This is a test"}],
mock_response="It's simple to use and easy to get started",
stream=streaming,
)
if streaming:
async for chunk in response:
print(chunk)
response_id = chunk.id
else:
response_id = response.id
await asyncio.sleep(2)
print(f"response: {response}")
total_objects, all_s3_keys = list_all_s3_objects("load-testing-oct")
# assert that atlest one key has response.id in it
assert any(response_id in key for key in all_s3_keys)
s3 = boto3.client("s3")
# delete all objects
for key in all_s3_keys:
s3.delete_object(Bucket="load-testing-oct", Key=key)
@pytest.mark.asyncio
@pytest.mark.parametrize("streaming", [True])
@pytest.mark.flaky(retries=3, delay=1)
async def test_basic_s3_v2_logging(streaming):
from unittest.mock import AsyncMock, MagicMock, patch
from litellm.integrations.s3_v2 import S3Logger
litellm.s3_callback_params = {
"s3_bucket_name": "load-testing-oct",
"s3_aws_secret_access_key": "test-secret",
"s3_aws_access_key_id": "test-key",
"s3_region_name": "us-west-2",
}
s3_v2_logger = S3Logger(s3_flush_interval=1)
litellm.callbacks = [s3_v2_logger]
uploaded_keys: list = []
original_upload = s3_v2_logger.async_upload_data_to_s3
async def mock_upload(batch_logging_element):
uploaded_keys.append(batch_logging_element.s3_object_key)
s3_v2_logger.async_upload_data_to_s3 = mock_upload
litellm.set_verbose = True
response_id = None
response = await litellm.acompletion(
model="gpt-5-mini",
messages=[{"role": "user", "content": "This is a test"}],
mock_response="It's simple to use and easy to get started",
stream=streaming,
)
if streaming:
async for chunk in response:
response_id = chunk.id
else:
response_id = response.id
await asyncio.sleep(5)
assert len(uploaded_keys) > 0, "S3 upload was never called"
assert any(
response_id in key for key in uploaded_keys
), f"Expected response_id={response_id} in one of the uploaded S3 keys: {uploaded_keys}"
@pytest.mark.asyncio
@pytest.mark.flaky(retries=3, delay=1)
async def test_basic_s3_v2_logging_failure():
"""Test that S3 v2 logger makes httpx PUT request when logging failures"""
from unittest.mock import AsyncMock, MagicMock, patch
from litellm.integrations.s3_v2 import S3Logger
# Create S3 logger with short flush interval
s3_v2_logger = S3Logger(s3_flush_interval=1)
# Mock the httpx client to capture the PUT request
mock_response = MagicMock()
mock_response.status_code = 200
mock_response.raise_for_status = MagicMock()
s3_v2_logger.async_httpx_client = AsyncMock()
s3_v2_logger.async_httpx_client.put.return_value = mock_response
# Track the upload method calls
original_upload = s3_v2_logger.async_upload_data_to_s3
upload_called = False
async def mock_upload(batch_logging_element):
nonlocal upload_called
upload_called = True
# Mock the upload process but still make the httpx call
url = f"https://test-bucket.s3.us-west-2.amazonaws.com/{batch_logging_element.s3_object_key}"
headers = {"Content-Type": "application/json"}
data = '{"model": "gpt-5-mini"}'
# Make the actual httpx call we want to test
await s3_v2_logger.async_httpx_client.put(url=url, headers=headers, data=data)
s3_v2_logger.async_upload_data_to_s3 = mock_upload
# Configure S3 callback params
litellm.callbacks = [s3_v2_logger]
litellm.s3_callback_params = {
"s3_bucket_name": "test-bucket",
"s3_aws_secret_access_key": "test-secret",
"s3_aws_access_key_id": "test-key",
"s3_region_name": "us-west-2",
}
litellm.set_verbose = True
# Trigger a failure by using invalid API key
try:
response = await litellm.acompletion(
model="gpt-5-mini",
api_key="invalid-api-key",
messages=[{"role": "user", "content": "This is a test"}],
mock_response=Exception("forced failure for S3 logging test"),
)
except Exception as e:
print(f"Expected error: {e}")
# Wait for logger to process the failure
await asyncio.sleep(5)
# Verify that our mock upload was called
assert upload_called, "S3 upload method was not called"
print("✓ S3 upload method was called")
# Verify that httpx PUT was called
s3_v2_logger.async_httpx_client.put.assert_called()
# Get the call arguments to verify the S3 URL
call_args = s3_v2_logger.async_httpx_client.put.call_args
assert call_args is not None
url = call_args[1]["url"] if "url" in call_args[1] else call_args[0][0]
# Verify the URL contains expected S3 endpoint
assert "test-bucket.s3.us-west-2.amazonaws.com" in url
print(f"✓ S3 PUT request made to: {url}")
# Verify headers include expected content type
headers = call_args[1]["headers"]
assert headers["Content-Type"] == "application/json"
print("✓ S3 request headers are correct")
# Verify JSON data was included
data = call_args[1]["data"]
assert data is not None
assert '"model": "gpt-5-mini"' in data
print("✓ S3 request data contains expected log payload")
def list_all_s3_objects(bucket_name):
s3 = boto3.client("s3")
all_s3_keys = []
paginator = s3.get_paginator("list_objects_v2")
total_objects = 0
for page in paginator.paginate(Bucket=bucket_name):
if "Contents" in page:
total_objects += len(page["Contents"])
all_s3_keys.extend([obj["Key"] for obj in page["Contents"]])
print(f"Total number of objects in {bucket_name}: {total_objects}")
print(all_s3_keys)
return total_objects, all_s3_keys
@pytest.mark.skip(reason="AWS Suspended Account")
def test_s3_logging():
# all s3 requests need to be in one test function
# since we are modifying stdout, and pytests runs tests in parallel
# on circle ci - we only test litellm.acompletion()
try:
# redirect stdout to log_file
litellm.cache = litellm.Cache(
type="s3",
s3_bucket_name="litellm-my-test-bucket-2",
s3_region_name="us-east-1",
)
litellm.success_callback = ["s3"]
litellm.s3_callback_params = {
"s3_bucket_name": "litellm-logs-2",
"s3_aws_secret_access_key": "os.environ/AWS_SECRET_ACCESS_KEY",
"s3_aws_access_key_id": "os.environ/AWS_ACCESS_KEY_ID",
}
litellm.set_verbose = True
print("Testing async s3 logging")
expected_keys = []
import time
curr_time = str(time.time())
async def _test():
return await litellm.acompletion(
model="gpt-5-mini",
messages=[{"role": "user", "content": f"This is a test {curr_time}"}],
max_tokens=10,
temperature=0.7,
user="ishaan-2",
)
response = asyncio.run(_test())
print(f"response: {response}")
expected_keys.append(response.id)
async def _test():
return await litellm.acompletion(
model="gpt-5-mini",
messages=[{"role": "user", "content": f"This is a test {curr_time}"}],
max_tokens=10,
temperature=0.7,
user="ishaan-2",
)
response = asyncio.run(_test())
expected_keys.append(response.id)
print(f"response: {response}")
time.sleep(5) # wait 5s for logs to land
import boto3
s3 = boto3.client("s3")
bucket_name = "litellm-logs-2"
# List objects in the bucket
response = s3.list_objects(Bucket=bucket_name)
# Sort the objects based on the LastModified timestamp
objects = sorted(
response["Contents"], key=lambda x: x["LastModified"], reverse=True
)
# Get the keys of the most recent objects
most_recent_keys = [obj["Key"] for obj in objects]
print(most_recent_keys)
# for each key, get the part before "-" as the key. Do it safely
cleaned_keys = []
for key in most_recent_keys:
split_key = key.split("_")
if len(split_key) < 2:
continue
cleaned_keys.append(split_key[1])
print("\n most recent keys", most_recent_keys)
print("\n cleaned keys", cleaned_keys)
print("\n Expected keys: ", expected_keys)
matches = 0
for key in expected_keys:
key += ".json"
assert key in cleaned_keys
if key in cleaned_keys:
matches += 1
# remove the match key
cleaned_keys.remove(key)
# this asserts we log, the first request + the 2nd cached request
print("we had two matches ! passed ", matches)
assert matches == 2
try:
# cleanup s3 bucket in test
for key in most_recent_keys:
s3.delete_object(Bucket=bucket_name, Key=key)
except Exception:
# don't let cleanup fail a test
pass
except Exception as e:
pytest.fail(f"An exception occurred - {e}")
finally:
# post, close log file and verify
# Reset stdout to the original value
print("Passed! Testing async s3 logging")
# test_s3_logging()
@pytest.mark.skip(reason="AWS Suspended Account")
def test_s3_logging_async():
# this tests time added to make s3 logging calls, vs just acompletion calls
try:
litellm.set_verbose = True
# Make 5 calls with an empty success_callback
litellm.success_callback = []
start_time_empty_callback = asyncio.run(make_async_calls())
print("done with no callback test")
print("starting s3 logging load test")
# Make 5 calls with success_callback set to "langfuse"
litellm.success_callback = ["s3"]
litellm.s3_callback_params = {
"s3_bucket_name": "litellm-logs-2",
"s3_aws_secret_access_key": "os.environ/AWS_SECRET_ACCESS_KEY",
"s3_aws_access_key_id": "os.environ/AWS_ACCESS_KEY_ID",
}
start_time_s3 = asyncio.run(make_async_calls())
print("done with s3 test")
# Compare the time for both scenarios
print(f"Time taken with success_callback='s3': {start_time_s3}")
print(f"Time taken with empty success_callback: {start_time_empty_callback}")
# assert the diff is not more than 1 second
assert abs(start_time_s3 - start_time_empty_callback) < 1
except litellm.Timeout as e:
pass
except Exception as e:
pytest.fail(f"An exception occurred - {e}")
async def make_async_calls():
tasks = []
for _ in range(5):
task = asyncio.create_task(
litellm.acompletion(
model="azure/gpt-4.1-mini",
messages=[{"role": "user", "content": "This is a test"}],
max_tokens=5,
temperature=0.7,
timeout=5,
user="langfuse_latency_test_user",
mock_response="It's simple to use and easy to get started",
)
)
tasks.append(task)
# Measure the start time before running the tasks
start_time = asyncio.get_event_loop().time()
# Wait for all tasks to complete
responses = await asyncio.gather(*tasks)
# Print the responses when tasks return
for idx, response in enumerate(responses):
print(f"Response from Task {idx + 1}: {response}")
# Calculate the total time taken
total_time = asyncio.get_event_loop().time() - start_time
return total_time
from litellm.integrations.s3_v2 import S3Logger
class TestS3Logger(S3Logger):
def __init__(self, *args, **kwargs):
self.recorded_requests = {}
self.logged_standard_logging_payload = None
super().__init__(*args, **kwargs)
async def async_log_success_event(self, kwargs, response_obj, start_time, end_time):
self.recorded_requests[response_obj["id"]] = start_time
print("recorded request", self.recorded_requests)
self.logged_standard_logging_payload = kwargs["standard_logging_object"]
return await super().async_log_success_event(
kwargs, response_obj, start_time, end_time
)