From ace27ab7fbf5605b5c23727a2bb00a423b03eb32 Mon Sep 17 00:00:00 2001 From: Quentin BEY Date: Fri, 5 Sep 2025 11:39:09 +0200 Subject: [PATCH] =?UTF-8?q?=F0=9F=93=88(langfuse)=20add=20light=20instrume?= =?UTF-8?q?ntation?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit This is a first implementation of Langfuse to see user interactions. Next step will be to allow prompt improvements. --- CHANGELOG.md | 1 + src/backend/chat/clients/pydantic_ai.py | 61 +++++++++++++++++++------ src/backend/conversations/settings.py | 32 +++++++++++++ src/backend/pyproject.toml | 1 + 4 files changed, 82 insertions(+), 13 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 57219fe..cf19042 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -29,6 +29,7 @@ and this project adheres to - ✨(backend) add feature flags from posthog #13 - ✨(user) allow to use conversation data for analytics #23 - ✨(chat) enforce response in user language #24 +- 📈(langfuse) add light instrumentation #26 [unreleased]: https://github.com/numerique-gouv/conversations/compare/HEAD...main diff --git a/src/backend/chat/clients/pydantic_ai.py b/src/backend/chat/clients/pydantic_ai.py index 7bfba88..a94e31b 100644 --- a/src/backend/chat/clients/pydantic_ai.py +++ b/src/backend/chat/clients/pydantic_ai.py @@ -11,7 +11,7 @@ import json import logging import time import uuid -from contextlib import AsyncExitStack +from contextlib import AsyncExitStack, ExitStack from typing import Dict, List, Optional, Tuple from django.conf import settings @@ -22,6 +22,7 @@ from django.utils.module_loading import import_string from django.utils.translation import gettext_lazy as _ from asgiref.sync import sync_to_async +from langfuse import get_client from pydantic import BaseModel from pydantic_ai import Agent, NativeOutput from pydantic_ai.messages import ( @@ -96,10 +97,13 @@ def _get_pydantic_agent(model_hrid, mcp_servers=None, **kwargs) -> Agent: ) -def _build_pydantic_agent(mcp_servers, model_hrid=None, language=None) -> Agent[None, str]: +def _build_pydantic_agent( + mcp_servers, model_hrid=None, language=None, instrument=False +) -> Agent[None, str]: """Create a Pydantic AI Agent instance with the configured settings.""" model_hrid = model_hrid or settings.LLM_DEFAULT_MODEL_HRID - agent = _get_pydantic_agent(model_hrid, mcp_servers) + + agent = _get_pydantic_agent(model_hrid, mcp_servers, instrument=instrument) @agent.system_prompt def add_the_date() -> str: @@ -120,7 +124,7 @@ def _build_pydantic_agent(mcp_servers, model_hrid=None, language=None) -> Agent[ return agent -def _build_routing_agent(model_hrid=None) -> Agent[None, str] | None: +def _build_routing_agent(model_hrid=None, instrument=False) -> Agent[None, str] | None: """ Create a Pydantic AI routing Agent instance with the configured settings. @@ -137,7 +141,11 @@ def _build_routing_agent(model_hrid=None) -> Agent[None, str] | None: model_hrid = model_hrid or settings.LLM_ROUTING_MODEL_HRID try: - agent = _get_pydantic_agent(model_hrid, output_type=NativeOutput([UserIntent])) + agent = _get_pydantic_agent( + model_hrid, + output_type=NativeOutput([UserIntent]), + instrument=instrument, + ) except ImproperlyConfigured: logger.info("AI routing model does not exist -> disabled") return None @@ -157,7 +165,7 @@ class UserIntent(BaseModel): attachment_summary: bool = False -class AIAgentService: +class AIAgentService: # pylint: disable=too-many-instance-attributes """Service class for AI-related operations (Pydantic-AI edition).""" def __init__(self, conversation, user, model_hrid=None, language=None): @@ -174,6 +182,8 @@ class AIAgentService: self.language = language # might be None self._last_stop_check = 0 + self._store_analytics = settings.LANGFUSE_ENABLED and user.allow_conversation_analytics + # Feature flags self._is_document_upload_enabled = is_feature_enabled(self.user, "document_upload") self._is_web_search_enabled = is_feature_enabled(self.user, "web_search") @@ -210,15 +220,23 @@ class AIAgentService: async def stream_text_async(self, messages: List[UIMessage], force_web_search: bool = False): """Return only the assistant text deltas (legacy text mode).""" await self._clean() - async for delta in self._run_agent(messages, force_web_search): - if delta["type"] == "0": - yield delta["payload"] + with ExitStack() as stack: + if self._store_analytics: + span = stack.enter_context(get_client().start_as_current_span(name="conversation")) + span.update_trace(user_id=str(self.user.sub), session_id=str(self.conversation.pk)) + async for delta in self._run_agent(messages, force_web_search): + if delta["type"] == "0": + yield delta["payload"] async def stream_data_async(self, messages: List[UIMessage], force_web_search: bool = False): """Return Vercel-AI-SDK formatted events.""" await self._clean() - async for delta in self._run_agent(messages, force_web_search): - yield f"{delta['type']}:{json.dumps(delta['payload'])}\n" + with ExitStack() as stack: + if self._store_analytics: + span = stack.enter_context(get_client().start_as_current_span(name="conversation")) + span.update_trace(user_id=str(self.user.sub), session_id=str(self.conversation.pk)) + async for delta in self._run_agent(messages, force_web_search): + yield f"{delta['type']}:{json.dumps(delta['payload'])}\n" async def _agent_stop_streaming(self, force_cache_check: Optional[bool] = False) -> None: """Check if the agent should stop streaming.""" @@ -268,7 +286,7 @@ class AIAgentService: ) return UserIntent() - agent = _build_routing_agent() + agent = _build_routing_agent(instrument=self._store_analytics) if not agent: return UserIntent() @@ -432,9 +450,20 @@ class AIAgentService: if messages[-1].role != "user": return + # Langfuse settings + if self._store_analytics: + langfuse = get_client() + langfuse.update_current_trace( + session_id=str(self.conversation.pk), + user_id=str(self.user.sub), + ) + history = ModelMessagesTypeAdapter.validate_python(self.conversation.pydantic_messages) user_prompt, input_images, input_documents = self.prepare_prompt(messages[-1]) + if self._store_analytics: + langfuse.update_current_trace(input=user_prompt) + usage = {"promptTokens": 0, "completionTokens": 0} # Feature flag management @@ -540,7 +569,10 @@ class AIAgentService: mcp_servers = [await stack.enter_async_context(mcp) for mcp in get_mcp_servers()] async with _build_pydantic_agent( - mcp_servers, model_hrid=self.model_hrid, language=self.language + mcp_servers, + model_hrid=self.model_hrid, + language=self.language, + instrument=self._store_analytics, ).iter( [user_prompt] + input_images, message_history=history, @@ -663,6 +695,9 @@ class AIAgentService: ui_sources=_ui_sources, ) + if self._store_analytics: + langfuse.update_current_trace(output=run.result.output) + # Vercel finish message yield { "type": "d", diff --git a/src/backend/conversations/settings.py b/src/backend/conversations/settings.py index d43643b..6b7d789 100755 --- a/src/backend/conversations/settings.py +++ b/src/backend/conversations/settings.py @@ -21,6 +21,7 @@ import posthog import sentry_sdk from configurations import Configuration, pristinemethod, values from corsheaders.defaults import default_headers +from langfuse import Langfuse from sentry_sdk.integrations.django import DjangoIntegration from sentry_sdk.integrations.logging import ignore_logger @@ -722,6 +723,24 @@ USER QUESTION: environ_prefix=None, ) + # LLM Instrumentation + LANGFUSE_ENABLED = values.BooleanValue( + default=False, environ_name="LANGFUSE_ENABLED", environ_prefix=None + ) + LANGFUSE_PUBLIC_KEY = values.Value( + None, environ_name="LANGFUSE_PUBLIC_KEY", environ_prefix=None + ) + LANGFUSE_SECRET_KEY = values.Value( + None, environ_name="LANGFUSE_SECRET_KEY", environ_prefix=None + ) + LANGFUSE_HOST = values.Value(None, environ_name="LANGFUSE_HOST", environ_prefix=None) + LANGFUSE_DEBUG = values.BooleanValue( + default=False, environ_name="LANGFUSE_DEBUG", environ_prefix=None + ) + LANGFUSE_MEDIA_UPLOAD_ENABLED = values.BooleanValue( + default=False, environ_name="LANGFUSE_MEDIA_UPLOAD_ENABLED", environ_prefix=None + ) + # pylint: disable=invalid-name @property def ENVIRONMENT(self): @@ -846,6 +865,19 @@ USER QUESTION: "OIDC_ALLOW_DUPLICATE_EMAILS cannot be set to True simultaneously. " ) + # Langfuse initialization + if cls.LANGFUSE_ENABLED: + if not cls.LANGFUSE_MEDIA_UPLOAD_ENABLED: + os.environ["LANGFUSE_MEDIA_UPLOAD_ENABLED"] = "false" + Langfuse( + public_key=cls.LANGFUSE_PUBLIC_KEY, + secret_key=cls.LANGFUSE_SECRET_KEY, + host=cls.LANGFUSE_HOST, + environment=cls.__name__.lower(), + release=get_release(), + debug=cls.LANGFUSE_DEBUG, + ) + class Build(Base): """Settings used when the application is built. diff --git a/src/backend/pyproject.toml b/src/backend/pyproject.toml index d682673..9661f75 100644 --- a/src/backend/pyproject.toml +++ b/src/backend/pyproject.toml @@ -47,6 +47,7 @@ dependencies = [ "factory_boy==3.3.3", "gunicorn==23.0.0", "jsonschema==4.24.0", + "langfuse==3.3.4", "lxml==5.4.0", "markdown==3.8", "markitdown==0.0.2",