mirror of
https://github.com/BerriAI/litellm.git
synced 2026-09-30 01:52:18 +00:00
fix(langfuse): truncate a single oversized span like v2 instead of dropping it, no retries on REST auth and project lookups
Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
parent
f0c38e1de2
commit
da3544653d
2 changed files with 209 additions and 21 deletions
|
|
@ -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 = "<truncated due to size exceeding limit>"
|
||||
_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:
|
||||
|
|
|
|||
|
|
@ -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"] == "<truncated due to size exceeding limit>"
|
||||
assert delivered["langfuse.trace.input"] == "<truncated due to size exceeding limit>"
|
||||
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 = "<truncated due to size exceeding limit>"
|
||||
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."""
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue