From 3e1e417620aca2001d382974e16a7844780ac9b8 Mon Sep 17 00:00:00 2001 From: LuisFigueroaG Date: Mon, 27 Jul 2026 21:13:42 -0400 Subject: [PATCH] feat(sdk): make image uploader pluggable --- .../tests/test_image_uploader.py | 120 ++++++++++++++++++ .../traceloop-sdk/traceloop/sdk/__init__.py | 11 +- .../traceloop/sdk/images/__init__.py | 6 + .../traceloop/sdk/images/image_uploader.py | 66 ++-------- .../sdk/images/traceloop_image_uploader.py | 54 ++++++++ 5 files changed, 202 insertions(+), 55 deletions(-) create mode 100644 packages/traceloop-sdk/tests/test_image_uploader.py create mode 100644 packages/traceloop-sdk/traceloop/sdk/images/__init__.py create mode 100644 packages/traceloop-sdk/traceloop/sdk/images/traceloop_image_uploader.py diff --git a/packages/traceloop-sdk/tests/test_image_uploader.py b/packages/traceloop-sdk/tests/test_image_uploader.py new file mode 100644 index 0000000000..732bddd753 --- /dev/null +++ b/packages/traceloop-sdk/tests/test_image_uploader.py @@ -0,0 +1,120 @@ +from unittest.mock import AsyncMock, patch + +import pytest +from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( + InMemorySpanExporter, +) + +from traceloop.sdk import ( + ImageUploader as SDKImageUploader, + Traceloop, + TraceloopImageUploader as SDKTraceloopImageUploader, +) +from traceloop.sdk.images import ImageUploader, TraceloopImageUploader +from traceloop.sdk.images.image_uploader import ( + ImageUploader as LegacyPathImageUploader, +) + + +class RecordingImageUploader(ImageUploader): + def __init__(self) -> None: + self.calls = [] + + async def aupload_base64_image(self, trace_id: str, span_id: str, image_name: str, image_file: str) -> str: + self.calls.append((trace_id, span_id, image_name, image_file)) + return f"custom://{trace_id}/{span_id}/{image_name}" + + +class FalseyImageUploader(RecordingImageUploader): + def __bool__(self) -> bool: + return False + + +def test_image_uploader_is_abstract() -> None: + with pytest.raises(TypeError): + ImageUploader() # type: ignore[abstract] + + +def test_sync_upload_delegates_to_async_implementation() -> None: + uploader = RecordingImageUploader() + + url = uploader.upload_base64_image("trace", "span", "image.png", "base64") + + assert url == "custom://trace/span/image.png" + assert uploader.calls == [("trace", "span", "image.png", "base64")] + + +@pytest.mark.asyncio +async def test_async_upload_uses_custom_implementation() -> None: + uploader = RecordingImageUploader() + + url = await uploader.aupload_base64_image("trace", "span", "image.png", "base64") + + assert url == "custom://trace/span/image.png" + assert uploader.calls == [("trace", "span", "image.png", "base64")] + + +@pytest.mark.asyncio +async def test_traceloop_uploader_preserves_backend_request_contract() -> None: + uploader = TraceloopImageUploader("https://api.example.com", "secret") + + with ( + patch("traceloop.sdk.images.traceloop_image_uploader.requests.post") as post, + patch.object(uploader, "_async_upload", new_callable=AsyncMock) as upload, + ): + post.return_value.json.return_value = {"url": "https://uploads.example.com/image"} + + url = await uploader.aupload_base64_image("trace", "span", "image.png", "base64") + + assert url == "https://uploads.example.com/image" + post.assert_called_once_with( + "https://api.example.com/v2/traces/trace/spans/span/images", + json={"image_name": "image.png"}, + headers={ + "Authorization": "Bearer secret", + "Content-Type": "application/json", + }, + ) + upload.assert_awaited_once_with("https://uploads.example.com/image", "base64") + + +def test_public_imports_remain_available() -> None: + assert SDKImageUploader is ImageUploader + assert LegacyPathImageUploader is ImageUploader + assert SDKTraceloopImageUploader is TraceloopImageUploader + + +def test_custom_uploader_is_injected_even_when_falsey() -> None: + uploader = FalseyImageUploader() + + with ( + patch("traceloop.sdk.TracerWrapper") as tracer_wrapper, + patch("traceloop.sdk.is_metrics_enabled", return_value=False), + patch("traceloop.sdk.is_logging_enabled", return_value=False), + ): + Traceloop.init( + exporter=InMemorySpanExporter(), + image_uploader=uploader, + resource_attributes={}, + ) + + assert tracer_wrapper.call_args.kwargs["image_uploader"] is uploader + + +def test_traceloop_uploader_is_the_default() -> None: + with ( + patch("traceloop.sdk.TracerWrapper") as tracer_wrapper, + patch("traceloop.sdk.is_metrics_enabled", return_value=False), + patch("traceloop.sdk.is_logging_enabled", return_value=False), + ): + Traceloop.init( + exporter=InMemorySpanExporter(), + api_endpoint="https://api.example.com", + api_key="secret", + resource_attributes={}, + ) + + uploader = tracer_wrapper.call_args.kwargs["image_uploader"] + assert isinstance(uploader, TraceloopImageUploader) + assert uploader.base_url == "https://api.example.com" + assert uploader.api_key == "secret" diff --git a/packages/traceloop-sdk/traceloop/sdk/__init__.py b/packages/traceloop-sdk/traceloop/sdk/__init__.py index 6663429d99..9484ada150 100644 --- a/packages/traceloop-sdk/traceloop/sdk/__init__.py +++ b/packages/traceloop-sdk/traceloop/sdk/__init__.py @@ -14,7 +14,7 @@ from opentelemetry.propagators.textmap import TextMapPropagator from opentelemetry.util.re import parse_env_headers -from traceloop.sdk.images.image_uploader import ImageUploader +from traceloop.sdk.images import ImageUploader, TraceloopImageUploader from traceloop.sdk.metrics.metrics import MetricsWrapper from traceloop.sdk.logging.logging import LoggerWrapper from traceloop.sdk.instruments import Instruments @@ -88,6 +88,9 @@ def init( events have nowhere to go and no prompt/completion data will be recorded. use_legacy_attributes: Deprecated alias for ``use_attributes``. Will be removed in a future release. + image_uploader: Custom image storage implementation. Subclass + :class:`ImageUploader` and pass an instance here to replace the + default Traceloop backend uploader. """ if use_attributes is not None and use_legacy_attributes is not None: raise TypeError( @@ -188,7 +191,11 @@ def init( exporter=exporter, sampler=sampler, should_enrich_metrics=should_enrich_metrics, - image_uploader=image_uploader or ImageUploader(api_endpoint, api_key), + image_uploader=( + image_uploader + if image_uploader is not None + else TraceloopImageUploader(api_endpoint, api_key) + ), instruments=instruments, block_instruments=block_instruments, span_postprocess_callback=span_postprocess_callback, diff --git a/packages/traceloop-sdk/traceloop/sdk/images/__init__.py b/packages/traceloop-sdk/traceloop/sdk/images/__init__.py new file mode 100644 index 0000000000..17dae4bb9b --- /dev/null +++ b/packages/traceloop-sdk/traceloop/sdk/images/__init__.py @@ -0,0 +1,6 @@ +"""Interfaces and implementations for uploading images referenced by spans.""" + +from traceloop.sdk.images.image_uploader import ImageUploader +from traceloop.sdk.images.traceloop_image_uploader import TraceloopImageUploader + +__all__ = ["ImageUploader", "TraceloopImageUploader"] diff --git a/packages/traceloop-sdk/traceloop/sdk/images/image_uploader.py b/packages/traceloop-sdk/traceloop/sdk/images/image_uploader.py index cdd47f194a..3b2a371c5a 100644 --- a/packages/traceloop-sdk/traceloop/sdk/images/image_uploader.py +++ b/packages/traceloop-sdk/traceloop/sdk/images/image_uploader.py @@ -1,59 +1,19 @@ -import aiohttp import asyncio -import logging +from abc import ABC, abstractmethod -import requests +class ImageUploader(ABC): + """Interface for storing base64-encoded images referenced by spans. -class ImageUploader: - def __init__(self, base_url: str, api_key: str) -> None: - self.base_url = base_url - self.api_key = api_key - self.logger = logging.getLogger(__name__) + Custom implementations can be supplied through the ``image_uploader`` argument + of :meth:`traceloop.sdk.Traceloop.init`. + """ - def upload_base64_image( - self, trace_id: str, span_id: str, image_name: str, image_file: str - ) -> None: - asyncio.run(self.aupload_base64_image(trace_id, span_id, image_name, image_file)) + def upload_base64_image(self, trace_id: str, span_id: str, image_name: str, image_file: str) -> str: + """Upload an image synchronously and return its destination URL.""" + return asyncio.run(self.aupload_base64_image(trace_id, span_id, image_name, image_file)) - async def aupload_base64_image( - self, trace_id: str, span_id: str, image_name: str, image_file: str - ) -> str: - url = self._get_image_url(trace_id, span_id, image_name) - - await self._async_upload(url, image_file) - - return url - - def _get_image_url(self, trace_id: str, span_id: str, image_name: str) -> str: - response = requests.post( - f"{self.base_url}/v2/traces/{trace_id}/spans/{span_id}/images", - json={ - "image_name": image_name, - }, - headers={ - "Authorization": f"Bearer {self.api_key}", - "Content-Type": "application/json", - }, - ) - - return response.json()["url"] # type: ignore[no-any-return] - - async def _async_upload(self, url: str, base64_image: str) -> None: - headers = { - "Authorization": f"Bearer {self.api_key}", - "Content-Type": "application/json", - } - payload = { - "image_data": base64_image, - } - - async with aiohttp.ClientSession() as session: - async with session.post(url, json=payload, headers=headers) as response: - if response.status < 200 or response.status >= 300: - self.logger.error( - f"Failed to upload image. Status code: {response.status}" - ) - self.logger.error(await response.text()) - else: - self.logger.info(f"Successfully uploaded image {url}") + @abstractmethod + async def aupload_base64_image(self, trace_id: str, span_id: str, image_name: str, image_file: str) -> str: + """Upload an image asynchronously and return its destination URL.""" + raise NotImplementedError diff --git a/packages/traceloop-sdk/traceloop/sdk/images/traceloop_image_uploader.py b/packages/traceloop-sdk/traceloop/sdk/images/traceloop_image_uploader.py new file mode 100644 index 0000000000..b607914a77 --- /dev/null +++ b/packages/traceloop-sdk/traceloop/sdk/images/traceloop_image_uploader.py @@ -0,0 +1,54 @@ +import logging +from typing import Optional + +import aiohttp +import requests + +from traceloop.sdk.images.image_uploader import ImageUploader + + +class TraceloopImageUploader(ImageUploader): + """Upload span images to the Traceloop backend.""" + + def __init__(self, base_url: str, api_key: Optional[str]) -> None: + self.base_url = base_url + self.api_key = api_key + self.logger = logging.getLogger(__name__) + + async def aupload_base64_image(self, trace_id: str, span_id: str, image_name: str, image_file: str) -> str: + url = self._get_image_url(trace_id, span_id, image_name) + + await self._async_upload(url, image_file) + + return url + + def _get_image_url(self, trace_id: str, span_id: str, image_name: str) -> str: + response = requests.post( + f"{self.base_url}/v2/traces/{trace_id}/spans/{span_id}/images", + json={ + "image_name": image_name, + }, + headers={ + "Authorization": f"Bearer {self.api_key}", + "Content-Type": "application/json", + }, + ) + + return response.json()["url"] # type: ignore[no-any-return] + + async def _async_upload(self, url: str, base64_image: str) -> None: + headers = { + "Authorization": f"Bearer {self.api_key}", + "Content-Type": "application/json", + } + payload = { + "image_data": base64_image, + } + + async with aiohttp.ClientSession() as session: + async with session.post(url, json=payload, headers=headers) as response: + if response.status < 200 or response.status >= 300: + self.logger.error(f"Failed to upload image. Status code: {response.status}") + self.logger.error(await response.text()) + else: + self.logger.info(f"Successfully uploaded image {url}")