From b31f57179d8036dcb8357089dee65291c59c48e9 Mon Sep 17 00:00:00 2001 From: Aditya Jha Date: Mon, 7 Sep 2026 16:33:12 +0530 Subject: [PATCH] Propagate OpenTelemetry context on inference API calls. Flow, Query Tuning, and Observability go through SingleStoreChatFactory without Analyst's trace hooks, so UMG never gets session/turn baggage and per-question cost rows stay empty. Co-authored-by: Cursor --- singlestoredb/ai/chat.py | 36 ++++++++++++++++++-- singlestoredb/tests/test_ai_chat.py | 52 +++++++++++++++++++++++++++++ 2 files changed, 86 insertions(+), 2 deletions(-) create mode 100644 singlestoredb/tests/test_ai_chat.py diff --git a/singlestoredb/ai/chat.py b/singlestoredb/ai/chat.py index 6636fe8d5..5e66d942b 100644 --- a/singlestoredb/ai/chat.py +++ b/singlestoredb/ai/chat.py @@ -9,6 +9,30 @@ from singlestoredb import manage_workspaces from singlestoredb.management.inference_api import InferenceAPIInfo + +def _inject_otel_headers(headers: Any) -> None: + try: + from opentelemetry.propagate import inject as otel_inject + except ImportError: + return + try: + otel_inject(headers) + except Exception: + return + + +def _httpx_inject_otel(request: httpx.Request) -> None: + _inject_otel_headers(request.headers) + + +def _attach_otel_request_hook( + client: Union[httpx.Client, httpx.AsyncClient], +) -> None: + hooks = client.event_hooks.setdefault('request', []) + if _httpx_inject_otel not in hooks: + hooks.append(_httpx_inject_otel) + + try: from langchain_openai import ChatOpenAI except ImportError: @@ -119,6 +143,7 @@ def _inject_headers(request: Any, **_ignored: Any) -> None: obo_val = obo_token_getter() if obo_val: request.headers['X-S2-OBO'] = obo_val + _inject_otel_headers(request.headers) request.headers.pop('X-Amz-Date', None) request.headers.pop('X-Amz-Security-Token', None) @@ -161,8 +186,15 @@ def _inject_headers(request: Any, **_ignored: Any) -> None: model=model_name, streaming=streaming, ) - if http_client is not None: - openai_kwargs['http_client'] = http_client + http_async_client = kwargs.pop('http_async_client', None) + if http_client is None: + http_client = httpx.Client(timeout=httpx.Timeout(None)) + _attach_otel_request_hook(http_client) + openai_kwargs['http_client'] = http_client + if http_async_client is None: + http_async_client = httpx.AsyncClient(timeout=httpx.Timeout(None)) + _attach_otel_request_hook(http_async_client) + openai_kwargs['http_async_client'] = http_async_client return ChatOpenAI( **openai_kwargs, **kwargs, diff --git a/singlestoredb/tests/test_ai_chat.py b/singlestoredb/tests/test_ai_chat.py new file mode 100644 index 000000000..93cfaef70 --- /dev/null +++ b/singlestoredb/tests/test_ai_chat.py @@ -0,0 +1,52 @@ +#!/usr/bin/env python +# type: ignore +"""SingleStoreChatFactory OTel header injection tests.""" +import unittest +from types import ModuleType +from unittest.mock import patch + +try: + import httpx + from singlestoredb.ai import chat as chat_mod +except ImportError: + httpx = None + chat_mod = None + + +@unittest.skipIf(chat_mod is None, 'singlestoredb.ai.chat dependencies missing') +class TestInjectOtelHeaders(unittest.TestCase): + + def test_calls_otel_inject(self): + headers = {} + + def fake_inject(target): + target['baggage'] = 'session=abc,turn=def' + + fake_propagate = ModuleType('opentelemetry.propagate') + fake_propagate.inject = fake_inject + fake_otel = ModuleType('opentelemetry') + with patch.dict( + 'sys.modules', + { + 'opentelemetry': fake_otel, + 'opentelemetry.propagate': fake_propagate, + }, + ): + chat_mod._inject_otel_headers(headers) + self.assertEqual(headers['baggage'], 'session=abc,turn=def') + + def test_httpx_hook_appends_once(self): + client = httpx.Client() + try: + chat_mod._attach_otel_request_hook(client) + chat_mod._attach_otel_request_hook(client) + self.assertEqual( + client.event_hooks['request'].count(chat_mod._httpx_inject_otel), + 1, + ) + finally: + client.close() + + +if __name__ == '__main__': + unittest.main()