litellm/tests/local_testing/test_opik.py
Yuneng Jiang 21092d633b
test(opik): stop the batching test racing its own 5-second flush timer
test_opik_logging_http_request asserted "nothing has been POSTed yet" roughly one
second into a window governed by OpikLogger's 5-second periodic flush. On a loaded
CI worker the five preceding acompletion calls eat that budget, the periodic flush
fires, and the assertion flips. Reproduced with no product changes at all: letting
5.5 seconds pass before the assertion drains the queue and sets mock_post.called,
which is exactly the failure CircleCI reports.

The test now pins flush_interval past anything the test can reach, so the two
batching assertions measure batching instead of wall clock, and drives the flush
path explicitly at the end rather than sleeping the interval. That last phase used
to be near-vacuous, since the size-triggered flush had already emptied the queue.

Assertions now match only calls to Opik's own /traces/batch and /spans/batch.
get_async_httpx_client caches one client per special provider, so the mock is
process-wide and any other logging callback's POST would otherwise count.

Dropped the teardown that closed that shared client, which broke every later test
in the same worker that logs through it, and the try/except that turned assertion
failures into a pytest.fail with no traceback.

Mutation checked: flushing on every event and never flushing on size both fail the
test.
2026-08-27 23:43:38 -07:00

325 lines
11 KiB
Python

