mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-01 02:02:20 +00:00
fix(langfuse): name the auth check failure, split a 413 export, wire LANGFUSE_DEBUG, stamp error output under a parent, send the ingestion version header
Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
parent
c662597f0d
commit
6b41d4ece1
6 changed files with 283 additions and 44 deletions
|
|
@ -302,7 +302,7 @@ class LangFuseLogger:
|
|||
) from e
|
||||
raise_if_unsupported_langfuse_version(self.langfuse_sdk_version)
|
||||
raise_if_unusable_prompt_cache_ttl()
|
||||
from litellm.integrations.langfuse.langfuse_sdk import configured_release
|
||||
from litellm.integrations.langfuse.langfuse_sdk import configured_release, enable_langfuse_debug_logging
|
||||
|
||||
self.public_key, self.secret_key, self.langfuse_host = resolve_langfuse_credentials(
|
||||
langfuse_public_key=langfuse_public_key,
|
||||
|
|
@ -318,6 +318,8 @@ class LangFuseLogger:
|
|||
self.langfuse_environment = self.resolve_deployment_environment()
|
||||
self.langfuse_release = configured_release()
|
||||
self.langfuse_debug = parse_langfuse_debug(os.getenv("LANGFUSE_DEBUG"))
|
||||
if self.langfuse_debug:
|
||||
enable_langfuse_debug_logging()
|
||||
self.langfuse_flush_interval = LangFuseLogger._get_langfuse_flush_interval(flush_interval)
|
||||
|
||||
if should_use_langfuse_mock():
|
||||
|
|
@ -751,10 +753,7 @@ class LangFuseLogger:
|
|||
for key in list(filter(lambda key: key.startswith("trace_"), clean_metadata.keys())):
|
||||
trace_params[key.replace("trace_", "")] = clean_metadata.pop(key, None)
|
||||
|
||||
if level == "ERROR":
|
||||
trace_params["status_message"] = masked_output
|
||||
else:
|
||||
trace_params["output"] = masked_output if not mask_output else "redacted-by-litellm"
|
||||
trace_params["output"] = masked_output if not mask_output else "redacted-by-litellm"
|
||||
|
||||
if debug is True or (isinstance(debug, str) and debug.lower() == "true"):
|
||||
debug_metadata: Final = {
|
||||
|
|
|
|||
|
|
@ -1,5 +1,6 @@
|
|||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import os
|
||||
import re
|
||||
import threading
|
||||
|
|
@ -37,6 +38,7 @@ from litellm.litellm_core_utils.safe_json_dumps import safe_dumps
|
|||
from litellm.llms.custom_httpx.http_handler import HTTPHandler, _get_httpx_client
|
||||
|
||||
__all__ = (
|
||||
"AuthCheckFailure",
|
||||
"DiscardingSpanExporter",
|
||||
"LangfuseApiClient",
|
||||
"LangfuseObservation",
|
||||
|
|
@ -51,6 +53,7 @@ __all__ = (
|
|||
"configured_release",
|
||||
"configured_sample_rate",
|
||||
"configured_timeout",
|
||||
"enable_langfuse_debug_logging",
|
||||
"flush_langfuse_tracing",
|
||||
"observation_attributes",
|
||||
"release_langfuse_tracing",
|
||||
|
|
@ -65,6 +68,13 @@ __all__ = (
|
|||
_TRACE_ID_PATTERN: Final = re.compile(r"^(?=.*[1-9a-f])[0-9a-f]{32}$")
|
||||
_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"
|
||||
_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)"
|
||||
)
|
||||
_langfuse_logger: Final = logging.getLogger("langfuse")
|
||||
_MAX_QUEUE_SIZE: Final = 100_000
|
||||
_DEFAULT_FLUSH_AT: Final = 512
|
||||
_CHANNEL_RETIRE_GRACE_SECONDS: Final = 60.0
|
||||
|
|
@ -541,7 +551,13 @@ class DiscardingSpanExporter(SpanExporter):
|
|||
return True
|
||||
|
||||
|
||||
_ExportOutcome = Literal["delivered", "retry", "rejected"]
|
||||
_ExportOutcome = Literal["delivered", "retry", "rejected", "too_large"]
|
||||
|
||||
|
||||
def enable_langfuse_debug_logging() -> None:
|
||||
"""What ``Langfuse(debug=True)`` does: a root handler if none exists, and the ``langfuse`` logger at DEBUG."""
|
||||
logging.basicConfig(format="%(asctime)s - %(name)s - %(levelname)s - %(message)s")
|
||||
_langfuse_logger.setLevel(logging.DEBUG)
|
||||
|
||||
|
||||
def _retryable_status(status: int) -> bool:
|
||||
|
|
@ -556,7 +572,8 @@ class LangfuseSpanExporter(SpanExporter):
|
|||
The handler carries litellm's TLS material (``ssl_verify``, CA bundle, client certificate) exactly
|
||||
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.
|
||||
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.
|
||||
"""
|
||||
|
||||
handler: HTTPHandler
|
||||
|
|
@ -569,19 +586,38 @@ class LangfuseSpanExporter(SpanExporter):
|
|||
body: Final = _encode(spans)
|
||||
if body is None:
|
||||
return SpanExportResult.FAILURE
|
||||
return self._send(body)
|
||||
outcome: Final = self._send(body)
|
||||
if outcome != "too_large":
|
||||
return SpanExportResult.SUCCESS if outcome == "delivered" else SpanExportResult.FAILURE
|
||||
if len(spans) == 1:
|
||||
verbose_logger.error(
|
||||
"Langfuse rejected a single %d byte span export to %s as too large, dropping it",
|
||||
len(body),
|
||||
self.endpoint,
|
||||
)
|
||||
return SpanExportResult.FAILURE
|
||||
verbose_logger.warning(
|
||||
"Langfuse rejected a %d byte export of %d spans as too large, resending in halves", len(body), len(spans)
|
||||
)
|
||||
half: Final = len(spans) // 2
|
||||
results: Final = (self.export(spans[:half]), self.export(spans[half:]))
|
||||
return (
|
||||
SpanExportResult.SUCCESS
|
||||
if all(r is SpanExportResult.SUCCESS for r in results)
|
||||
else SpanExportResult.FAILURE
|
||||
)
|
||||
|
||||
def _send(self, body: bytes) -> SpanExportResult:
|
||||
def _send(self, body: bytes) -> _ExportOutcome:
|
||||
for delay in self.delays:
|
||||
outcome: _ExportOutcome = self._post(body)
|
||||
if outcome != "retry":
|
||||
return SpanExportResult.SUCCESS if outcome == "delivered" else SpanExportResult.FAILURE
|
||||
return outcome
|
||||
verbose_logger.warning("Langfuse export to %s failed, retrying in %ss", self.endpoint, delay)
|
||||
sleep(delay)
|
||||
last: Final = self._post(body)
|
||||
if last == "retry":
|
||||
verbose_logger.error("Langfuse export to %s failed after %d retries", self.endpoint, len(self.delays))
|
||||
return SpanExportResult.SUCCESS if last == "delivered" else SpanExportResult.FAILURE
|
||||
return last
|
||||
|
||||
def _post(self, body: bytes) -> _ExportOutcome:
|
||||
try:
|
||||
|
|
@ -590,11 +626,19 @@ class LangfuseSpanExporter(SpanExporter):
|
|||
status: Final = error.response.status_code
|
||||
if _retryable_status(status):
|
||||
return "retry"
|
||||
verbose_logger.error("Langfuse rejected an export to %s with HTTP %d", self.endpoint, status)
|
||||
if status == 413:
|
||||
return "too_large"
|
||||
verbose_logger.error(
|
||||
"Langfuse rejected an export to %s with HTTP %d%s",
|
||||
self.endpoint,
|
||||
status,
|
||||
_SERVER_FLOOR_HINT if status == 404 else "",
|
||||
)
|
||||
return "rejected"
|
||||
except (httpx.TransportError, litellm.Timeout) as error:
|
||||
verbose_logger.warning("Langfuse export to %s raised %s", self.endpoint, error)
|
||||
return "retry"
|
||||
_langfuse_logger.debug("Exported %d bytes of spans to %s", len(body), self.endpoint)
|
||||
return "delivered"
|
||||
|
||||
def shutdown(self) -> None:
|
||||
|
|
@ -623,7 +667,9 @@ def _encodes(span: ReadableSpan) -> bool:
|
|||
|
||||
|
||||
def _build_span_exporter(*, public_key: str, secret_key: str, base_url: str) -> LangfuseSpanExporter:
|
||||
"""Endpoint, headers, timeout and retries mirror the v2 SDK's so the server treats the spans as SDK traffic."""
|
||||
"""Endpoint, headers and export path are the v4 SDK span processor's, so the server treats the spans as SDK
|
||||
traffic; the 20 s timeout and the retry count are what the v2 consumer used. The ingestion-version header is
|
||||
the one Langfuse's compatibility matrix asks a v4 producer to send."""
|
||||
export_path: Final = os.getenv("LANGFUSE_OTEL_TRACES_EXPORT_PATH") or "/api/public/otel/v1/traces"
|
||||
encoded_auth: Final = b64encode(f"{public_key}:{secret_key}".encode()).decode("ascii")
|
||||
return LangfuseSpanExporter(
|
||||
|
|
@ -636,6 +682,7 @@ def _build_span_exporter(*, public_key: str, secret_key: str, base_url: str) ->
|
|||
"x-langfuse-sdk-name": "python",
|
||||
"x-langfuse-sdk-version": version("langfuse"),
|
||||
"x-langfuse-public-key": public_key,
|
||||
_LANGFUSE_INGESTION_VERSION_HEADER: _LANGFUSE_INGESTION_VERSION,
|
||||
}
|
||||
),
|
||||
timeout=configured_timeout(),
|
||||
|
|
@ -885,6 +932,16 @@ def _prompt_client(prompt: Prompt) -> PromptClient:
|
|||
return ChatPromptClient(prompt) if isinstance(prompt, Prompt_Chat) else TextPromptClient(prompt)
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class AuthCheckFailure:
|
||||
reason: str
|
||||
|
||||
|
||||
def _auth_check_failure(reason: str) -> AuthCheckFailure:
|
||||
verbose_logger.warning("Langfuse auth check failed: %s", reason)
|
||||
return AuthCheckFailure(reason)
|
||||
|
||||
|
||||
class LangfuseApiClient:
|
||||
"""litellm's handle on one Langfuse project over its REST API: prompts, ``auth_check`` and the project id.
|
||||
|
||||
|
|
@ -908,12 +965,19 @@ class LangfuseApiClient:
|
|||
self._refreshing: Final[set[_PromptKey]] = set()
|
||||
self._lock: Final = threading.Lock()
|
||||
|
||||
def auth_check(self) -> bool:
|
||||
def auth_check(self) -> AuthCheckFailure | None:
|
||||
"""``None`` when the keys reach a project; otherwise the reason, which is also logged.
|
||||
|
||||
Mirrors the SDK's ``Langfuse.auth_check``: a 200 with no project is a failure too, and a server
|
||||
error or a transport failure is reported as itself rather than as bad credentials.
|
||||
"""
|
||||
try:
|
||||
self.api.projects.get()
|
||||
except Exception:
|
||||
return False
|
||||
return True
|
||||
projects: Final = self.api.projects.get().data
|
||||
except Exception as error: # noqa: BLE001 # ApiError, httpx transport errors or a body the response model rejects
|
||||
return _auth_check_failure(str(error) or type(error).__name__)
|
||||
if not projects:
|
||||
return _auth_check_failure("no project found for the keys provided")
|
||||
return None
|
||||
|
||||
def project_id(self) -> str | None:
|
||||
projects: Final = self.api.projects.get().data
|
||||
|
|
@ -963,7 +1027,7 @@ def build_langfuse_client(
|
|||
"""The REST client for prompt management, ``auth_check`` and the Slack project link.
|
||||
|
||||
Missing keys are passed through as absent credentials: the server answers 401, which
|
||||
``auth_check`` reports as ``False`` rather than raising at construction.
|
||||
``auth_check`` reports as a failure rather than raising at construction.
|
||||
"""
|
||||
return LangfuseApiClient(
|
||||
LangfuseAPI(
|
||||
|
|
|
|||
|
|
@ -394,10 +394,9 @@ async def health_services_endpoint(
|
|||
from litellm.integrations.langfuse.langfuse import LangFuseLogger
|
||||
|
||||
langfuse_logger: Final = LangFuseLogger()
|
||||
if langfuse_logger.api_client.auth_check() is False:
|
||||
raise ValueError(
|
||||
"langfuse auth_check failed - verify LANGFUSE_PUBLIC_KEY and LANGFUSE_SECRET_KEY are set correctly"
|
||||
)
|
||||
auth_failure: Final = langfuse_logger.api_client.auth_check()
|
||||
if auth_failure is not None:
|
||||
raise ValueError(f"langfuse auth_check failed: {auth_failure.reason}")
|
||||
_ = litellm.completion(
|
||||
model="openai/litellm-mock-response-model",
|
||||
messages=[{"role": "user", "content": "Hey, how's it going?"}],
|
||||
|
|
|
|||
|
|
@ -34,12 +34,14 @@ from litellm.integrations.langfuse.langfuse_sdk import (
|
|||
LangfuseSpanExporter,
|
||||
LangfuseTracing,
|
||||
_build_span_exporter,
|
||||
_encode,
|
||||
acquire_langfuse_tracing,
|
||||
build_langfuse_client,
|
||||
build_langfuse_tracing,
|
||||
configured_flush_at,
|
||||
configured_prompt_cache_ttl,
|
||||
configured_sample_rate,
|
||||
enable_langfuse_debug_logging,
|
||||
flush_langfuse_tracing,
|
||||
observation_attributes,
|
||||
release_langfuse_tracing,
|
||||
|
|
@ -978,7 +980,7 @@ def test_rest_client_authenticates_with_the_credentials_it_was_built_with():
|
|||
httpx_client=_recording_transport(requests),
|
||||
)
|
||||
|
||||
assert rotated.auth_check() is False
|
||||
assert rotated.auth_check() is not None
|
||||
assert requests[-1].url.host == "127.0.0.1" and requests[-1].url.port == 2
|
||||
assert requests[-1].headers["authorization"] == "Basic " + b64encode(b"pk-rest-test:sk-second").decode()
|
||||
|
||||
|
|
@ -988,7 +990,50 @@ def test_rest_client_without_keys_fails_auth_check_instead_of_raising(monkeypatc
|
|||
for name in ("LANGFUSE_PUBLIC_KEY", "LANGFUSE_SECRET_KEY"):
|
||||
monkeypatch.delenv(name, raising=False)
|
||||
client = build_langfuse_client(public_key=None, secret_key=None, base_url="http://127.0.0.1:1", httpx_client=None)
|
||||
assert client.auth_check() is False
|
||||
assert client.auth_check() is not None
|
||||
|
||||
|
||||
def test_auth_check_names_the_servers_rejection(caplog):
|
||||
"""``/health/services`` used to print the 401 verbatim; a generic credentials message hides a 403 or a 500."""
|
||||
client = build_langfuse_client(
|
||||
public_key="pk",
|
||||
secret_key="sk",
|
||||
base_url="http://127.0.0.1:1",
|
||||
httpx_client=_recording_transport([], status=401),
|
||||
)
|
||||
with caplog.at_level(logging.WARNING, logger="LiteLLM"):
|
||||
failure = client.auth_check()
|
||||
assert failure is not None
|
||||
assert "status_code: 401" in failure.reason and "unauthorized" in failure.reason
|
||||
assert failure.reason in caplog.text
|
||||
|
||||
|
||||
def test_auth_check_names_an_unreachable_destination_rather_than_the_keys():
|
||||
def refuse(request: httpx.Request) -> httpx.Response:
|
||||
raise httpx.ConnectError("connection refused by lf.internal.example", request=request)
|
||||
|
||||
client = build_langfuse_client(
|
||||
public_key="pk",
|
||||
secret_key="sk",
|
||||
base_url="http://lf.internal.example",
|
||||
httpx_client=httpx.Client(transport=httpx.MockTransport(refuse)),
|
||||
)
|
||||
failure = client.auth_check()
|
||||
assert failure is not None
|
||||
assert "connection refused by lf.internal.example" in failure.reason
|
||||
|
||||
|
||||
def test_auth_check_fails_when_the_keys_reach_no_project():
|
||||
"""A 200 with an empty project list is what the SDK's own ``auth_check`` raises on; it is not a pass."""
|
||||
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(lambda _: httpx.Response(200, json={"data": []}))),
|
||||
)
|
||||
failure = client.auth_check()
|
||||
assert failure is not None
|
||||
assert "no project" in failure.reason
|
||||
|
||||
|
||||
def test_rest_client_reports_the_project_id_and_a_passing_auth_check():
|
||||
|
|
@ -1000,7 +1045,7 @@ def test_rest_client_reports_the_project_id_and_a_passing_auth_check():
|
|||
httpx_client=_recording_transport(requests, status=200),
|
||||
)
|
||||
assert client.project_id() == "proj-under-test"
|
||||
assert client.auth_check() is True
|
||||
assert client.auth_check() is None
|
||||
|
||||
|
||||
def test_rest_client_leaves_a_host_applications_langfuse_client_alone():
|
||||
|
|
@ -1127,7 +1172,62 @@ def test_exporter_gives_up_after_the_last_delay(monkeypatch):
|
|||
assert slept == [1.0, 2.0]
|
||||
|
||||
|
||||
@pytest.mark.parametrize("status", [400, 401, 403, 404, 413, 422, 499])
|
||||
def _exporter_with_body_cap(max_bytes: int, *, deliveries: list[int]):
|
||||
"""A destination that answers 413 to any body over ``max_bytes``, the way an ingress with a body limit does."""
|
||||
from litellm.llms.custom_httpx.http_handler import HTTPHandler
|
||||
|
||||
def transport(request: httpx.Request) -> httpx.Response:
|
||||
if len(request.content) > max_bytes:
|
||||
return httpx.Response(413, request=request)
|
||||
deliveries.append(len(request.content))
|
||||
return httpx.Response(200, request=request)
|
||||
|
||||
return 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=(),
|
||||
)
|
||||
|
||||
|
||||
def test_exporter_splits_a_batch_the_destination_finds_too_large(monkeypatch):
|
||||
"""One 413 used to drop every span in the batch; the v2 consumer sized its batches by bytes before posting."""
|
||||
monkeypatch.setattr("litellm.integrations.langfuse.langfuse_sdk.sleep", lambda _: None)
|
||||
spans = tuple(_finished_span() for _ in range(8))
|
||||
whole = _encode(spans)
|
||||
assert whole is not None
|
||||
deliveries: list[int] = []
|
||||
exporter = _exporter_with_body_cap(len(whole) // 2, deliveries=deliveries)
|
||||
|
||||
assert exporter.export(spans) is SpanExportResult.SUCCESS
|
||||
assert len(deliveries) >= 2
|
||||
assert all(size <= len(whole) // 2 for size in deliveries)
|
||||
assert (
|
||||
sum(deliveries) >= len(whole) - 8 * 8
|
||||
) # each half repeats the resource and scope envelope, spans are not lost
|
||||
|
||||
|
||||
def test_exporter_drops_only_the_single_span_that_alone_exceeds_the_cap(monkeypatch, caplog):
|
||||
monkeypatch.setattr("litellm.integrations.langfuse.langfuse_sdk.sleep", lambda _: None)
|
||||
provider = TracerProvider()
|
||||
huge = provider.get_tracer("t").start_span("generation", attributes={"body": "x" * 4000})
|
||||
huge.end()
|
||||
small = tuple(_finished_span() for _ in range(3))
|
||||
single_small = _encode(small[:1])
|
||||
assert single_small is not None
|
||||
deliveries: list[int] = []
|
||||
exporter = _exporter_with_body_cap(len(single_small) * 3, deliveries=deliveries)
|
||||
|
||||
with caplog.at_level(logging.ERROR, logger="LiteLLM"):
|
||||
result = exporter.export((*small, huge))
|
||||
|
||||
assert result is SpanExportResult.FAILURE
|
||||
assert len(deliveries) >= 1 and all(size <= len(single_small) * 3 for size in deliveries)
|
||||
assert "single" in caplog.text and "too large" 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."""
|
||||
slept = []
|
||||
|
|
@ -1139,6 +1239,17 @@ def test_exporter_does_not_retry_a_rejected_batch(monkeypatch, status):
|
|||
assert slept == []
|
||||
|
||||
|
||||
def test_exporter_names_the_server_floor_when_the_otlp_route_is_missing(monkeypatch, caplog):
|
||||
"""A Langfuse server too old to serve the OTLP route answers 404; a bare status leaves the operator guessing."""
|
||||
monkeypatch.setattr("litellm.integrations.langfuse.langfuse_sdk.sleep", lambda _: None)
|
||||
exporter, _ = _exporter_over([404])
|
||||
|
||||
with caplog.at_level(logging.ERROR, logger="LiteLLM"):
|
||||
assert exporter.export((_finished_span(),)) is SpanExportResult.FAILURE
|
||||
|
||||
assert "HTTP 404" in caplog.text and "3.63.0" in caplog.text
|
||||
|
||||
|
||||
def _finished_span_named(name: object):
|
||||
provider = TracerProvider()
|
||||
span = provider.get_tracer("t").start_span("placeholder")
|
||||
|
|
@ -1192,6 +1303,25 @@ def test_built_exporter_uses_the_shared_litellm_handler_and_langfuse_headers(mon
|
|||
assert exporter.headers["Authorization"] == "Basic " + b64encode(b"pk:sk").decode()
|
||||
assert exporter.headers["x-langfuse-public-key"] == "pk"
|
||||
assert exporter.headers["x-langfuse-sdk-version"] == installed_langfuse_version()
|
||||
assert exporter.headers["x-langfuse-ingestion-version"] == "4"
|
||||
|
||||
|
||||
def test_enable_langfuse_debug_logging_makes_deliveries_visible_on_the_langfuse_logger(caplog):
|
||||
"""``LANGFUSE_DEBUG`` turned on the v2 SDK's own logger; it has to do the same for litellm's export channel."""
|
||||
exporter, _ = _exporter_over([200])
|
||||
langfuse_logger = logging.getLogger("langfuse")
|
||||
level_before = langfuse_logger.level
|
||||
try:
|
||||
with caplog.at_level(logging.INFO, logger="langfuse"):
|
||||
assert exporter.export((_finished_span(),)) is SpanExportResult.SUCCESS
|
||||
assert "Exported" not in caplog.text
|
||||
enable_langfuse_debug_logging()
|
||||
assert langfuse_logger.level == logging.DEBUG
|
||||
exporter_after, _ = _exporter_over([200])
|
||||
assert exporter_after.export((_finished_span(),)) is SpanExportResult.SUCCESS
|
||||
assert "Exported" in caplog.text and "lf.internal.example" in caplog.text
|
||||
finally:
|
||||
langfuse_logger.setLevel(level_before)
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
|
|
|
|||
|
|
@ -1019,7 +1019,9 @@ def test_failure_handler_langfuse_kwargs_excludes_original_response():
|
|||
|
||||
try:
|
||||
# Mock LangFuseHandler to return our capturing mock logger
|
||||
with patch("litellm.litellm_core_utils.litellm_logging.LangFuseHandler") as mock_handler_class: # test-quality-ok: route the request to the capturing logger; the real handler builds live clients
|
||||
with (
|
||||
patch("litellm.litellm_core_utils.litellm_logging.LangFuseHandler") as mock_handler_class
|
||||
): # test-quality-ok: route the request to the capturing logger; the real handler builds live clients
|
||||
mock_handler_class.get_langfuse_logger_for_request.return_value = mock_langfuse_logger
|
||||
|
||||
# Call the actual failure_handler
|
||||
|
|
@ -1086,7 +1088,9 @@ async def test_async_log_failure_event_logs_to_langfuse():
|
|||
"generation_id": "mock-gen",
|
||||
}
|
||||
|
||||
with patch("litellm.integrations.langfuse.langfuse_prompt_management.LangFuseHandler") as mock_handler: # test-quality-ok: route the request to the capturing logger; the real handler builds live clients
|
||||
with (
|
||||
patch("litellm.integrations.langfuse.langfuse_prompt_management.LangFuseHandler") as mock_handler
|
||||
): # test-quality-ok: route the request to the capturing logger; the real handler builds live clients
|
||||
mock_handler.get_langfuse_logger_for_request.return_value = mock_logger
|
||||
|
||||
kwargs = {
|
||||
|
|
@ -1151,7 +1155,9 @@ async def test_async_log_failure_event_works_without_standard_logging_object():
|
|||
"generation_id": "mock-gen",
|
||||
}
|
||||
|
||||
with patch("litellm.integrations.langfuse.langfuse_prompt_management.LangFuseHandler") as mock_handler: # test-quality-ok: route the request to the capturing logger; the real handler builds live clients
|
||||
with (
|
||||
patch("litellm.integrations.langfuse.langfuse_prompt_management.LangFuseHandler") as mock_handler
|
||||
): # test-quality-ok: route the request to the capturing logger; the real handler builds live clients
|
||||
mock_handler.get_langfuse_logger_for_request.return_value = mock_logger
|
||||
|
||||
kwargs = {
|
||||
|
|
@ -1288,7 +1294,7 @@ def test_max_langfuse_clients_limit():
|
|||
assert litellm.initialized_langfuse_clients == 2
|
||||
|
||||
# Third client should fail with exception
|
||||
with pytest.raises(Exception, match='Max langfuse clients reached') as exc_info:
|
||||
with pytest.raises(Exception, match="Max langfuse clients reached") as exc_info:
|
||||
logger3 = LangFuseLogger(
|
||||
langfuse_public_key="test_key_3",
|
||||
langfuse_secret="test_secret_3",
|
||||
|
|
@ -1384,9 +1390,7 @@ def test_dynamic_langfuse_environment_triggers_dynamic_logger():
|
|||
|
||||
assert LangFuseHandler._dynamic_langfuse_credentials_are_passed(params) is True
|
||||
|
||||
config = LangFuseHandler.get_dynamic_langfuse_logging_config(
|
||||
standard_callback_dynamic_params=params
|
||||
)
|
||||
config = LangFuseHandler.get_dynamic_langfuse_logging_config(standard_callback_dynamic_params=params)
|
||||
assert config["langfuse_environment"] == "team-a-env"
|
||||
|
||||
|
||||
|
|
@ -1413,7 +1417,7 @@ def test_langfuse_rest_client_survives_httpx_cache_eviction(monkeypatch):
|
|||
assert litellm.in_memory_llm_clients_cache.get_cache("httpx_client") is None
|
||||
assert handler_ref() is not None, "logger must keep the handler that owns the client behind its REST API"
|
||||
assert not logger.langfuse_client.is_closed
|
||||
assert logger.api_client.auth_check() is False
|
||||
assert logger.api_client.auth_check() is not None
|
||||
|
||||
|
||||
def test_langfuse_logger_reuses_the_shared_cached_client(monkeypatch):
|
||||
|
|
@ -1947,6 +1951,32 @@ def test_a_fresh_trace_under_a_callers_parent_still_carries_its_own_input_and_ou
|
|||
assert "the-output" in str(span.attributes["langfuse.trace.output"])
|
||||
|
||||
|
||||
def test_a_failed_call_under_a_callers_parent_stamps_the_error_as_the_trace_output():
|
||||
"""The ERROR branch used to write a trace-level ``status_message``, a field the v4 trace schema does not
|
||||
have, and skip ``output``; the generation's parent is the caller's, so nothing else fills the trace."""
|
||||
logger, exporter = _steering_logger()
|
||||
now = datetime.datetime.now()
|
||||
|
||||
logger.log_event_on_langfuse(
|
||||
kwargs={
|
||||
"call_type": "completion",
|
||||
"litellm_params": {"metadata": {"parent_observation_id": "0123456789abcdef"}},
|
||||
"messages": [{"role": "user", "content": "the-input"}],
|
||||
"optional_params": {},
|
||||
},
|
||||
response_obj=None,
|
||||
start_time=now,
|
||||
end_time=now,
|
||||
level="ERROR",
|
||||
status_message="provider said no",
|
||||
)
|
||||
span = _exported_span(logger, exporter)
|
||||
|
||||
assert span.parent is not None
|
||||
assert "provider said no" in str(span.attributes["langfuse.trace.output"])
|
||||
assert span.attributes["langfuse.observation.status_message"] == "provider said no"
|
||||
|
||||
|
||||
def test_a_fresh_trace_root_leaves_the_duplicate_io_to_langfuse():
|
||||
rig = _steering_logger()
|
||||
|
||||
|
|
@ -2227,6 +2257,24 @@ def test_langfuse_debug_env_string_false_stays_off(monkeypatch):
|
|||
assert LangFuseLogger().langfuse_debug is False
|
||||
|
||||
|
||||
def test_langfuse_debug_env_true_turns_on_the_langfuse_logger(monkeypatch):
|
||||
"""``LANGFUSE_DEBUG=true`` reached the v2 client as ``debug=`` and switched the SDK's logger to DEBUG;
|
||||
a parsed flag that nothing reads would make the variable a silent no-op."""
|
||||
monkeypatch.setenv("LANGFUSE_PUBLIC_KEY", "pk-debug-wire-test")
|
||||
monkeypatch.setenv("LANGFUSE_SECRET_KEY", "sk-debug-wire-test")
|
||||
monkeypatch.setenv("LANGFUSE_MOCK", "true")
|
||||
monkeypatch.setenv("LANGFUSE_DEBUG", "true")
|
||||
monkeypatch.setattr(litellm, "initialized_langfuse_clients", litellm.initialized_langfuse_clients)
|
||||
langfuse_logger = logging.getLogger("langfuse")
|
||||
level_before = langfuse_logger.level
|
||||
langfuse_logger.setLevel(logging.WARNING)
|
||||
try:
|
||||
assert LangFuseLogger().langfuse_debug is True
|
||||
assert langfuse_logger.level == logging.DEBUG
|
||||
finally:
|
||||
langfuse_logger.setLevel(level_before)
|
||||
|
||||
|
||||
def test_explicit_langfuse_host_beats_the_v4_base_url_env(monkeypatch):
|
||||
"""Per-key/per-team ``langfuse_host`` must win over LANGFUSE_BASE_URL.
|
||||
|
||||
|
|
@ -2269,7 +2317,7 @@ def test_resolve_credentials_falls_back_to_langfuse_base_url(monkeypatch):
|
|||
|
||||
|
||||
def test_version_gate_rejects_v5_prereleases():
|
||||
""""5.0.0rc1" sorts below "5", so a plain version comparison would admit it."""
|
||||
""" "5.0.0rc1" sorts below "5", so a plain version comparison would admit it."""
|
||||
langfuse_module.raise_if_unsupported_langfuse_version("4.7")
|
||||
with pytest.raises(ImportError):
|
||||
langfuse_module.raise_if_unsupported_langfuse_version("5.0.0rc1")
|
||||
|
|
|
|||
|
|
@ -4183,15 +4183,14 @@ async def test_health_services_endpoint_pointfive_blocks_non_admin(monkeypatch,
|
|||
|
||||
@pytest.mark.asyncio
|
||||
async def test_health_services_endpoint_langfuse_missing_keys_errors(monkeypatch):
|
||||
"""A disabled v4 client returns False from auth_check instead of raising.
|
||||
|
||||
v2 raised out of ``auth_check`` on missing keys, so the endpoint errored;
|
||||
the endpoint must not report success when the return value says the check
|
||||
failed.
|
||||
"""
|
||||
for key in ("LANGFUSE_PUBLIC_KEY", "LANGFUSE_SECRET_KEY", "LANGFUSE_HOST", "LANGFUSE_BASE_URL", "LANGFUSE_MOCK"):
|
||||
"""v2 raised out of ``auth_check`` and the endpoint printed the server's answer; the v4 check
|
||||
returns the failure as a value, and the endpoint has to error with that reason rather than a
|
||||
generic credentials message that reads the same for an outage and a bad key."""
|
||||
for key in ("LANGFUSE_PUBLIC_KEY", "LANGFUSE_SECRET_KEY", "LANGFUSE_BASE_URL", "LANGFUSE_MOCK"):
|
||||
monkeypatch.delenv(key, raising=False)
|
||||
monkeypatch.setenv("LANGFUSE_HOST", "http://127.0.0.1:9")
|
||||
monkeypatch.setattr(litellm, "initialized_langfuse_clients", litellm.initialized_langfuse_clients)
|
||||
|
||||
with pytest.raises(ProxyException, match="auth_check failed"):
|
||||
with pytest.raises(ProxyException, match="auth_check failed") as raised:
|
||||
await health_services_endpoint(service="langfuse")
|
||||
assert "127.0.0.1" in str(raised.value.message) or "refused" in str(raised.value.message).lower()
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue