diff --git a/docs/my-website/docs/observability/opik_integration.md b/docs/my-website/docs/observability/opik_integration.md index b4bcef53937..1ba1c2de210 100644 --- a/docs/my-website/docs/observability/opik_integration.md +++ b/docs/my-website/docs/observability/opik_integration.md @@ -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?" + } + ] +}' +``` diff --git a/litellm/integrations/opik/opik.py b/litellm/integrations/opik/opik.py index 9f90d2384d8..9fa3482f663 100644 --- a/litellm/integrations/opik/opik.py +++ b/litellm/integrations/opik/opik.py @@ -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, } ) diff --git a/tests/local_testing/test_opik.py b/tests/local_testing/test_opik.py index 62c17fbb2d5..7e3fa84df1a 100644 --- a/tests/local_testing/test_opik.py +++ b/tests/local_testing/test_opik.py @@ -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}"