Merge pull request #14888 from mrFranklin/feat/improve-opik

feat: improve opik integration code
This commit is contained in:
Krish Dholakia 2025-09-25 23:40:41 -07:00 • committed by GitHub
commit 79ebb2c95e
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
3 changed files with 81 additions and 32 deletions

View file

@ -140,6 +140,7 @@ These can be passed inside metadata with the `opik` key.
- `project_name` - Name of the Opik project to send data to.
- `current_span_data` - The current span data to be used for tracing.
- `tags` - Tags to be used for tracing.
- `thread_id` - The thread id to group together multiple related traces.
### Usage
@ -159,8 +160,10 @@ response = litellm.completion(
messages=messages,
metadata = {
"opik": {
"project_name": "your-opik-project-name",
"current_span_data": get_current_span_data(),
"tags": ["streaming-test"],
"thread_id": "your-thread-id"
},
}
)
@ -174,7 +177,7 @@ curl -L -X POST 'http://0.0.0.0:4000/v1/chat/completions' \
-H 'Content-Type: application/json' \
-H 'Authorization: Bearer sk-1234' \
-d '{
"model": "gpt-3.5-turbo-testing",
"model": "gpt-3.5-turbo",
"messages": [
{
"role": "user",
@ -183,8 +186,10 @@ curl -L -X POST 'http://0.0.0.0:4000/v1/chat/completions' \
],
"metadata": {
"opik": {
"project_name": "your-opik-project-name",
"current_span_data": "...",
"tags": ["streaming-test"],
"thread_id": "your-thread-id"
},
}
}'
@ -195,12 +200,25 @@ curl -L -X POST 'http://0.0.0.0:4000/v1/chat/completions' \
You can also pass the fields as part of the request header with a `opik_*` prefix:
```shell
curl --location --request POST 'http://0.0.0.0:4000/chat/completions' \
--header 'Content-Type: application/json' \
--header 'Authorization: Bearer sk-1234' \
--header 'opik_project_name: your-opik-project-name' \
--header 'opik_thread_id: your-thread-id' \
--header 'opik_tags: ["streaming-test"]' \
--data '{
"model": "gpt-3.5-turbo",
"messages": [
{
"role": "user",
"content": "What's the weather like in Boston today?"
}
]
}'
```

View file

@ -192,9 +192,25 @@ class OpikLogger(CustomBatchLogger):
# Extract opik metadata
litellm_opik_metadata = litellm_params_metadata.get("opik", {})
# Use standard_logging_object to create metadata and input/output data
standard_logging_object = kwargs.get("standard_logging_object", None)
if standard_logging_object is None:
verbose_logger.debug(
"OpikLogger skipping event; no standard_logging_object found"
)
return []
# Update litellm_opik_metadata with opik metadata from requester
standard_logging_metadata = standard_logging_object.get("metadata", {}) or {}
requester_metadata = standard_logging_metadata.get("requester_metadata", {}) or {}
requester_opik_metadata = requester_metadata.get("opik", {}) or {}
litellm_opik_metadata.update(requester_opik_metadata)
verbose_logger.debug(
f"litellm_opik_metadata - {json.dumps(litellm_opik_metadata, default=str)}"
)
project_name = litellm_opik_metadata.get("project_name", self.opik_project_name)
# Extract trace_id and parent_span_id
@ -208,19 +224,33 @@ class OpikLogger(CustomBatchLogger):
else:
trace_id = None
parent_span_id = None
# Create Opik tags
opik_tags = litellm_opik_metadata.get("tags", [])
if kwargs.get("custom_llm_provider"):
opik_tags.append(kwargs["custom_llm_provider"])
# Get thread_id if present
thread_id = litellm_opik_metadata.get("thread_id", None)
# Use standard_logging_object to create metadata and input/output data
standard_logging_object = kwargs.get("standard_logging_object", None)
if standard_logging_object is None:
verbose_logger.debug(
"OpikLogger skipping event; no standard_logging_object found"
)
return []
# Override with any opik_ headers from proxy request
proxy_server_request = _litellm_params.get("proxy_server_request", {}) or {}
proxy_headers = proxy_server_request.get("headers", {}) or {}
for key, value in proxy_headers.items():
if key.startswith("opik_"):
param_key = key.replace("opik_", "", 1)
if param_key == "project_name" and value:
project_name = value
elif param_key == "thread_id" and value:
thread_id = value
elif param_key == "tags" and value:
try:
parsed_tags = json.loads(value)
if isinstance(parsed_tags, list):
opik_tags.extend(parsed_tags)
except (json.JSONDecodeError, TypeError):
pass
# Create input and output data
input_data = standard_logging_object.get("messages", {})
output_data = standard_logging_object.get("response", {})
@ -243,7 +273,7 @@ class OpikLogger(CustomBatchLogger):
del metadata["current_span_data"]
metadata["created_from"] = "litellm"
metadata.update(standard_logging_object.get("metadata", {}))
metadata.update(standard_logging_metadata)
if "call_type" in standard_logging_object:
metadata["type"] = standard_logging_object["call_type"]
if "status" in standard_logging_object:
@ -286,20 +316,20 @@ class OpikLogger(CustomBatchLogger):
verbose_logger.debug(
f"OpikLogger creating payload for trace with id {trace_id}"
)
payload.append(
{
"project_name": project_name,
"id": trace_id,
"name": trace_name,
"start_time": start_time.astimezone(timezone.utc).isoformat().replace("+00:00", "Z"),
"end_time": end_time.astimezone(timezone.utc).isoformat().replace("+00:00", "Z"),
"input": input_data,
"output": output_data,
"metadata": metadata,
"tags": opik_tags,
}
)
payload.append(
{
"project_name": project_name,
"id": trace_id,
"name": trace_name,
"start_time": start_time.astimezone(timezone.utc).isoformat().replace("+00:00", "Z"),
"end_time": end_time.astimezone(timezone.utc).isoformat().replace("+00:00", "Z"),
"input": input_data,
"output": output_data,
"metadata": metadata,
"tags": opik_tags,
"thread_id": thread_id,
}
)
span_id = create_uuid7()
verbose_logger.debug(
@ -319,6 +349,7 @@ class OpikLogger(CustomBatchLogger):
"output": output_data,
"metadata": metadata,
"tags": opik_tags,
"thread_id": thread_id,
"usage": usage,
}
)

View file

@ -120,9 +120,9 @@ def test_sync_opik_logging_http_request():
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
time.sleep(1)
# 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}"