diff --git a/litellm/integrations/langfuse/langfuse.py b/litellm/integrations/langfuse/langfuse.py index 182c8863760..e5142dc692a 100644 --- a/litellm/integrations/langfuse/langfuse.py +++ b/litellm/integrations/langfuse/langfuse.py @@ -1,10 +1,10 @@ #### What this does #### # On success, logs events to Langfuse import copy -import inspect import os import traceback -from typing import TYPE_CHECKING, Any, Dict, Optional +from typing import TYPE_CHECKING, Any +from collections.abc import MutableSequence, MutableSet, MutableMapping from packaging.version import Version from pydantic import BaseModel @@ -14,7 +14,7 @@ from litellm._logging import verbose_logger from litellm.litellm_core_utils.redact_messages import redact_user_api_key_info from litellm.secret_managers.main import str_to_bool from litellm.types.integrations.langfuse import * -from litellm.types.utils import StandardCallbackDynamicParams, StandardLoggingPayload +from litellm.types.utils import StandardLoggingPayload if TYPE_CHECKING: from litellm.litellm_core_utils.litellm_logging import DynamicLoggingCache @@ -355,6 +355,37 @@ class LangFuseLogger: ) ) + def _prepare_metadata(self, metadata) -> Any: + try: + return copy.deepcopy(metadata) # Avoid modifying the original metadata + except (TypeError, copy.Error) as e: + verbose_logger.warning(f"Langfuse Layer Error - {e}") + + new_metadata = {} + + # if metadata is not a MutableMapping, return an empty dict since we can't call items() on it + if not isinstance(metadata, MutableMapping): + verbose_logger.warning("Langfuse Layer Logging - metadata is not a MutableMapping, returning empty dict") + return new_metadata + + for key, value in metadata.items(): + try: + if isinstance(value, MutableMapping): + new_metadata[key] = self._prepare_metadata(value) + elif isinstance(value, (MutableSequence, MutableSet)): + new_metadata[key] = type(value)( + self._prepare_metadata(v) if isinstance(v, MutableMapping) else copy.deepcopy(v) + for v in value + ) + elif isinstance(value, BaseModel): + new_metadata[key] = value.model_dump() + else: + new_metadata[key] = copy.deepcopy(value) + except (TypeError, copy.Error): + verbose_logger.warning(f"Langfuse Layer Error - Couldn't copy metadata key: {key} - {traceback.format_exc()}") + + return new_metadata + def _log_langfuse_v2( # noqa: PLR0915 self, user_id, @@ -373,40 +404,19 @@ class LangFuseLogger: ) -> tuple: import langfuse + print_verbose("Langfuse Layer Logging - logging to langfuse v2") + try: - tags = [] - try: - optional_params.pop("metadata") - metadata = copy.deepcopy( - metadata - ) # Avoid modifying the original metadata - except Exception: - new_metadata = {} - for key, value in metadata.items(): - if ( - isinstance(value, list) - or isinstance(value, dict) - or isinstance(value, str) - or isinstance(value, int) - or isinstance(value, float) - ): - new_metadata[key] = copy.deepcopy(value) - elif isinstance(value, BaseModel): - new_metadata[key] = value.model_dump() - metadata = new_metadata + metadata = self._prepare_metadata(metadata) - supports_tags = Version(langfuse.version.__version__) >= Version("2.6.3") - supports_prompt = Version(langfuse.version.__version__) >= Version("2.7.3") - supports_costs = Version(langfuse.version.__version__) >= Version("2.7.3") - supports_completion_start_time = Version( - langfuse.version.__version__ - ) >= Version("2.7.3") + langfuse_version = Version(langfuse.version.__version__) - print_verbose("Langfuse Layer Logging - logging to langfuse v2 ") + supports_tags = langfuse_version >= Version("2.6.3") + supports_prompt = langfuse_version >= Version("2.7.3") + supports_costs = langfuse_version >= Version("2.7.3") + supports_completion_start_time = langfuse_version >= Version("2.7.3") - if supports_tags: - metadata_tags = metadata.pop("tags", []) - tags = metadata_tags + tags = metadata.pop("tags", []) if supports_tags else [] # Clean Metadata before logging - never log raw metadata # the raw metadata can contain circular references which leads to infinite recursion diff --git a/tests/logging_callback_tests/test_langfuse_unit_tests.py b/tests/logging_callback_tests/test_langfuse_unit_tests.py index 2a6cbe00a90..20b33f81b52 100644 --- a/tests/logging_callback_tests/test_langfuse_unit_tests.py +++ b/tests/logging_callback_tests/test_langfuse_unit_tests.py @@ -1,19 +1,13 @@ -import json import os import sys +import threading from datetime import datetime -from pydantic.main import Model - sys.path.insert( 0, os.path.abspath("../..") ) # Adds the parent directory to the system-path import pytest -import litellm -import asyncio -import logging -from litellm._logging import verbose_logger from litellm.integrations.langfuse.langfuse import ( LangFuseLogger, ) @@ -217,3 +211,27 @@ def test_get_langfuse_logger_for_request_with_cached_logger(): assert result == cached_logger mock_cache.get_cache.assert_called_once() + +@pytest.mark.parametrize("metadata", [ + {'a': 1, 'b': 2, 'c': 3}, + {'a': {'nested_a': 1}, 'b': {'nested_b': 2}}, + {'a': [1, 2, 3], 'b': {4, 5, 6}}, + {'a': (1, 2), 'b': frozenset([3, 4]), 'c': {'d': [5, 6]}}, + {'lock': threading.Lock()}, + {'func': lambda x: x + 1}, + { + 'int': 42, + 'str': 'hello', + 'list': [1, 2, 3], + 'set': {4, 5}, + 'dict': {'nested': 'value'}, + 'non_copyable': threading.Lock(), + 'function': print + }, + ['list', 'not', 'a', 'dict'], + {'timestamp': datetime.now()}, + {}, + None, +]) +def test_langfuse_logger_prepare_metadata(metadata): + global_langfuse_logger._prepare_metadata(metadata)