import io
import os
import asyncio
import logging
import pytest
import litellm
from litellm._logging import verbose_logger
from unittest.mock import AsyncMock, Mock
verbose_logger.setLevel(logging.DEBUG)
litellm.set_verbose = True
import time
INTERVAL_TOO_LONG_TO_FIRE_DURING_THIS_TEST = 3600
@pytest.mark.asyncio
async def test_opik_logging_http_request():
"""
- Test that HTTP requests are made to Opik
- Traces and spans are batched correctly
"""
from litellm.integrations.opik.opik import OpikLogger
os.environ["OPIK_URL_OVERRIDE"] = "https://fake.comet.com/opik/api"
os.environ["OPIK_API_KEY"] = "anything"
os.environ["OPIK_WORKSPACE"] = "anything"
test_opik_logger = OpikLogger()
test_opik_logger.flush_interval = INTERVAL_TOO_LONG_TO_FIRE_DURING_THIS_TEST
test_opik_logger.batch_size = 12
litellm.callbacks = [test_opik_logger]
mock_post = AsyncMock(return_value=Mock(status_code=202, text="Accepted"))
test_opik_logger.async_httpx_client.post = mock_post
def opik_batch_calls():
return [
call
for call in mock_post.call_args_list
if "/traces/batch" in str(call) or "/spans/batch" in str(call)
]
for _ in range(5):
await litellm.acompletion(
model="gpt-3.5-turbo",
messages=[{"role": "user", "content": "Test message"}],
max_tokens=10,
temperature=0.2,
mock_response="This is a mock response",
)
await asyncio.sleep(1)
assert opik_batch_calls() == [], "events below batch_size must stay queued"
assert len(test_opik_logger.log_queue) == 10
for _ in range(3):
await litellm.acompletion(
model="gpt-3.5-turbo",
messages=[{"role": "user", "content": "Test message"}],
max_tokens=10,
temperature=0.2,
mock_response="This is a mock response",
)
await asyncio.sleep(1)
assert opik_batch_calls(), "crossing batch_size must flush the queue"
events_left_over_after_the_size_triggered_flush = len(test_opik_logger.log_queue)
assert 0 < events_left_over_after_the_size_triggered_flush < test_opik_logger.batch_size
calls_before_periodic_flush = len(opik_batch_calls())
await test_opik_logger.flush_queue()
assert len(opik_batch_calls()) > calls_before_periodic_flush
assert len(test_opik_logger.log_queue) == 0
def test_sync_opik_logging_http_request():
"""
- Test that HTTP requests are made to Opik
- Traces and spans are batched correctly
"""
try:
from litellm.integrations.opik.opik import OpikLogger
os.environ["OPIK_URL_OVERRIDE"] = "https://fake.comet.com/opik/api"
os.environ["OPIK_API_KEY"] = "anything"
os.environ["OPIK_WORKSPACE"] = "anything"
# Initialize OpikLogger
test_opik_logger = OpikLogger()
litellm.callbacks = [test_opik_logger]
litellm.set_verbose = True
# Create a mock for the clients's post method
mock_post = Mock()
mock_post.return_value.status_code = 204
mock_post.return_value.text = "Accepted"
test_opik_logger.sync_httpx_client.post = mock_post
# Make multiple calls to ensure we don't hit the batch size
for _ in range(5):
response = litellm.completion(
model="gpt-3.5-turbo",
messages=[{"role": "user", "content": "Test message"}],
max_tokens=10,
temperature=0.2,
mock_response="This is a mock response",
)
# Need to wait for a short amount of time as the log_success callback is called in a different thread. One or two seconds is often not enough.
time.sleep(3)
# Check that 5 spans and 5 traces were sent
assert (
mock_post.call_count == 10
), f"Expected 10 HTTP requests, but got {mock_post.call_count}"
except Exception as e:
pytest.fail(f"Error occurred: {e}")
@pytest.mark.asyncio
@pytest.mark.skip(reason="local-only test, to test if everything works fine.")
async def test_opik_logging():
try:
from litellm.integrations.opik.opik import OpikLogger
# Initialize OpikLogger
test_opik_logger = OpikLogger()
litellm.callbacks = [test_opik_logger]
litellm.set_verbose = True
# Log a chat completion call
response = await litellm.acompletion(
model="gpt-3.5-turbo",
messages=[{"role": "user", "content": "What LLM are you ?"}],
max_tokens=10,
temperature=0.2,
metadata={"opik": {"custom_field": "custom_value"}},
)
print("Non-streaming response:", response)
# Log a streaming completion call
stream_response = await litellm.acompletion(
model="gpt-3.5-turbo",
messages=[
{"role": "user", "content": "Stream = True - What llm are you ?"}
],
max_tokens=10,
temperature=0.2,
stream=True,
metadata={"opik": {"custom_field": "custom_value"}},
)
print("Streaming response:")
async for chunk in stream_response:
print(chunk.choices[0].delta.content, end="", flush=True)
print() # New line after streaming response
await asyncio.sleep(2)
assert len(test_opik_logger.log_queue) == 4
await asyncio.sleep(test_opik_logger.flush_interval + 1)
assert len(test_opik_logger.log_queue) == 0
except Exception as e:
pytest.fail(f"Error occurred: {e}")
def test_opik_attach_to_existing_trace():
"""
Test attaching spans to existing trace (regression fix for PR #14888)
- When trace_id is provided via current_span_data, only create a span
- Do NOT create a new trace (this was the bug)
- Verify span has correct trace_id and parent_span_id
"""
try:
from litellm.integrations.opik.opik import OpikLogger
os.environ["OPIK_URL_OVERRIDE"] = "https://fake.comet.com/opik/api"
os.environ["OPIK_API_KEY"] = "anything"
os.environ["OPIK_WORKSPACE"] = "anything"
# Initialize OpikLogger
test_opik_logger = OpikLogger()
litellm.callbacks = [test_opik_logger]
litellm.set_verbose = True
# Create a mock for the sync client's post method
mock_post = Mock()
mock_post.return_value.status_code = 204
mock_post.return_value.text = "Accepted"
test_opik_logger.sync_httpx_client.post = mock_post
# Simulate an existing trace and parent span
existing_trace_id = "existing-trace-12345"
existing_parent_span_id = "existing-span-67890"
# Make a completion call with existing trace_id
response = litellm.completion(
model="gpt-3.5-turbo",
messages=[{"role": "user", "content": "Test message"}],
max_tokens=10,
temperature=0.2,
mock_response="This is a mock response",
metadata={
"opik": {
"current_span_data": {
"trace_id": existing_trace_id,
"id": existing_parent_span_id,
},
"tags": ["test-attach-span"],
}
},
)
# Need to wait for a short amount of time as the log_success callback is called in a different thread
time.sleep(3)
# Check the calls made to the mock
calls_made = mock_post.call_args_list
trace_calls = [call for call in calls_made if "/traces/batch" in str(call)]
span_calls = [call for call in calls_made if "/spans/batch" in str(call)]
# With the fix, when trace_id is provided, we should NOT create a new trace
assert (
len(trace_calls) == 0
), f"Expected 0 trace calls when attaching to existing trace, but got {len(trace_calls)}"
assert (
len(span_calls) == 1
), f"Expected exactly 1 span call, but got {len(span_calls)}"
# Verify span has correct trace_id and parent_span_id
span_payload = span_calls[0][1]["json"]["spans"][0]
assert (
span_payload["trace_id"] == existing_trace_id
), f"Expected trace_id to be {existing_trace_id}, but got {span_payload['trace_id']}"
assert (
span_payload["parent_span_id"] == existing_parent_span_id
), f"Expected parent_span_id to be {existing_parent_span_id}, but got {span_payload['parent_span_id']}"
assert (
"test-attach-span" in span_payload["tags"]
), f"Expected 'test-attach-span' tag in {span_payload['tags']}"
except Exception as e:
pytest.fail(f"Error occurred: {e}")
def test_opik_create_new_trace():
"""
Test normal trace creation when no trace_id is provided
- When NO trace_id is provided, create both a new trace and a new span
- Verify the span references the created trace
- Verify tags are included in both trace and span
"""
try:
from litellm.integrations.opik.opik import OpikLogger
os.environ["OPIK_URL_OVERRIDE"] = "https://fake.comet.com/opik/api"
os.environ["OPIK_API_KEY"] = "anything"
os.environ["OPIK_WORKSPACE"] = "anything"
# Initialize OpikLogger
test_opik_logger = OpikLogger()
litellm.callbacks = [test_opik_logger]
litellm.set_verbose = True
# Create a mock for the sync client's post method
mock_post = Mock()
mock_post.return_value.status_code = 204
mock_post.return_value.text = "Accepted"
test_opik_logger.sync_httpx_client.post = mock_post
# Make a completion call WITHOUT providing trace_id
response = litellm.completion(
model="gpt-3.5-turbo",
messages=[{"role": "user", "content": "Test message"}],
max_tokens=10,
temperature=0.2,
mock_response="This is a mock response",
metadata={"opik": {"tags": ["test-new-trace"]}},
)
# Need to wait for a short amount of time as the log_success callback is called in a different thread
time.sleep(3)
# Check the calls made to the mock
calls_made = mock_post.call_args_list
trace_calls = [call for call in calls_made if "/traces/batch" in str(call)]
span_calls = [call for call in calls_made if "/spans/batch" in str(call)]
# Without trace_id provided, we should create both a new trace and a new span
assert (
len(trace_calls) == 1
), f"Expected exactly 1 trace call, but got {len(trace_calls)}"
assert (
len(span_calls) == 1
), f"Expected exactly 1 span call, but got {len(span_calls)}"
# Verify the span references the created trace
trace_payload = trace_calls[0][1]["json"]["traces"][0]
span_payload = span_calls[0][1]["json"]["spans"][0]
assert (
span_payload["trace_id"] == trace_payload["id"]
), "Span should reference the created trace"
# Verify tags are included in both trace and span
assert (
"test-new-trace" in trace_payload["tags"]
), f"Expected 'test-new-trace' tag in trace tags"
assert (
"test-new-trace" in span_payload["tags"]
), f"Expected 'test-new-trace' tag in span tags"
except Exception as e:
pytest.fail(f"Error occurred: {e}")