Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
120 changes: 120 additions & 0 deletions packages/traceloop-sdk/tests/test_image_uploader.py
Original file line number Diff line number Diff line change
@@ -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"
11 changes: 9 additions & 2 deletions packages/traceloop-sdk/traceloop/sdk/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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,
Expand Down
6 changes: 6 additions & 0 deletions packages/traceloop-sdk/traceloop/sdk/images/__init__.py
Original file line number Diff line number Diff line change
@@ -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"]
66 changes: 13 additions & 53 deletions packages/traceloop-sdk/traceloop/sdk/images/image_uploader.py
Original file line number Diff line number Diff line change
@@ -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
Original file line number Diff line number Diff line change
@@ -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}")