diff --git a/litellm/integrations/langfuse/langfuse_sdk.py b/litellm/integrations/langfuse/langfuse_sdk.py index cbbcbb5dbd3..0fcc7fdc573 100644 --- a/litellm/integrations/langfuse/langfuse_sdk.py +++ b/litellm/integrations/langfuse/langfuse_sdk.py @@ -22,6 +22,7 @@ import opentelemetry.trace as otel_trace from langfuse import LangfuseOtelSpanAttributes from langfuse.api import LangfuseAPI, Prompt, Prompt_Chat from langfuse.api.core.api_error import ApiError +from langfuse.api.core.request_options import RequestOptions from langfuse.model import BasePromptClient, ChatPromptClient, PromptClient, TextPromptClient from opentelemetry.context import Context from opentelemetry.exporter.otlp.proto.common.trace_encoder import encode_spans @@ -73,6 +74,13 @@ _OBSERVATION_ID_PATTERN: Final = re.compile(r"^(?=.*[1-9a-f])[0-9a-f]{16}$") _TRACER_NAME: Final = "langfuse-sdk" _LANGFUSE_INGESTION_VERSION_HEADER: Final = "x-langfuse-ingestion-version" _LANGFUSE_INGESTION_VERSION: Final = "4" +_NO_REST_RETRIES: Final = RequestOptions(max_retries=0) +_TRUNCATION_MARKER: Final = "" +_TRUNCATION_GROUPS: Final = ( + (LangfuseOtelSpanAttributes.OBSERVATION_INPUT, LangfuseOtelSpanAttributes.TRACE_INPUT), + (LangfuseOtelSpanAttributes.OBSERVATION_OUTPUT, LangfuseOtelSpanAttributes.TRACE_OUTPUT), + (LangfuseOtelSpanAttributes.OBSERVATION_METADATA, LangfuseOtelSpanAttributes.TRACE_METADATA), +) _SERVER_FLOOR_HINT: Final = ( "; the OTLP traces route needs a self-hosted Langfuse server on 3.63.0 or newer " "(https://langfuse.com/self-hosting/upgrade/versioning#sdk-server)" @@ -577,8 +585,51 @@ class _Halving: settled: tuple[SpanExportResult, ...] = () -def _halves(batch: _Batch) -> tuple[_Batch, _Batch]: - return batch[: len(batch) // 2], batch[len(batch) // 2 :] +def _smaller(batch: _Batch) -> tuple[_Batch, ...]: + """What to send after a 413: the two halves of a batch, or a single span with its largest field truncated.""" + match batch: + case (only,): + truncated: Final = _truncated(only) + return () if truncated is None else ((truncated,),) + case _: + return batch[: len(batch) // 2], batch[len(batch) // 2 :] + + +def _in_group(key: str, group: tuple[str, ...]) -> bool: + return any(key == prefix or key.startswith(prefix + ".") for prefix in group) + + +def _group_size(attributes: Mapping[str, AttributeValue], group: tuple[str, ...]) -> int: + return sum( + len(str(value)) for key, value in attributes.items() if _in_group(key, group) and value != _TRUNCATION_MARKER + ) + + +def _truncated(span: ReadableSpan) -> ReadableSpan | None: + """The span with its largest remaining input, output or metadata replaced by the marker the v2 consumer wrote + when an event went over ``LANGFUSE_MAX_EVENT_SIZE_BYTES``, or ``None`` once all three are gone.""" + attributes: Final = span.attributes or MappingProxyType({}) + largest: Final = max(_TRUNCATION_GROUPS, key=lambda group: _group_size(attributes, group)) + if _group_size(attributes, largest) == 0: + return None + kept: Final = {key: value for key, value in attributes.items() if not _in_group(key, largest)} + marked: Final = { + prefix: _TRUNCATION_MARKER for prefix in largest if any(_in_group(key, (prefix,)) for key in attributes) + } + return ReadableSpan( + name=span.name, + context=span.context, + parent=span.parent, + resource=span.resource, + attributes=MappingProxyType({**kept, **marked}), + events=span.events, + links=span.links, + kind=span.kind, + status=span.status, + start_time=span.start_time, + end_time=span.end_time, + instrumentation_scope=span.instrumentation_scope, + ) def enable_langfuse_debug_logging() -> None: @@ -600,7 +651,8 @@ class LangfuseSpanExporter(SpanExporter): as v2's injected httpx client did. A connect or read failure and a retryable status are re-sent after each delay, matching the v2 ingestion consumer; ``BatchSpanProcessor`` would otherwise drop the whole batch on the first exception. A 413 splits the batch in halves until each body fits or a single span - is left, which stands in for the byte ceiling the v2 consumer applied before posting. + is left; that span is re-sent with its input, output and metadata replaced by the v2 consumer's + truncation marker, largest first, and dropped only when the fully truncated span is still refused. """ handler: HTTPHandler @@ -610,9 +662,9 @@ class LangfuseSpanExporter(SpanExporter): delays: Sequence[float] = (1.0, 2.0, 4.0) def export(self, spans: Sequence[ReadableSpan]) -> SpanExportResult: - """Halving a batch of n spans settles every span within ``n.bit_length()`` rounds, so the rounds are a fixed - fold rather than a recursion.""" - rounds: Final = range(len(spans).bit_length() + 1) + """Halving a batch of n spans settles every span within ``n.bit_length()`` rounds plus one per truncation + step, so the rounds are a fixed fold rather than a recursion.""" + rounds: Final = range(len(spans).bit_length() + 1 + len(_TRUNCATION_GROUPS)) final: Final = reduce(lambda halving, _: self._round(halving), rounds, _Halving(pending=(tuple(spans),))) return ( SpanExportResult.SUCCESS @@ -623,7 +675,7 @@ class LangfuseSpanExporter(SpanExporter): def _round(self, halving: _Halving) -> _Halving: sent: Final = tuple((batch, self._send_batch(batch)) for batch in halving.pending) return _Halving( - pending=tuple(half for batch, outcome in sent if outcome == "too_large" for half in _halves(batch)), + pending=tuple(part for batch, outcome in sent if outcome == "too_large" for part in _smaller(batch)), settled=halving.settled + tuple( SpanExportResult.SUCCESS if outcome == "delivered" else SpanExportResult.FAILURE @@ -633,24 +685,38 @@ class LangfuseSpanExporter(SpanExporter): ) def _send_batch(self, batch: _Batch) -> _ExportOutcome: - """A 413 on more than one span asks for halves; on a single span the span is dropped and reported.""" + """A 413 on more than one span asks for halves; on a single span it asks for a truncation, and the span is + dropped and reported once nothing is left to truncate.""" body: Final = _encode(batch) if body is None: return "rejected" outcome: Final = self._send(body) if outcome != "too_large": return outcome - if len(batch) == 1: - verbose_logger.error( - "Langfuse rejected a single %d byte span export to %s as too large, dropping it", - len(body), - self.endpoint, - ) - return "rejected" - verbose_logger.warning( - "Langfuse rejected a %d byte export of %d spans as too large, resending in halves", len(body), len(batch) - ) - return "too_large" + match batch: + case (only,) if _truncated(only) is None: + verbose_logger.error( + "Langfuse rejected a single %d byte span export to %s as too large, dropping it", + len(body), + self.endpoint, + ) + return "rejected" + case (_,): + verbose_logger.warning( + "Langfuse rejected a single %d byte span export to %s as too large, resending it with its " + "largest field replaced by %r", + len(body), + self.endpoint, + _TRUNCATION_MARKER, + ) + return "too_large" + case _: + verbose_logger.warning( + "Langfuse rejected a %d byte export of %d spans as too large, resending in halves", + len(body), + len(batch), + ) + return "too_large" def _send(self, body: bytes) -> _ExportOutcome: for delay in self.delays: @@ -1032,7 +1098,7 @@ class LangfuseApiClient: error or a transport failure is reported as itself rather than as bad credentials. """ try: - projects: Final = self.api.projects.get().data + projects: Final = self.api.projects.get(request_options=_NO_REST_RETRIES).data except ApiError as error: return _auth_check_failure(_api_error_reason(error)) except Exception as error: # noqa: BLE001 # httpx transport errors or a body the response model rejects @@ -1042,7 +1108,7 @@ class LangfuseApiClient: return None def project_id(self) -> str | None: - projects: Final = self.api.projects.get().data + projects: Final = self.api.projects.get(request_options=_NO_REST_RETRIES).data return projects[0].id if projects else None def get_prompt(self, name: str, *, label: str | None = None, version: int | None = None) -> PromptClient: diff --git a/tests/test_litellm/integrations/langfuse/test_langfuse_sdk.py b/tests/test_litellm/integrations/langfuse/test_langfuse_sdk.py index dfe5f284226..1e5f7c07a26 100644 --- a/tests/test_litellm/integrations/langfuse/test_langfuse_sdk.py +++ b/tests/test_litellm/integrations/langfuse/test_langfuse_sdk.py @@ -19,6 +19,8 @@ import httpx import opentelemetry.trace as otel_trace import pytest from langfuse import LangfuseOtelSpanAttributes as A +from langfuse.api.core.api_error import ApiError +from opentelemetry.proto.collector.trace.v1.trace_service_pb2 import ExportTraceServiceRequest from opentelemetry.sdk.trace import SpanProcessor, TracerProvider from opentelemetry.sdk.trace.export import SpanExporter, SpanExportResult from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter @@ -52,6 +54,7 @@ from litellm.integrations.langfuse.langfuse_sdk import ( to_unix_nanos, trace_attributes, ) +from litellm.llms.custom_httpx.http_handler import HTTPHandler CALL_START = datetime(2024, 3, 1, 12, 0, 0, tzinfo=timezone.utc) FIRST_TOKEN = CALL_START + timedelta(seconds=5) @@ -1036,6 +1039,31 @@ def test_auth_check_fails_when_the_keys_reach_no_project(): assert "no project" in failure.reason +@pytest.mark.parametrize("status", [500, 503, 429], ids=["http-500", "http-503", "http-429"]) +def test_auth_check_and_project_id_make_one_round_trip_when_langfuse_is_down(status): + """Both run on the event loop; the generated client's default retries sleep for seconds, or for Retry-After.""" + requests: list[httpx.Request] = [] + + def fail(request: httpx.Request) -> httpx.Response: + requests.append(request) + return httpx.Response(status, request=request, headers={"retry-after": "20"}, json={"message": "down"}) + + client = build_langfuse_client( + public_key="pk", + secret_key="sk", + base_url="http://127.0.0.1:1", + httpx_client=httpx.Client(transport=httpx.MockTransport(fail)), + ) + + started = monotonic() + failure = client.auth_check() + with pytest.raises(ApiError): + client.project_id() + assert failure is not None and f"status_code: {status}" in failure.reason + assert len(requests) == 2 + assert monotonic() - started < 0.5 + + def test_rest_client_reports_the_project_id_and_a_passing_auth_check(): requests: list[httpx.Request] = [] client = build_langfuse_client( @@ -1227,6 +1255,100 @@ def test_exporter_drops_only_the_single_span_that_alone_exceeds_the_cap(monkeypa assert "single" in caplog.text and "too large" in caplog.text +def _decoded_attributes(body: bytes) -> dict[str, str]: + decoded = ExportTraceServiceRequest() + decoded.ParseFromString(body) + return { + attribute.key: attribute.value.string_value + for attribute in decoded.resource_spans[0].scope_spans[0].spans[0].attributes + } + + +def _generation_span(**attributes: str): + provider = TracerProvider() + span = provider.get_tracer("t").start_span("generation", attributes=attributes) + span.end() + return span + + +def test_exporter_truncates_a_single_oversized_span_the_way_v2_did_instead_of_dropping_it(monkeypatch, caplog): + """v2 replaced the largest of input, output and metadata with a marker and still delivered the observation; a + vision request over a self-hosted ingress cap used to lose the whole generation, model and usage included.""" + monkeypatch.setattr("litellm.integrations.langfuse.langfuse_sdk.sleep", lambda _: None) + span = _generation_span( + **{ + "langfuse.observation.input": "data:image/png;base64," + "A" * 6000, + "langfuse.trace.input": "data:image/png;base64," + "A" * 200, + "langfuse.observation.output": "o" * 1000, + "langfuse.observation.metadata.team": "m" * 100, + "langfuse.observation.model.name": "gpt-4o", + } + ) + bodies: list[bytes] = [] + + def transport(request: httpx.Request) -> httpx.Response: + if len(request.content) > 2000: + return httpx.Response(413, request=request) + bodies.append(request.content) + return httpx.Response(200, request=request) + + exporter = LangfuseSpanExporter( + handler=HTTPHandler(client=httpx.Client(transport=httpx.MockTransport(transport))), + endpoint="https://lf.internal.example/api/public/otel/v1/traces", + headers=MappingProxyType({}), + timeout=5.0, + delays=(), + ) + with caplog.at_level(logging.WARNING, logger="LiteLLM"): + result = exporter.export((span,)) + + assert result is SpanExportResult.SUCCESS + delivered = _decoded_attributes(bodies[-1]) + assert delivered["langfuse.observation.input"] == "" + assert delivered["langfuse.trace.input"] == "" + assert delivered["langfuse.observation.output"] == "o" * 1000 + assert delivered["langfuse.observation.metadata.team"] == "m" * 100 + assert delivered["langfuse.observation.model.name"] == "gpt-4o" + assert "dropping it" not in caplog.text and "truncated" in caplog.text + + +def test_exporter_truncates_largest_first_and_drops_only_when_nothing_is_left(monkeypatch, caplog): + monkeypatch.setattr("litellm.integrations.langfuse.langfuse_sdk.sleep", lambda _: None) + span = _generation_span( + **{ + "langfuse.observation.input": "i" * 3000, + "langfuse.observation.output": "o" * 2000, + "langfuse.observation.metadata.a": "m" * 500, + "langfuse.observation.metadata.b": "m" * 500, + } + ) + posted: list[dict[str, str]] = [] + + def always_too_large(request: httpx.Request) -> httpx.Response: + posted.append(_decoded_attributes(request.content)) + return httpx.Response(413, request=request) + + exporter = LangfuseSpanExporter( + handler=HTTPHandler(client=httpx.Client(transport=httpx.MockTransport(always_too_large))), + endpoint="https://lf.internal.example/api/public/otel/v1/traces", + headers=MappingProxyType({}), + timeout=5.0, + delays=(), + ) + with caplog.at_level(logging.ERROR, logger="LiteLLM"): + assert exporter.export((span,)) is SpanExportResult.FAILURE + + marker = "" + assert [sorted(key for key, value in body.items() if value == marker) for body in posted] == [ + [], + ["langfuse.observation.input"], + ["langfuse.observation.input", "langfuse.observation.output"], + ["langfuse.observation.input", "langfuse.observation.metadata", "langfuse.observation.output"], + ] + assert "langfuse.observation.metadata.a" not in posted[-1] + assert "dropping it" in caplog.text + + @pytest.mark.parametrize("status", [400, 401, 403, 404, 422, 499]) def test_exporter_does_not_retry_a_rejected_batch(monkeypatch, status): """Bad credentials or a bad payload will not get better on the next attempt, so retrying only delays the flush."""