Refactor langflow processing and langfuse callback initialization

This commit is contained in:
Gabriel Luiz Freitas Almeida 2024-02-09 12:00:24 -03:00
commit 2b646429c0
2 changed files with 6 additions and 5 deletions

View file

@ -2,10 +2,13 @@ from typing import TYPE_CHECKING, List, Union
from langchain.agents.agent import AgentExecutor from langchain.agents.agent import AgentExecutor
from langchain.callbacks.base import BaseCallbackHandler from langchain.callbacks.base import BaseCallbackHandler
from loguru import logger
from langflow.api.v1.callback import AsyncStreamingLLMCallbackHandler, StreamingLLMCallbackHandler from langflow.api.v1.callback import AsyncStreamingLLMCallbackHandler, StreamingLLMCallbackHandler
from langflow.processing.process import fix_memory_inputs, format_actions from langflow.processing.process import fix_memory_inputs, format_actions
from langflow.services.deps import get_plugins_service from langflow.services.deps import get_plugins_service
from loguru import logger from langflow.processing.process import fix_memory_inputs, format_actions
from langflow.services.deps import get_plugins_service
if TYPE_CHECKING: if TYPE_CHECKING:
from langfuse.callback import CallbackHandler # type: ignore from langfuse.callback import CallbackHandler # type: ignore
@ -28,13 +31,12 @@ def setup_callbacks(sync, trace_id, **kwargs):
def get_langfuse_callback(trace_id): def get_langfuse_callback(trace_id):
from langflow.services.deps import get_plugins_service from langflow.services.deps import get_plugins_service
from langfuse.callback import CreateTrace
logger.debug("Initializing langfuse callback") logger.debug("Initializing langfuse callback")
if langfuse := get_plugins_service().get("langfuse"): if langfuse := get_plugins_service().get("langfuse"):
logger.debug("Langfuse credentials found") logger.debug("Langfuse credentials found")
try: try:
trace = langfuse.trace(CreateTrace(name="langflow-" + trace_id, id=trace_id)) trace = langfuse.trace(name="langflow-" + trace_id, id=trace_id)
return trace.getNewHandler() return trace.getNewHandler()
except Exception as exc: except Exception as exc:
logger.error(f"Error initializing langfuse callback: {exc}") logger.error(f"Error initializing langfuse callback: {exc}")

View file

@ -64,14 +64,13 @@ class LangfusePlugin(CallbackPlugin):
def get_callback(self, _id: Optional[str] = None): def get_callback(self, _id: Optional[str] = None):
if _id is None: if _id is None:
_id = "default" _id = "default"
from langfuse.callback import CreateTrace # type: ignore
logger.debug("Initializing langfuse callback") logger.debug("Initializing langfuse callback")
try: try:
langfuse_instance = self.get() langfuse_instance = self.get()
if langfuse_instance is not None and hasattr(langfuse_instance, "trace"): if langfuse_instance is not None and hasattr(langfuse_instance, "trace"):
trace = langfuse_instance.trace(CreateTrace(name="langflow-" + _id, id=_id)) trace = langfuse_instance.trace(name="langflow-" + _id, id=_id)
if trace: if trace:
return trace.getNewHandler() return trace.getNewHandler()