feat(otel): add http/json export protocol for OTel v2 traces (#40290)

* feat(otel): add http/json export protocol for OTel v2 traces

OTEL_EXPORTER_OTLP_PROTOCOL=http/json was accepted but routed to the protobuf
OTLP/HTTP exporter, so collectors that only decode JSON rejected every batch.
Route it to an OTLP/JSON span exporter that reuses the SDK HTTP transport and
expose the protocol as a select field on the OpenTelemetry callback in the
admin UI.

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* refactor(otel): walk the fixed OTLP shape instead of recursing when hex-encoding ids

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* fix(ui): map stored callback variables onto their form fields when editing a callback

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

---------

Co-authored-by: yassin <yassin@berri.ai>
Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
devin-ai-integration[bot] 2026-09-08 15:45:32 -07:00 • committed by GitHub
parent 5b2b5420af
commit 183d05ae05
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
9 changed files with 334 additions and 28 deletions

View file

@ -367,6 +367,13 @@
"ui_name": "Headers",
"description": "Headers for OTEL exporter (e.g., x-honeycomb-team=YOUR_API_KEY)",
"required": false
},
"otel_exporter_otlp_protocol": {
"type": "select",
"ui_name": "Export Protocol",
"description": "OTLP wire format for trace exports. Use http/json for collectors that cannot decode protobuf",
"options": ["http/protobuf", "http/json"],
"required": false
}
},
"description": "OpenTelemetry Logging Integration"

View file

@ -69,7 +69,7 @@ class ExporterSpec(BaseModel):
kind: str = Field(
default="console",
description="console | in_memory | otlp_http | otlp_grpc | <factory kind>",
description="console | in_memory | otlp_http | http/json | otlp_grpc | <factory kind>",
)
endpoint: str | None = None
traces_endpoint: str | None = Field(

View file

@ -0,0 +1,70 @@
"""OTLP/HTTP span exporter that sends the OTLP/JSON encoding instead of protobuf.
The SDK only ships a protobuf OTLP/HTTP exporter; this reuses its transport and
retry loop and swaps the payload for OTLP/JSON (enums as integers, ids as hex).
"""
import base64
import json
from collections.abc import Mapping, Sequence
from types import MappingProxyType
from typing import Final, TypeAlias
from google.protobuf.json_format import MessageToDict
from opentelemetry.exporter.otlp.proto.common.trace_encoder import encode_spans
from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter
from opentelemetry.sdk.trace import ReadableSpan
JSON_CONTENT_TYPE: Final = "application/json"
_HEX_ID_KEYS: Final = frozenset({"traceId", "spanId", "parentSpanId"})
_JsonValue: TypeAlias = "Mapping[str, _JsonValue] | Sequence[_JsonValue] | str | int | float | bool | None"
_JsonObject: TypeAlias = Mapping[str, "_JsonValue"]
def _objects(node: _JsonObject, key: str) -> tuple[_JsonObject, ...]:
items: Final = node.get(key)
if isinstance(items, str) or not isinstance(items, Sequence):
return ()
return tuple(item for item in items if isinstance(item, Mapping))
def _hex_ids(node: _JsonObject) -> _JsonObject:
return MappingProxyType(
{
key: base64.b64decode(item).hex() if key in _HEX_ID_KEYS and isinstance(item, str) else item
for key, item in node.items()
}
)
def _hex_span(span: _JsonObject) -> _JsonObject:
links: Final = _objects(span, "links")
if not links:
return _hex_ids(span)
return MappingProxyType({**_hex_ids(span), "links": tuple(_hex_ids(link) for link in links)})
def _hex_scope_spans(scope: _JsonObject) -> _JsonObject:
return MappingProxyType({**scope, "spans": tuple(_hex_span(span) for span in _objects(scope, "spans"))})
def _hex_resource_spans(resource: _JsonObject) -> _JsonObject:
scope_spans: Final = tuple(_hex_scope_spans(scope) for scope in _objects(resource, "scopeSpans"))
return MappingProxyType({**resource, "scopeSpans": scope_spans})
def encode_spans_json(spans: Sequence[ReadableSpan]) -> bytes:
payload: Final[_JsonObject] = MessageToDict(encode_spans(spans), use_integers_for_enums=True)
resource_spans: Final = tuple(_hex_resource_spans(resource) for resource in _objects(payload, "resourceSpans"))
hexed: Final[_JsonObject] = MappingProxyType({**payload, "resourceSpans": resource_spans})
return json.dumps(hexed, default=dict, separators=(",", ":")).encode()
class OTLPJsonSpanExporter(OTLPSpanExporter):
def __init__(self, endpoint: str | None, headers: dict[str, str]) -> None: # mutable-ok: SDK __init__ takes Dict
super().__init__(endpoint=endpoint, headers=headers)
self._session.headers["Content-Type"] = JSON_CONTENT_TYPE
def _serialize_spans(self, spans: Sequence[ReadableSpan]) -> bytes:
return encode_spans_json(spans)

View file

@ -136,7 +136,8 @@ def parse_headers(raw: str | None) -> dict[str, str]:
_IN_MEMORY_KINDS: Final = ("in_memory", "inmemory", "memory")
_OTLP_HTTP_KINDS: Final = ("otlp_http", "http", "http/protobuf", "http/json")
_OTLP_HTTP_JSON_KINDS: Final = ("http/json",)
_OTLP_HTTP_KINDS: Final = ("otlp_http", "http", "http/protobuf", *_OTLP_HTTP_JSON_KINDS)
_OTLP_GRPC_KINDS: Final = ("otlp_grpc", "grpc")
@ -164,6 +165,13 @@ def _exporter_from_spec(spec: ExporterSpec) -> SpanExporter:
return factory(spec)
if kind in _IN_MEMORY_KINDS:
return InMemorySpanExporter()
if kind in _OTLP_HTTP_JSON_KINDS:
from litellm.integrations.otel.plumbing.otlp_json import OTLPJsonSpanExporter
return OTLPJsonSpanExporter(
endpoint=spec.traces_endpoint or _otlp_traces_endpoint(spec.endpoint),
headers=parse_headers(spec.headers),
)
if kind in _OTLP_HTTP_KINDS:
from opentelemetry.exporter.otlp.proto.http.trace_exporter import (
OTLPSpanExporter as HTTPExporter,

View file

@ -3584,6 +3584,7 @@ class AllCallbacks(LiteLLMPydanticObjectBase):
ui_callback_name="OpenTelemetry",
litellm_callback_params=[
"OTEL_EXPORTER",
"OTEL_EXPORTER_OTLP_PROTOCOL",
"OTEL_ENDPOINT",
"OTEL_TRACES_ENDPOINT",
"OTEL_HEADERS",

View file

@ -6,14 +6,18 @@ import json
import threading
from collections.abc import Iterator
from dataclasses import replace
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from http.server import BaseHTTPRequestHandler, HTTPServer, ThreadingHTTPServer
import pytest
pytest.importorskip("opentelemetry")
from opentelemetry.proto.collector.trace.v1.trace_service_pb2 import ( # noqa: E402
ExportTraceServiceRequest,
)
from opentelemetry.sdk.metrics import MeterProvider # noqa: E402
from opentelemetry.sdk.metrics.export import InMemoryMetricReader # noqa: E402
from opentelemetry.sdk.trace import TracerProvider # noqa: E402
from opentelemetry.sdk.trace.export import ( # noqa: E402
BatchSpanProcessor,
ConsoleSpanExporter,
@ -541,6 +545,89 @@ def test_build_span_exporter_variants():
assert "OTLPSpanExporter" in type(http_exporter).__name__
def _export_one_trace_to_local_collector(exporter_kind: str) -> tuple[list[dict], tuple[int, int, int]]:
"""Run a parent/child trace through the configured exporter against a
throwaway HTTP collector. Returns the requests as the collector saw them
(child first, since it ends first) and (trace_id, parent span_id, child span_id)."""
received: list[dict] = []
class Collector(BaseHTTPRequestHandler):
def do_POST(self):
body = self.rfile.read(int(self.headers["Content-Length"]))
received.append({"path": self.path, "headers": dict(self.headers), "body": body})
self.send_response(200)
self.end_headers()
def log_message(self, *_args):
pass
server = HTTPServer(("127.0.0.1", 0), Collector)
threading.Thread(target=server.serve_forever, daemon=True).start()
try:
config = OpenTelemetryV2Config(
exporter=exporter_kind,
endpoint=f"http://127.0.0.1:{server.server_port}",
headers="x-collector-token=secret",
)
provider = TracerProvider()
provider.add_span_processor(SimpleSpanProcessor(providers.build_span_exporter(config)))
tracer = provider.get_tracer("test")
with tracer.start_as_current_span("parent", kind=SpanKind.SERVER) as parent:
with tracer.start_as_current_span("child") as child:
ids = (
parent.get_span_context().trace_id,
parent.get_span_context().span_id,
child.get_span_context().span_id,
)
provider.shutdown()
finally:
server.shutdown()
server.server_close()
assert len(received) == 2
return received, ids
def _only_span(request: dict) -> dict:
scope_spans = json.loads(request["body"])["resourceSpans"][0]["scopeSpans"][0]["spans"]
assert len(scope_spans) == 1
return scope_spans[0]
def test_http_json_exporter_posts_otlp_json_to_traces_endpoint():
"""``http/json`` must put the OTLP/JSON mapping on the wire (camelCase
fields, integer enums, hex ids) with a JSON content type, so collectors that
cannot decode protobuf can ingest the trace. Headers still travel."""
(child_request, parent_request), (trace_id, parent_id, child_id) = _export_one_trace_to_local_collector("http/json")
assert parent_request["path"] == "/v1/traces"
assert parent_request["headers"]["Content-Type"] == "application/json"
assert parent_request["headers"]["x-collector-token"] == "secret"
parent = _only_span(parent_request)
assert parent["name"] == "parent"
assert parent["kind"] == 2
assert parent["traceId"] == format(trace_id, "032x")
assert parent["spanId"] == format(parent_id, "016x")
assert "parentSpanId" not in parent
child = _only_span(child_request)
assert child["traceId"] == format(trace_id, "032x")
assert child["spanId"] == format(child_id, "016x")
assert child["parentSpanId"] == format(parent_id, "016x")
def test_http_protobuf_exporter_still_posts_protobuf():
(_child_request, parent_request), (trace_id, _parent_id, _child_id) = _export_one_trace_to_local_collector(
"http/protobuf"
)
assert parent_request["path"] == "/v1/traces"
assert parent_request["headers"]["Content-Type"] == "application/x-protobuf"
assert format(trace_id, "032x").encode() not in parent_request["body"]
decoded = ExportTraceServiceRequest.FromString(parent_request["body"])
span = decoded.resource_spans[0].scope_spans[0].spans[0]
assert span.name == "parent"
assert span.trace_id == trace_id.to_bytes(16, "big")
@pytest.fixture
def otlp_collector() -> Iterator[tuple[str, list[str]]]:
received_paths: list[str] = []
@ -612,6 +699,21 @@ def test_traces_endpoint_per_exporter_coexists_with_default_normalization(otlp_c
assert sorted(received_paths) == ["/services/collector/traces", "/v1/traces"]
def test_http_json_exporter_honors_traces_endpoint(otlp_collector):
base_url, received_paths = otlp_collector
cfg = OpenTelemetryV2Config(
exporters=[
{
"kind": "http/json",
"endpoint": base_url,
"traces_endpoint": f"{base_url}/services/collector/traces",
}
]
)
_export_one_span(cfg)
assert received_paths == ["/services/collector/traces"]
def test_otlp_metric_exporter_uses_cumulative_histogram_temporality():
"""Histograms must export as cumulative, not delta.

View file

@ -158,6 +158,7 @@ export const CALLBACK_CONFIGS: CallbackConfig[] = [
dynamic_params: {
otel_endpoint: "text",
otel_headers: "text",
otel_exporter_otlp_protocol: "select",
},
description: "OpenTelemetry Logging Integration",
},

View file

@ -219,6 +219,89 @@ describe("Settings", () => {
expect(vi.mocked(setCallbacksCall)).not.toHaveBeenCalled();
});
const mockOtelCallback = (variables: Record<string, string | null>) => {
mockGetCallbacksCall.mockResolvedValue({
callbacks: [{ name: "otel", variables }],
available_callbacks: {
otel: {
litellm_callback_name: "otel",
litellm_callback_params: ["OTEL_EXPORTER", "OTEL_EXPORTER_OTLP_PROTOCOL", "OTEL_ENDPOINT", "OTEL_HEADERS"],
ui_callback_name: "OpenTelemetry",
},
},
alerts: [],
});
mockGetCallbackConfigsCall.mockResolvedValue([
{
id: "otel",
displayName: "Open Telemetry",
dynamic_params: {
otel_endpoint: { type: "text", ui_name: "Endpoint URL", required: true },
otel_exporter_otlp_protocol: {
type: "select",
ui_name: "Export Protocol",
options: ["http/protobuf", "http/json"],
required: false,
},
},
},
]);
};
const openOtelEditModal = async () => {
const user = userEvent.setup();
render(<Settings {...defaultProps} />);
await user.click(await screen.findByTestId("callback-actions-otel-success"));
await user.click(await screen.findByTestId("callback-action-edit"));
return user;
};
it("should post the chosen export protocol when a select dynamic param is saved", async () => {
mockOtelCallback({ OTEL_ENDPOINT: "http://collector:4318" });
const user = await openOtelEditModal();
expect(await screen.findByLabelText("Endpoint URL")).toHaveValue("http://collector:4318");
await user.click(screen.getByLabelText("Export Protocol"));
await user.click(await screen.findByRole("option", { name: "http/json" }));
expect(screen.getByLabelText("Export Protocol")).toHaveTextContent("http/json");
await user.click(within(screen.getByRole("dialog")).getByRole("button", { name: "Save Changes" }));
await waitFor(() => {
expect(vi.mocked(setCallbacksCall)).toHaveBeenCalledWith(
"token",
expect.objectContaining({
environment_variables: expect.objectContaining({
callback: "otel",
otel_endpoint: "http://collector:4318",
otel_exporter_otlp_protocol: "http/json",
}),
}),
);
});
});
it("should show the saved export protocol in the edit modal and keep it on an unchanged save", async () => {
mockOtelCallback({ OTEL_ENDPOINT: "http://collector:4318", OTEL_EXPORTER_OTLP_PROTOCOL: "http/json" });
const user = await openOtelEditModal();
expect(await screen.findByLabelText("Endpoint URL")).toHaveValue("http://collector:4318");
expect(screen.getByLabelText("Export Protocol")).toHaveTextContent("http/json");
await user.click(within(screen.getByRole("dialog")).getByRole("button", { name: "Save Changes" }));
await waitFor(() => {
expect(vi.mocked(setCallbacksCall)).toHaveBeenCalledWith("token", {
environment_variables: {
callback: "otel",
otel_endpoint: "http://collector:4318",
otel_exporter_otlp_protocol: "http/json",
},
litellm_settings: { success_callback: ["otel"] },
});
});
});
it("should send the typed webhook url for an alert type when the alerting tab is saved", async () => {
const user = userEvent.setup();
render(<Settings {...defaultProps} />);

View file

@ -14,6 +14,7 @@ import {
} from "@/components/ui/combobox";
import { Dialog, DialogContent, DialogHeader, DialogTitle } from "@/components/ui/dialog";
import { Input } from "@/components/ui/input";
import { Select, SelectContent, SelectItem, SelectTrigger, SelectValue } from "@/components/ui/select";
import { Switch } from "@/components/ui/switch";
import { Table, TableBody, TableCell, TableHead, TableHeader, TableRow } from "@/components/ui/table";
import { Tabs, TabsContent, TabsList, TabsTrigger } from "@/components/ui/tabs";
@ -59,7 +60,7 @@ interface DynamicParamsFieldsProps {
}
const DynamicParamsFields: React.FC<DynamicParamsFieldsProps> = ({ params, callbackConfigs, selectedCallback }) => {
const { register, formState } = useFormContext<CallbackFormValues>();
const { register, control, formState } = useFormContext<CallbackFormValues>();
const fieldIdPrefix = React.useId();
if (!params || params.length === 0) {
@ -74,37 +75,63 @@ const DynamicParamsFields: React.FC<DynamicParamsFieldsProps> = ({ params, callb
const paramType = paramConfig.type || "text";
const fieldLabel = paramConfig.ui_name || param.replace(/_/g, " ").replace(/\b\w/g, (l) => l.toUpperCase());
const isRequired = paramConfig.required || false;
const selectOptions: string[] = Array.isArray(paramConfig.options) ? paramConfig.options : [];
const isSelect = paramType === "select" && selectOptions.length > 0;
const fieldId = `${fieldIdPrefix}-${param}`;
const registration = register(
param,
isRequired ? { required: `Please enter the ${fieldLabel.toLowerCase()}` } : undefined,
);
const validationRules = isRequired ? { required: `Please enter the ${fieldLabel.toLowerCase()}` } : undefined;
const registration = isSelect ? undefined : register(param, validationRules);
return (
<Field key={param} className="mb-4">
<FieldLabel htmlFor={fieldId}>
<span className="text-sm font-medium text-foreground">{fieldLabel} </span>
</FieldLabel>
{paramType === "password" ? (
<Input
id={fieldId}
type="password"
placeholder={`Enter your ${fieldLabel.toLowerCase()}`}
{...registration}
{isSelect && (
<Controller
control={control}
name={param}
rules={validationRules}
render={({ field }) => (
<Select
items={selectOptions.map((option) => ({ label: option, value: option }))}
value={field.value || null}
onValueChange={(selected: string | null) => field.onChange(selected ?? "")}
>
<SelectTrigger id={fieldId} className="w-full" onBlur={field.onBlur}>
<SelectValue placeholder={`Select ${fieldLabel.toLowerCase()}`} />
</SelectTrigger>
<SelectContent>
{selectOptions.map((option) => (
<SelectItem key={option} value={option}>
{option}
</SelectItem>
))}
</SelectContent>
</Select>
)}
/>
) : paramType === "number" ? (
<Input
id={fieldId}
type="number"
placeholder={`Enter ${fieldLabel.toLowerCase()}`}
min={0}
max={1}
step={0.1}
{...registration}
/>
) : (
<Input id={fieldId} placeholder={`Enter your ${fieldLabel.toLowerCase()}`} {...registration} />
)}
{!isSelect &&
(paramType === "password" ? (
<Input
id={fieldId}
type="password"
placeholder={`Enter your ${fieldLabel.toLowerCase()}`}
{...registration}
/>
) : paramType === "number" ? (
<Input
id={fieldId}
type="number"
placeholder={`Enter ${fieldLabel.toLowerCase()}`}
min={0}
max={1}
step={0.1}
{...registration}
/>
) : (
<Input id={fieldId} placeholder={`Enter your ${fieldLabel.toLowerCase()}`} {...registration} />
))}
<FieldError errors={[formState.errors[param]]} />
</Field>
);
@ -271,15 +298,22 @@ const Settings: React.FC<SettingsPageProps> = ({ accessToken, userRole, userID,
useEffect(() => {
if (showEditCallback && selectedEditCallback) {
const params = getDynamicParamsForCallback(
selectedEditCallback.name,
callbackConfigs,
selectedEditCallback.variables,
);
const fieldNameFor = (variable: string) =>
params.find((param) => param.toUpperCase() === variable.toUpperCase()) ?? variable;
const normalized = Object.fromEntries(
Object.entries(selectedEditCallback.variables || {}).map(([k, v]) => [k, v ?? ""]),
Object.entries(selectedEditCallback.variables || {}).map(([k, v]) => [fieldNameFor(k), v ?? ""]),
);
editForm.reset({
...normalized,
callback: selectedEditCallback.name,
});
}
}, [showEditCallback, selectedEditCallback, editForm]);
}, [showEditCallback, selectedEditCallback, editForm, callbackConfigs]);
const handleSwitchChange = (alertName: string) => {
if (activeAlerts.includes(alertName)) {