ref: Add ruff rules for blind exceptions (BLE) (#4076)
Add ruff rules for blind exceptions (BLE)
This commit is contained in:
parent
d0bfac3e7e
commit
66be632086
78 changed files with 206 additions and 133 deletions
|
|
@ -10,6 +10,7 @@ import click
|
||||||
import httpx
|
import httpx
|
||||||
import typer
|
import typer
|
||||||
from dotenv import load_dotenv
|
from dotenv import load_dotenv
|
||||||
|
from httpx import HTTPError
|
||||||
from multiprocess import cpu_count
|
from multiprocess import cpu_count
|
||||||
from multiprocess.context import Process
|
from multiprocess.context import Process
|
||||||
from packaging import version as pkg_version
|
from packaging import version as pkg_version
|
||||||
|
|
@ -218,7 +219,7 @@ def run(
|
||||||
if process is not None:
|
if process is not None:
|
||||||
process.terminate()
|
process.terminate()
|
||||||
sys.exit(0)
|
sys.exit(0)
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
logger.exception(e)
|
logger.exception(e)
|
||||||
sys.exit(1)
|
sys.exit(1)
|
||||||
|
|
||||||
|
|
@ -231,7 +232,10 @@ def wait_for_server_ready(host, port):
|
||||||
while status_code != httpx.codes.OK:
|
while status_code != httpx.codes.OK:
|
||||||
try:
|
try:
|
||||||
status_code = httpx.get(f"http://{host}:{port}/health").status_code
|
status_code = httpx.get(f"http://{host}:{port}/health").status_code
|
||||||
except Exception:
|
except HTTPError:
|
||||||
|
time.sleep(1)
|
||||||
|
except Exception: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error while waiting for the server to become ready.")
|
||||||
time.sleep(1)
|
time.sleep(1)
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -48,16 +48,16 @@ async def health_check(
|
||||||
stmt = select(Flow).where(Flow.id == uuid.uuid4())
|
stmt = select(Flow).where(Flow.id == uuid.uuid4())
|
||||||
session.exec(stmt).first()
|
session.exec(stmt).first()
|
||||||
response.db = "ok"
|
response.db = "ok"
|
||||||
except Exception as e:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(e)
|
logger.exception("Error checking database")
|
||||||
|
|
||||||
try:
|
try:
|
||||||
chat = get_chat_service()
|
chat = get_chat_service()
|
||||||
await chat.set_cache("health_check", str(user_id))
|
await chat.set_cache("health_check", str(user_id))
|
||||||
await chat.get_cache("health_check")
|
await chat.get_cache("health_check")
|
||||||
response.chat = "ok"
|
response.chat = "ok"
|
||||||
except Exception as e:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(e)
|
logger.exception("Error checking chat service")
|
||||||
|
|
||||||
if response.has_error():
|
if response.has_error():
|
||||||
raise HTTPException(status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail=response.model_dump())
|
raise HTTPException(status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail=response.model_dump())
|
||||||
|
|
|
||||||
|
|
@ -74,7 +74,7 @@ class AsyncStreamingLLMCallbackHandleSIO(AsyncCallbackHandler):
|
||||||
# This is to emulate the stream of tokens
|
# This is to emulate the stream of tokens
|
||||||
for resp in resps:
|
for resp in resps:
|
||||||
await self.socketio_service.emit_token(to=self.sid, data=resp.model_dump())
|
await self.socketio_service.emit_token(to=self.sid, data=resp.model_dump())
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Error sending response")
|
logger.exception("Error sending response")
|
||||||
|
|
||||||
async def on_tool_error(
|
async def on_tool_error(
|
||||||
|
|
|
||||||
|
|
@ -60,7 +60,7 @@ async def try_running_celery_task(vertex, user_id):
|
||||||
|
|
||||||
task = build_vertex.delay(vertex)
|
task = build_vertex.delay(vertex)
|
||||||
vertex.task_id = task.id
|
vertex.task_id = task.id
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.opt(exception=True).debug("Error running task in celery")
|
logger.opt(exception=True).debug("Error running task in celery")
|
||||||
vertex.task_id = None
|
vertex.task_id = None
|
||||||
await vertex.build(user_id=user_id)
|
await vertex.build(user_id=user_id)
|
||||||
|
|
@ -167,8 +167,8 @@ async def build_flow(
|
||||||
if stop_component_id or start_component_id:
|
if stop_component_id or start_component_id:
|
||||||
try:
|
try:
|
||||||
first_layer = graph.sort_vertices(stop_component_id, start_component_id)
|
first_layer = graph.sort_vertices(stop_component_id, start_component_id)
|
||||||
except Exception as exc:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(exc)
|
logger.exception("Error sorting vertices")
|
||||||
first_layer = graph.sort_vertices()
|
first_layer = graph.sort_vertices()
|
||||||
else:
|
else:
|
||||||
first_layer = graph.sort_vertices()
|
first_layer = graph.sort_vertices()
|
||||||
|
|
@ -233,7 +233,7 @@ async def build_flow(
|
||||||
top_level_vertices = graph.get_top_level_vertices(next_runnable_vertices)
|
top_level_vertices = graph.get_top_level_vertices(next_runnable_vertices)
|
||||||
|
|
||||||
result_data_response = ResultDataResponse.model_validate(result_dict, from_attributes=True)
|
result_data_response = ResultDataResponse.model_validate(result_dict, from_attributes=True)
|
||||||
except Exception as exc:
|
except Exception as exc: # noqa: BLE001
|
||||||
if isinstance(exc, ComponentBuildException):
|
if isinstance(exc, ComponentBuildException):
|
||||||
params = exc.message
|
params = exc.message
|
||||||
tb = exc.formatted_traceback
|
tb = exc.formatted_traceback
|
||||||
|
|
@ -515,7 +515,7 @@ async def build_vertex(
|
||||||
next_runnable_vertices = await graph.get_next_runnable_vertices(lock, vertex=vertex, cache=False)
|
next_runnable_vertices = await graph.get_next_runnable_vertices(lock, vertex=vertex, cache=False)
|
||||||
top_level_vertices = graph.get_top_level_vertices(next_runnable_vertices)
|
top_level_vertices = graph.get_top_level_vertices(next_runnable_vertices)
|
||||||
result_data_response = ResultDataResponse.model_validate(result_dict, from_attributes=True)
|
result_data_response = ResultDataResponse.model_validate(result_dict, from_attributes=True)
|
||||||
except Exception as exc:
|
except Exception as exc: # noqa: BLE001
|
||||||
if isinstance(exc, ComponentBuildException):
|
if isinstance(exc, ComponentBuildException):
|
||||||
params = exc.message
|
params = exc.message
|
||||||
tb = exc.formatted_traceback
|
tb = exc.formatted_traceback
|
||||||
|
|
@ -690,7 +690,7 @@ async def build_vertex_stream(
|
||||||
msg = f"No result found for vertex {vertex_id}"
|
msg = f"No result found for vertex {vertex_id}"
|
||||||
raise ValueError(msg)
|
raise ValueError(msg)
|
||||||
|
|
||||||
except Exception as exc:
|
except Exception as exc: # noqa: BLE001
|
||||||
logger.exception("Error building Component")
|
logger.exception("Error building Component")
|
||||||
exc_message = parse_exception(exc)
|
exc_message = parse_exception(exc)
|
||||||
if exc_message == "The message must be an iterator or an async iterator.":
|
if exc_message == "The message must be an iterator or an async iterator.":
|
||||||
|
|
|
||||||
|
|
@ -168,7 +168,7 @@ async def simple_run_flow_task(
|
||||||
api_key_user=api_key_user,
|
api_key_user=api_key_user,
|
||||||
)
|
)
|
||||||
|
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(f"Error running flow {flow.id} task")
|
logger.exception(f"Error running flow {flow.id} task")
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -170,8 +170,8 @@ def read_flows(
|
||||||
for example_flow in example_flows:
|
for example_flow in example_flows:
|
||||||
if example_flow.id not in flow_ids:
|
if example_flow.id not in flow_ids:
|
||||||
flows.append(example_flow)
|
flows.append(example_flow)
|
||||||
except Exception as e:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(e)
|
logger.exception("Error getting example flows")
|
||||||
|
|
||||||
if remove_example_flows:
|
if remove_example_flows:
|
||||||
flows = [flow for flow in flows if flow.folder_id != folder.id]
|
flows = [flow for flow in flows if flow.folder_id != folder.id]
|
||||||
|
|
|
||||||
|
|
@ -42,7 +42,7 @@ def get_optional_user_store_api_key(
|
||||||
return None
|
return None
|
||||||
try:
|
try:
|
||||||
return auth_utils.decrypt_api_key(user.store_api_key, settings_service)
|
return auth_utils.decrypt_api_key(user.store_api_key, settings_service)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Failed to decrypt API key")
|
logger.exception("Failed to decrypt API key")
|
||||||
return user.store_api_key
|
return user.store_api_key
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -17,7 +17,8 @@ def post_validate_code(code: Code):
|
||||||
imports=errors.get("imports", {}),
|
imports=errors.get("imports", {}),
|
||||||
function=errors.get("function", {}),
|
function=errors.get("function", {}),
|
||||||
)
|
)
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error validating code")
|
||||||
return HTTPException(status_code=500, detail=str(e))
|
return HTTPException(status_code=500, detail=str(e))
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -102,7 +102,7 @@ class FlowTool(BaseTool):
|
||||||
tweaks = self.build_tweaks_dict(args, kwargs)
|
tweaks = self.build_tweaks_dict(args, kwargs)
|
||||||
try:
|
try:
|
||||||
run_id = self.graph.run_id if self.graph else None
|
run_id = self.graph.run_id if self.graph else None
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.opt(exception=True).warning("Failed to set run_id")
|
logger.opt(exception=True).warning("Failed to set run_id")
|
||||||
run_id = None
|
run_id = None
|
||||||
run_outputs = await run_flow(
|
run_outputs = await run_flow(
|
||||||
|
|
|
||||||
|
|
@ -4,6 +4,7 @@ from typing import Any
|
||||||
import requests
|
import requests
|
||||||
from bs4 import BeautifulSoup
|
from bs4 import BeautifulSoup
|
||||||
from langchain.tools import StructuredTool
|
from langchain.tools import StructuredTool
|
||||||
|
from loguru import logger
|
||||||
from markdown import markdown
|
from markdown import markdown
|
||||||
from pydantic import BaseModel, Field
|
from pydantic import BaseModel, Field
|
||||||
|
|
||||||
|
|
@ -82,7 +83,8 @@ class AddContentToPage(LCToolComponent):
|
||||||
if hasattr(e, "response") and e.response is not None:
|
if hasattr(e, "response") and e.response is not None:
|
||||||
error_message += f" Status code: {e.response.status_code}, Response: {e.response.text}"
|
error_message += f" Status code: {e.response.status_code}, Response: {e.response.text}"
|
||||||
return error_message
|
return error_message
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error adding content to Notion page")
|
||||||
return f"Error: An unexpected error occurred while adding content to Notion page. {e}"
|
return f"Error: An unexpected error occurred while adding content to Notion page. {e}"
|
||||||
|
|
||||||
def process_node(self, node):
|
def process_node(self, node):
|
||||||
|
|
|
||||||
|
|
@ -1,5 +1,6 @@
|
||||||
import requests
|
import requests
|
||||||
from langchain.tools import StructuredTool
|
from langchain.tools import StructuredTool
|
||||||
|
from loguru import logger
|
||||||
from pydantic import BaseModel, Field
|
from pydantic import BaseModel, Field
|
||||||
|
|
||||||
from langflow.base.langchain_utilities.model import LCToolComponent
|
from langflow.base.langchain_utilities.model import LCToolComponent
|
||||||
|
|
@ -62,5 +63,6 @@ class NotionDatabaseProperties(LCToolComponent):
|
||||||
return f"Error fetching Notion database properties: {e}"
|
return f"Error fetching Notion database properties: {e}"
|
||||||
except ValueError as e:
|
except ValueError as e:
|
||||||
return f"Error parsing Notion API response: {e}"
|
return f"Error parsing Notion API response: {e}"
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error fetching Notion database properties")
|
||||||
return f"An unexpected error occurred: {e}"
|
return f"An unexpected error occurred: {e}"
|
||||||
|
|
|
||||||
|
|
@ -3,6 +3,7 @@ from typing import Any
|
||||||
|
|
||||||
import requests
|
import requests
|
||||||
from langchain.tools import StructuredTool
|
from langchain.tools import StructuredTool
|
||||||
|
from loguru import logger
|
||||||
from pydantic import BaseModel, Field
|
from pydantic import BaseModel, Field
|
||||||
|
|
||||||
from langflow.base.langchain_utilities.model import LCToolComponent
|
from langflow.base.langchain_utilities.model import LCToolComponent
|
||||||
|
|
@ -116,5 +117,6 @@ class NotionListPages(LCToolComponent):
|
||||||
return f"Error querying Notion database: {e}"
|
return f"Error querying Notion database: {e}"
|
||||||
except KeyError:
|
except KeyError:
|
||||||
return "Unexpected response format from Notion API"
|
return "Unexpected response format from Notion API"
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error querying Notion database")
|
||||||
return f"An unexpected error occurred: {e}"
|
return f"An unexpected error occurred: {e}"
|
||||||
|
|
|
||||||
|
|
@ -1,5 +1,6 @@
|
||||||
import requests
|
import requests
|
||||||
from langchain.tools import StructuredTool
|
from langchain.tools import StructuredTool
|
||||||
|
from loguru import logger
|
||||||
from pydantic import BaseModel, Field
|
from pydantic import BaseModel, Field
|
||||||
|
|
||||||
from langflow.base.langchain_utilities.model import LCToolComponent
|
from langflow.base.langchain_utilities.model import LCToolComponent
|
||||||
|
|
@ -63,7 +64,8 @@ class NotionPageContent(LCToolComponent):
|
||||||
if hasattr(e, "response") and e.response is not None:
|
if hasattr(e, "response") and e.response is not None:
|
||||||
error_message += f" Status code: {e.response.status_code}, Response: {e.response.text}"
|
error_message += f" Status code: {e.response.status_code}, Response: {e.response.text}"
|
||||||
return error_message
|
return error_message
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error retrieving Notion page content")
|
||||||
return f"Error: An unexpected error occurred while retrieving Notion page content. {e}"
|
return f"Error: An unexpected error occurred while retrieving Notion page content. {e}"
|
||||||
|
|
||||||
def parse_blocks(self, blocks: list) -> str:
|
def parse_blocks(self, blocks: list) -> str:
|
||||||
|
|
|
||||||
|
|
@ -104,7 +104,7 @@ class NotionPageUpdate(LCToolComponent):
|
||||||
error_message = f"An error occurred while making the request: {e}"
|
error_message = f"An error occurred while making the request: {e}"
|
||||||
logger.exception(error_message)
|
logger.exception(error_message)
|
||||||
return error_message
|
return error_message
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
error_message = f"An unexpected error occurred: {e}"
|
error_message = f"An unexpected error occurred: {e}"
|
||||||
logger.exception(error_message)
|
logger.exception(error_message)
|
||||||
return error_message
|
return error_message
|
||||||
|
|
|
||||||
|
|
@ -1,4 +1,5 @@
|
||||||
import assemblyai as aai
|
import assemblyai as aai
|
||||||
|
from loguru import logger
|
||||||
|
|
||||||
from langflow.custom import Component
|
from langflow.custom import Component
|
||||||
from langflow.io import DataInput, DropdownInput, IntInput, Output, SecretStrInput
|
from langflow.io import DataInput, DropdownInput, IntInput, Output, SecretStrInput
|
||||||
|
|
@ -53,8 +54,9 @@ class AssemblyAIGetSubtitles(Component):
|
||||||
try:
|
try:
|
||||||
transcript_id = self.transcription_result.data["id"]
|
transcript_id = self.transcription_result.data["id"]
|
||||||
transcript = aai.Transcript.get_by_id(transcript_id)
|
transcript = aai.Transcript.get_by_id(transcript_id)
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
error = f"Getting transcription failed: {e}"
|
error = f"Getting transcription failed: {e}"
|
||||||
|
logger.opt(exception=True).debug(error)
|
||||||
self.status = error
|
self.status = error
|
||||||
return Data(data={"error": error})
|
return Data(data={"error": error})
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -1,4 +1,5 @@
|
||||||
import assemblyai as aai
|
import assemblyai as aai
|
||||||
|
from loguru import logger
|
||||||
|
|
||||||
from langflow.custom import Component
|
from langflow.custom import Component
|
||||||
from langflow.io import DataInput, DropdownInput, FloatInput, IntInput, MultilineInput, Output, SecretStrInput
|
from langflow.io import DataInput, DropdownInput, FloatInput, IntInput, MultilineInput, Output, SecretStrInput
|
||||||
|
|
@ -134,7 +135,8 @@ class AssemblyAILeMUR(Component):
|
||||||
result = Data(data=response)
|
result = Data(data=response)
|
||||||
self.status = result
|
self.status = result
|
||||||
return result
|
return result
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error running LeMUR")
|
||||||
error = f"An Error happened: {e}"
|
error = f"An Error happened: {e}"
|
||||||
self.status = error
|
self.status = error
|
||||||
return Data(data={"error": error})
|
return Data(data={"error": error})
|
||||||
|
|
|
||||||
|
|
@ -1,4 +1,5 @@
|
||||||
import assemblyai as aai
|
import assemblyai as aai
|
||||||
|
from loguru import logger
|
||||||
|
|
||||||
from langflow.custom import Component
|
from langflow.custom import Component
|
||||||
from langflow.io import BoolInput, DropdownInput, IntInput, MessageTextInput, Output, SecretStrInput
|
from langflow.io import BoolInput, DropdownInput, IntInput, MessageTextInput, Output, SecretStrInput
|
||||||
|
|
@ -85,7 +86,8 @@ class AssemblyAIListTranscripts(Component):
|
||||||
|
|
||||||
self.status = transcripts
|
self.status = transcripts
|
||||||
return transcripts
|
return transcripts
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error listing transcripts")
|
||||||
error_data = Data(data={"error": f"An error occurred: {e}"})
|
error_data = Data(data={"error": f"An error occurred: {e}"})
|
||||||
self.status = [error_data]
|
self.status = [error_data]
|
||||||
return [error_data]
|
return [error_data]
|
||||||
|
|
|
||||||
|
|
@ -1,4 +1,5 @@
|
||||||
import assemblyai as aai
|
import assemblyai as aai
|
||||||
|
from loguru import logger
|
||||||
|
|
||||||
from langflow.custom import Component
|
from langflow.custom import Component
|
||||||
from langflow.field_typing.range_spec import RangeSpec
|
from langflow.field_typing.range_spec import RangeSpec
|
||||||
|
|
@ -49,8 +50,9 @@ class AssemblyAITranscriptionJobPoller(Component):
|
||||||
|
|
||||||
try:
|
try:
|
||||||
transcript = aai.Transcript.get_by_id(self.transcript_id.data["transcript_id"])
|
transcript = aai.Transcript.get_by_id(self.transcript_id.data["transcript_id"])
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
error = f"Getting transcription failed: {e}"
|
error = f"Getting transcription failed: {e}"
|
||||||
|
logger.opt(exception=True).debug(error)
|
||||||
self.status = error
|
self.status = error
|
||||||
return Data(data={"error": error})
|
return Data(data={"error": error})
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -180,6 +180,7 @@ class AssemblyAITranscriptionJobCreator(Component):
|
||||||
result = Data(data={"transcript_id": transcript.id})
|
result = Data(data={"transcript_id": transcript.id})
|
||||||
self.status = result
|
self.status = result
|
||||||
return result
|
return result
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error submitting transcription job")
|
||||||
self.status = f"An error occurred: {e}"
|
self.status = f"An error occurred: {e}"
|
||||||
return Data(data={"error": f"An error occurred: {e}"})
|
return Data(data={"error": f"An error occurred: {e}"})
|
||||||
|
|
|
||||||
|
|
@ -131,7 +131,8 @@ class APIRequestComponent(Component):
|
||||||
response = await client.request(method, url, headers=headers, json=data, timeout=timeout)
|
response = await client.request(method, url, headers=headers, json=data, timeout=timeout)
|
||||||
try:
|
try:
|
||||||
result = response.json()
|
result = response.json()
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error decoding JSON response")
|
||||||
result = response.text
|
result = response.text
|
||||||
return Data(
|
return Data(
|
||||||
data={
|
data={
|
||||||
|
|
@ -150,7 +151,8 @@ class APIRequestComponent(Component):
|
||||||
"error": "Request timed out",
|
"error": "Request timed out",
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
except Exception as exc:
|
except Exception as exc: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug(f"Error making request to {url}")
|
||||||
return Data(
|
return Data(
|
||||||
data={
|
data={
|
||||||
"source": url,
|
"source": url,
|
||||||
|
|
|
||||||
|
|
@ -54,7 +54,7 @@ class SubFlowComponent(CustomComponent):
|
||||||
inputs = get_flow_inputs(graph)
|
inputs = get_flow_inputs(graph)
|
||||||
# Add inputs to the build config
|
# Add inputs to the build config
|
||||||
build_config = self.add_inputs_to_build_config(inputs, build_config)
|
build_config = self.add_inputs_to_build_config(inputs, build_config)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(f"Error getting flow {field_value}")
|
logger.exception(f"Error getting flow {field_value}")
|
||||||
|
|
||||||
return build_config
|
return build_config
|
||||||
|
|
|
||||||
|
|
@ -1,6 +1,8 @@
|
||||||
from datetime import datetime
|
from datetime import datetime
|
||||||
from zoneinfo import ZoneInfo
|
from zoneinfo import ZoneInfo
|
||||||
|
|
||||||
|
from loguru import logger
|
||||||
|
|
||||||
from langflow.custom import Component
|
from langflow.custom import Component
|
||||||
from langflow.io import DropdownInput, Output
|
from langflow.io import DropdownInput, Output
|
||||||
from langflow.schema.message import Message
|
from langflow.schema.message import Message
|
||||||
|
|
@ -67,7 +69,8 @@ class CurrentDateComponent(Component):
|
||||||
result = f"Current date and time in {self.timezone}: {current_date}"
|
result = f"Current date and time in {self.timezone}: {current_date}"
|
||||||
self.status = result
|
self.status = result
|
||||||
return Message(text=result)
|
return Message(text=result)
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error getting current date")
|
||||||
error_message = f"Error: {e}"
|
error_message = f"Error: {e}"
|
||||||
self.status = error_message
|
self.status = error_message
|
||||||
return Message(text=error_message)
|
return Message(text=error_message)
|
||||||
|
|
|
||||||
|
|
@ -1,3 +1,5 @@
|
||||||
|
from loguru import logger
|
||||||
|
|
||||||
from langflow.custom import Component
|
from langflow.custom import Component
|
||||||
from langflow.io import MessageInput, Output
|
from langflow.io import MessageInput, Output
|
||||||
from langflow.schema import Data
|
from langflow.schema import Data
|
||||||
|
|
@ -34,7 +36,8 @@ class MessageToDataComponent(Component):
|
||||||
|
|
||||||
self.status = "Successfully converted Message to Data"
|
self.status = "Successfully converted Message to Data"
|
||||||
return data
|
return data
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
error_message = f"Error converting Message to Data: {e}"
|
error_message = f"Error converting Message to Data: {e}"
|
||||||
|
logger.opt(exception=True).debug(error_message)
|
||||||
self.status = error_message
|
self.status = error_message
|
||||||
return Data(data={"error": error_message})
|
return Data(data={"error": error_message})
|
||||||
|
|
|
||||||
|
|
@ -85,7 +85,7 @@ class FlowToolComponent(LCToolComponent):
|
||||||
graph = Graph.from_payload(flow_data.data["data"])
|
graph = Graph.from_payload(flow_data.data["data"])
|
||||||
try:
|
try:
|
||||||
graph.set_run_id(self.graph.run_id)
|
graph.set_run_id(self.graph.run_id)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.opt(exception=True).warning("Failed to set run_id")
|
logger.opt(exception=True).warning("Failed to set run_id")
|
||||||
inputs = get_flow_inputs(graph)
|
inputs = get_flow_inputs(graph)
|
||||||
tool = FlowTool(
|
tool = FlowTool(
|
||||||
|
|
|
||||||
|
|
@ -1,5 +1,7 @@
|
||||||
from collections.abc import Callable
|
from collections.abc import Callable
|
||||||
|
|
||||||
|
from loguru import logger
|
||||||
|
|
||||||
from langflow.custom import Component
|
from langflow.custom import Component
|
||||||
from langflow.custom.utils import get_function
|
from langflow.custom.utils import get_function
|
||||||
from langflow.io import CodeInput, Output
|
from langflow.io import CodeInput, Output
|
||||||
|
|
@ -54,7 +56,8 @@ class PythonFunctionComponent(Component):
|
||||||
try:
|
try:
|
||||||
func = get_function(function_code)
|
func = get_function(function_code)
|
||||||
return func()
|
return func()
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error executing function")
|
||||||
return f"Error executing function: {e}"
|
return f"Error executing function: {e}"
|
||||||
|
|
||||||
def execute_function_data(self) -> list[Data]:
|
def execute_function_data(self) -> list[Data]:
|
||||||
|
|
|
||||||
|
|
@ -46,7 +46,7 @@ class SubFlowComponent(Component):
|
||||||
inputs = get_flow_inputs(graph)
|
inputs = get_flow_inputs(graph)
|
||||||
# Add inputs to the build config
|
# Add inputs to the build config
|
||||||
build_config = self.add_inputs_to_build_config(inputs, build_config)
|
build_config = self.add_inputs_to_build_config(inputs, build_config)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(f"Error getting flow {field_value}")
|
logger.exception(f"Error getting flow {field_value}")
|
||||||
|
|
||||||
return build_config
|
return build_config
|
||||||
|
|
|
||||||
|
|
@ -65,7 +65,8 @@ class ComposioAPIComponent(LCToolComponent):
|
||||||
try:
|
try:
|
||||||
entity.get_connection(app=app)
|
entity.get_connection(app=app)
|
||||||
return f"{app} CONNECTED"
|
return f"{app} CONNECTED"
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Authorization error")
|
||||||
return self._handle_authorization_failure(toolset, entity, app)
|
return self._handle_authorization_failure(toolset, entity, app)
|
||||||
|
|
||||||
def _handle_authorization_failure(self, toolset: ComposioToolSet, entity: Any, app: str) -> str:
|
def _handle_authorization_failure(self, toolset: ComposioToolSet, entity: Any, app: str) -> str:
|
||||||
|
|
@ -85,7 +86,7 @@ class ComposioAPIComponent(LCToolComponent):
|
||||||
if auth_schemes[0].auth_mode == "API_KEY":
|
if auth_schemes[0].auth_mode == "API_KEY":
|
||||||
return self._process_api_key_auth(entity, app)
|
return self._process_api_key_auth(entity, app)
|
||||||
return self._initiate_default_connection(entity, app)
|
return self._initiate_default_connection(entity, app)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Authorization error")
|
logger.exception("Authorization error")
|
||||||
return "Error"
|
return "Error"
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -2,6 +2,7 @@ import ast
|
||||||
import operator
|
import operator
|
||||||
|
|
||||||
from langchain.tools import StructuredTool
|
from langchain.tools import StructuredTool
|
||||||
|
from loguru import logger
|
||||||
from pydantic import BaseModel, Field
|
from pydantic import BaseModel, Field
|
||||||
|
|
||||||
from langflow.base.langchain_utilities.model import LCToolComponent
|
from langflow.base.langchain_utilities.model import LCToolComponent
|
||||||
|
|
@ -76,7 +77,8 @@ class CalculatorToolComponent(LCToolComponent):
|
||||||
error_message = "Error: Division by zero"
|
error_message = "Error: Division by zero"
|
||||||
self.status = error_message
|
self.status = error_message
|
||||||
return [Data(data={"error": error_message})]
|
return [Data(data={"error": error_message})]
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error evaluating expression")
|
||||||
error_message = f"Error: {e}"
|
error_message = f"Error: {e}"
|
||||||
self.status = error_message
|
self.status = error_message
|
||||||
return [Data(data={"error": error_message})]
|
return [Data(data={"error": error_message})]
|
||||||
|
|
|
||||||
|
|
@ -4,6 +4,7 @@ from typing import Any
|
||||||
|
|
||||||
from langchain.agents import Tool
|
from langchain.agents import Tool
|
||||||
from langchain_core.tools import StructuredTool
|
from langchain_core.tools import StructuredTool
|
||||||
|
from loguru import logger
|
||||||
from pydantic.v1 import Field, create_model
|
from pydantic.v1 import Field, create_model
|
||||||
from pydantic.v1.fields import Undefined
|
from pydantic.v1.fields import Undefined
|
||||||
|
|
||||||
|
|
@ -119,8 +120,9 @@ class PythonCodeStructuredTool(LCToolComponent):
|
||||||
build_config["_functions"]["value"] = json.dumps(named_functions)
|
build_config["_functions"]["value"] = json.dumps(named_functions)
|
||||||
build_config["_classes"]["value"] = json.dumps(classes)
|
build_config["_classes"]["value"] = json.dumps(classes)
|
||||||
build_config["tool_function"]["options"] = names
|
build_config["tool_function"]["options"] = names
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
self.status = f"Failed to extract names: {e}"
|
self.status = f"Failed to extract names: {e}"
|
||||||
|
logger.opt(exception=True).debug(self.status)
|
||||||
build_config["tool_function"]["options"] = ["Failed to parse", str(e)]
|
build_config["tool_function"]["options"] = ["Failed to parse", str(e)]
|
||||||
return build_config
|
return build_config
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -2,6 +2,7 @@ import importlib
|
||||||
|
|
||||||
from langchain.tools import StructuredTool
|
from langchain.tools import StructuredTool
|
||||||
from langchain_experimental.utilities import PythonREPL
|
from langchain_experimental.utilities import PythonREPL
|
||||||
|
from loguru import logger
|
||||||
from pydantic import BaseModel, Field
|
from pydantic import BaseModel, Field
|
||||||
|
|
||||||
from langflow.base.langchain_utilities.model import LCToolComponent
|
from langflow.base.langchain_utilities.model import LCToolComponent
|
||||||
|
|
@ -73,7 +74,8 @@ class PythonREPLToolComponent(LCToolComponent):
|
||||||
def run_python_code(code: str) -> str:
|
def run_python_code(code: str) -> str:
|
||||||
try:
|
try:
|
||||||
return python_repl.run(code)
|
return python_repl.run(code)
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error running Python code")
|
||||||
return f"Error: {e}"
|
return f"Error: {e}"
|
||||||
|
|
||||||
tool = StructuredTool.from_function(
|
tool = StructuredTool.from_function(
|
||||||
|
|
|
||||||
|
|
@ -5,6 +5,7 @@ from typing import Any
|
||||||
import requests
|
import requests
|
||||||
from langchain.agents import Tool
|
from langchain.agents import Tool
|
||||||
from langchain_core.tools import StructuredTool
|
from langchain_core.tools import StructuredTool
|
||||||
|
from loguru import logger
|
||||||
from pydantic.v1 import Field, create_model
|
from pydantic.v1 import Field, create_model
|
||||||
|
|
||||||
from langflow.base.langchain_utilities.model import LCToolComponent
|
from langflow.base.langchain_utilities.model import LCToolComponent
|
||||||
|
|
@ -72,8 +73,9 @@ class SearXNGToolComponent(LCToolComponent):
|
||||||
build_config["categories"]["value"].remove(selected_category)
|
build_config["categories"]["value"].remove(selected_category)
|
||||||
languages = list(data["locales"])
|
languages = list(data["locales"])
|
||||||
build_config["language"]["options"] = languages.copy()
|
build_config["language"]["options"] = languages.copy()
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
self.status = f"Failed to extract names: {e}"
|
self.status = f"Failed to extract names: {e}"
|
||||||
|
logger.opt(exception=True).debug(self.status)
|
||||||
build_config["categories"]["options"] = ["Failed to parse", str(e)]
|
build_config["categories"]["options"] = ["Failed to parse", str(e)]
|
||||||
return build_config
|
return build_config
|
||||||
|
|
||||||
|
|
@ -107,7 +109,8 @@ class SearXNGToolComponent(LCToolComponent):
|
||||||
|
|
||||||
num_results = min(SearxSearch._max_results, len(response["results"]))
|
num_results = min(SearxSearch._max_results, len(response["results"]))
|
||||||
return [response["results"][i] for i in range(num_results)]
|
return [response["results"][i] for i in range(num_results)]
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error running SearXNG Search")
|
||||||
return [f"Failed to search: {e}"]
|
return [f"Failed to search: {e}"]
|
||||||
|
|
||||||
SearxSearch._url = self.url
|
SearxSearch._url = self.url
|
||||||
|
|
|
||||||
|
|
@ -2,6 +2,7 @@ from typing import Any
|
||||||
|
|
||||||
from langchain.tools import StructuredTool
|
from langchain.tools import StructuredTool
|
||||||
from langchain_community.utilities.serpapi import SerpAPIWrapper
|
from langchain_community.utilities.serpapi import SerpAPIWrapper
|
||||||
|
from loguru import logger
|
||||||
from pydantic import BaseModel, Field
|
from pydantic import BaseModel, Field
|
||||||
|
|
||||||
from langflow.base.langchain_utilities.model import LCToolComponent
|
from langflow.base.langchain_utilities.model import LCToolComponent
|
||||||
|
|
@ -87,6 +88,7 @@ class SerpAPIComponent(LCToolComponent):
|
||||||
|
|
||||||
self.status = data_list
|
self.status = data_list
|
||||||
return data_list
|
return data_list
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error running SerpAPI")
|
||||||
self.status = f"Error: {e}"
|
self.status = f"Error: {e}"
|
||||||
return [Data(data={"error": str(e)}, text=str(e))]
|
return [Data(data={"error": str(e)}, text=str(e))]
|
||||||
|
|
|
||||||
|
|
@ -2,6 +2,7 @@ from typing import Any
|
||||||
|
|
||||||
import httpx
|
import httpx
|
||||||
from langchain.tools import StructuredTool
|
from langchain.tools import StructuredTool
|
||||||
|
from loguru import logger
|
||||||
from pydantic import BaseModel, Field
|
from pydantic import BaseModel, Field
|
||||||
|
|
||||||
from langflow.base.langchain_utilities.model import LCToolComponent
|
from langflow.base.langchain_utilities.model import LCToolComponent
|
||||||
|
|
@ -154,7 +155,8 @@ Note: Check 'Advanced' for all options.
|
||||||
error_message = f"HTTP error: {e.response.status_code} - {e.response.text}"
|
error_message = f"HTTP error: {e.response.status_code} - {e.response.text}"
|
||||||
self.status = error_message
|
self.status = error_message
|
||||||
return [Data(data={"error": error_message})]
|
return [Data(data={"error": error_message})]
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error running Tavily Search")
|
||||||
error_message = f"Unexpected error: {e}"
|
error_message = f"Unexpected error: {e}"
|
||||||
self.status = error_message
|
self.status = error_message
|
||||||
return [Data(data={"error": error_message})]
|
return [Data(data={"error": error_message})]
|
||||||
|
|
|
||||||
|
|
@ -3,6 +3,7 @@ import pprint
|
||||||
|
|
||||||
import yfinance as yf
|
import yfinance as yf
|
||||||
from langchain.tools import StructuredTool
|
from langchain.tools import StructuredTool
|
||||||
|
from loguru import logger
|
||||||
from pydantic import BaseModel, Field
|
from pydantic import BaseModel, Field
|
||||||
|
|
||||||
from langflow.base.langchain_utilities.model import LCToolComponent
|
from langflow.base.langchain_utilities.model import LCToolComponent
|
||||||
|
|
@ -95,7 +96,8 @@ class YfinanceToolComponent(LCToolComponent):
|
||||||
|
|
||||||
return data_list
|
return data_list
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
error_message = f"Error retrieving data: {e}"
|
error_message = f"Error retrieving data: {e}"
|
||||||
|
logger.opt(exception=True).debug(error_message)
|
||||||
self.status = error_message
|
self.status = error_message
|
||||||
return [Data(data={"error": error_message})]
|
return [Data(data={"error": error_message})]
|
||||||
|
|
|
||||||
|
|
@ -342,7 +342,7 @@ class CodeParser:
|
||||||
for import_node in import_nodes:
|
for import_node in import_nodes:
|
||||||
self.parse_imports(import_node)
|
self.parse_imports(import_node)
|
||||||
nodes.append(class_node)
|
nodes.append(class_node)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Error finding base class node")
|
logger.exception("Error finding base class node")
|
||||||
nodes.insert(0, node)
|
nodes.insert(0, node)
|
||||||
class_details = ClassCodeDetails(
|
class_details = ClassCodeDetails(
|
||||||
|
|
|
||||||
|
|
@ -8,6 +8,7 @@ from typing import TYPE_CHECKING, Any, ClassVar, get_type_hints
|
||||||
|
|
||||||
import nanoid
|
import nanoid
|
||||||
import yaml
|
import yaml
|
||||||
|
from loguru import logger
|
||||||
from pydantic import BaseModel
|
from pydantic import BaseModel
|
||||||
|
|
||||||
from langflow.base.tools.constants import TOOL_OUTPUT_NAME
|
from langflow.base.tools.constants import TOOL_OUTPUT_NAME
|
||||||
|
|
@ -338,7 +339,8 @@ class Component(CustomComponent):
|
||||||
try:
|
try:
|
||||||
source_code = inspect.getsource(method)
|
source_code = inspect.getsource(method)
|
||||||
ast_tree = ast.parse(dedent(source_code))
|
ast_tree = ast.parse(dedent(source_code))
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug(f"Could not get source code for method {method}")
|
||||||
source_code = self._code
|
source_code = self._code
|
||||||
ast_tree = ast.parse(dedent(source_code))
|
ast_tree = ast.parse(dedent(source_code))
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -77,9 +77,8 @@ class DirectoryReader:
|
||||||
if component["error"] if with_errors else not component["error"]:
|
if component["error"] if with_errors else not component["error"]:
|
||||||
component_tuple = (*build_component(component), component)
|
component_tuple = (*build_component(component), component)
|
||||||
components.append(component_tuple)
|
components.append(component_tuple)
|
||||||
except Exception as e:
|
except Exception: # noqa: BLE001
|
||||||
logger.debug(f"Error while loading component {component['name']}")
|
logger.opt(exception=True).debug(f"Error while loading component {component['name']}")
|
||||||
logger.debug(e)
|
|
||||||
continue
|
continue
|
||||||
items.append({"name": menu["name"], "path": menu["path"], "components": components})
|
items.append({"name": menu["name"], "path": menu["path"], "components": components})
|
||||||
filtered = [menu for menu in items if menu["components"]]
|
filtered = [menu for menu in items if menu["components"]]
|
||||||
|
|
@ -215,7 +214,7 @@ class DirectoryReader:
|
||||||
"""
|
"""
|
||||||
try:
|
try:
|
||||||
file_content = self.read_file_content(file_path)
|
file_content = self.read_file_content(file_path)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(f"Error while reading file {file_path}")
|
logger.exception(f"Error while reading file {file_path}")
|
||||||
return False, f"Could not read {file_path}"
|
return False, f"Could not read {file_path}"
|
||||||
|
|
||||||
|
|
@ -270,7 +269,8 @@ class DirectoryReader:
|
||||||
if validation_result:
|
if validation_result:
|
||||||
try:
|
try:
|
||||||
output_types = self.get_output_types_from_code(result_content)
|
output_types = self.get_output_types_from_code(result_content)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error while getting output types from code")
|
||||||
output_types = [component_name_camelcase]
|
output_types = [component_name_camelcase]
|
||||||
else:
|
else:
|
||||||
output_types = [component_name_camelcase]
|
output_types = [component_name_camelcase]
|
||||||
|
|
@ -292,7 +292,7 @@ class DirectoryReader:
|
||||||
async def process_file_async(self, file_path):
|
async def process_file_async(self, file_path):
|
||||||
try:
|
try:
|
||||||
file_content = self.read_file_content(file_path)
|
file_content = self.read_file_content(file_path)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(f"Error while reading file {file_path}")
|
logger.exception(f"Error while reading file {file_path}")
|
||||||
return False, f"Could not read {file_path}"
|
return False, f"Could not read {file_path}"
|
||||||
|
|
||||||
|
|
@ -346,7 +346,7 @@ class DirectoryReader:
|
||||||
if validation_result:
|
if validation_result:
|
||||||
try:
|
try:
|
||||||
output_types = await self.get_output_types_from_code_async(result_content)
|
output_types = await self.get_output_types_from_code_async(result_content)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Error while getting output types from code")
|
logger.exception("Error while getting output types from code")
|
||||||
output_types = [component_name_camelcase]
|
output_types = [component_name_camelcase]
|
||||||
else:
|
else:
|
||||||
|
|
|
||||||
|
|
@ -132,7 +132,7 @@ def build_invalid_menu_items(menu_item):
|
||||||
component_name, component_template = build_invalid_component(component)
|
component_name, component_template = build_invalid_component(component)
|
||||||
menu_items[component_name] = component_template
|
menu_items[component_name] = component_template
|
||||||
logger.debug(f"Added {component_name} to invalid menu.")
|
logger.debug(f"Added {component_name} to invalid menu.")
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(f"Error while creating custom component [{component_name}]")
|
logger.exception(f"Error while creating custom component [{component_name}]")
|
||||||
return menu_items
|
return menu_items
|
||||||
|
|
||||||
|
|
@ -165,6 +165,6 @@ def build_menu_items(menu_item):
|
||||||
for component_name, component_template, component in menu_item["components"]:
|
for component_name, component_template, component in menu_item["components"]:
|
||||||
try:
|
try:
|
||||||
menu_items[component_name] = component_template
|
menu_items[component_name] = component_template
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(f"Error while building custom component {component['output_types']}")
|
logger.exception(f"Error while building custom component {component['output_types']}")
|
||||||
return menu_items
|
return menu_items
|
||||||
|
|
|
||||||
|
|
@ -121,7 +121,7 @@ class Graph:
|
||||||
self._snapshots: list[dict[str, Any]] = []
|
self._snapshots: list[dict[str, Any]] = []
|
||||||
try:
|
try:
|
||||||
self.tracing_service: TracingService | None = get_tracing_service()
|
self.tracing_service: TracingService | None = get_tracing_service()
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Error getting tracing service")
|
logger.exception("Error getting tracing service")
|
||||||
self.tracing_service = None
|
self.tracing_service = None
|
||||||
if start is not None and end is not None:
|
if start is not None and end is not None:
|
||||||
|
|
@ -676,8 +676,8 @@ class Graph:
|
||||||
cache_service = get_chat_service()
|
cache_service = get_chat_service()
|
||||||
if self.flow_id:
|
if self.flow_id:
|
||||||
await cache_service.set_cache(self.flow_id, self)
|
await cache_service.set_cache(self.flow_id, self)
|
||||||
except Exception as exc:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(exc)
|
logger.exception("Error setting cache")
|
||||||
|
|
||||||
try:
|
try:
|
||||||
# Prioritize the webhook component if it exists
|
# Prioritize the webhook component if it exists
|
||||||
|
|
@ -1368,7 +1368,8 @@ class Graph:
|
||||||
vertex._finalize_build()
|
vertex._finalize_build()
|
||||||
if vertex.result is not None:
|
if vertex.result is not None:
|
||||||
vertex.result.used_frozen_result = True
|
vertex.result.used_frozen_result = True
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error finalizing build")
|
||||||
should_build = True
|
should_build = True
|
||||||
except KeyError:
|
except KeyError:
|
||||||
should_build = True
|
should_build = True
|
||||||
|
|
@ -1747,8 +1748,8 @@ class Graph:
|
||||||
if stop_component_id or start_component_id:
|
if stop_component_id or start_component_id:
|
||||||
try:
|
try:
|
||||||
first_layer = self.sort_vertices(stop_component_id, start_component_id)
|
first_layer = self.sort_vertices(stop_component_id, start_component_id)
|
||||||
except Exception as exc:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(exc)
|
logger.exception("Error sorting vertices")
|
||||||
first_layer = self.sort_vertices()
|
first_layer = self.sort_vertices()
|
||||||
else:
|
else:
|
||||||
first_layer = self.sort_vertices()
|
first_layer = self.sort_vertices()
|
||||||
|
|
|
||||||
|
|
@ -16,7 +16,7 @@ class GraphStateManager:
|
||||||
def __init__(self):
|
def __init__(self):
|
||||||
try:
|
try:
|
||||||
self.state_service: StateService = get_state_service()
|
self.state_service: StateService = get_state_service()
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.opt(exception=True).debug("Error getting state service. Defaulting to InMemoryStateService")
|
logger.opt(exception=True).debug("Error getting state service. Defaulting to InMemoryStateService")
|
||||||
from langflow.services.state.service import InMemoryStateService
|
from langflow.services.state.service import InMemoryStateService
|
||||||
|
|
||||||
|
|
@ -42,6 +42,6 @@ class GraphStateManager:
|
||||||
for callback in self.observers[key]:
|
for callback in self.observers[key]:
|
||||||
try:
|
try:
|
||||||
callback(key, new_state, append=True)
|
callback(key, new_state, append=True)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(f"Error in observer {callback} for key {key}")
|
logger.exception(f"Error in observer {callback} for key {key}")
|
||||||
logger.warning("Callbacks not implemented yet")
|
logger.warning("Callbacks not implemented yet")
|
||||||
|
|
|
||||||
|
|
@ -155,7 +155,7 @@ async def log_transaction(
|
||||||
with session_getter(get_db_service()) as session:
|
with session_getter(get_db_service()) as session:
|
||||||
inserted = crud_log_transaction(session, transaction)
|
inserted = crud_log_transaction(session, transaction)
|
||||||
logger.debug(f"Logged transaction: {inserted.id}")
|
logger.debug(f"Logged transaction: {inserted.id}")
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Error logging transaction")
|
logger.exception("Error logging transaction")
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -183,7 +183,7 @@ def log_vertex_build(
|
||||||
with session_getter(get_db_service()) as session:
|
with session_getter(get_db_service()) as session:
|
||||||
inserted = crud_log_vertex_build(session, vertex_build)
|
inserted = crud_log_vertex_build(session, vertex_build)
|
||||||
logger.debug(f"Logged vertex build: {inserted.build_id}")
|
logger.debug(f"Logged vertex build: {inserted.build_id}")
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Error logging vertex build")
|
logger.exception("Error logging vertex build")
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -371,7 +371,8 @@ class Vertex:
|
||||||
if field.get("type") == "code":
|
if field.get("type") == "code":
|
||||||
try:
|
try:
|
||||||
params[field_name] = ast.literal_eval(val) if val else None
|
params[field_name] = ast.literal_eval(val) if val else None
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug(f"Error evaluating code for {field_name}")
|
||||||
params[field_name] = val
|
params[field_name] = val
|
||||||
elif field.get("type") in ["dict", "NestedDict"]:
|
elif field.get("type") in ["dict", "NestedDict"]:
|
||||||
# When dict comes from the frontend it comes as a
|
# When dict comes from the frontend it comes as a
|
||||||
|
|
|
||||||
|
|
@ -382,7 +382,7 @@ def copy_profile_pictures():
|
||||||
shutil.copytree(origin, target, dirs_exist_ok=True)
|
shutil.copytree(origin, target, dirs_exist_ok=True)
|
||||||
logger.debug(f"Folder copied from '{origin}' to '{target}'")
|
logger.debug(f"Folder copied from '{origin}' to '{target}'")
|
||||||
|
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Error copying the folder")
|
logger.exception("Error copying the folder")
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -602,8 +602,8 @@ async def create_or_update_starter_projects(get_all_components_coro: Awaitable[d
|
||||||
updated_project_data = update_edges_with_latest_component_versions(updated_project_data)
|
updated_project_data = update_edges_with_latest_component_versions(updated_project_data)
|
||||||
try:
|
try:
|
||||||
Graph.from_payload(updated_project_data)
|
Graph.from_payload(updated_project_data)
|
||||||
except Exception as e:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(e)
|
logger.exception(f"Error loading project {project_name}")
|
||||||
if updated_project_data != project_data:
|
if updated_project_data != project_data:
|
||||||
project_data = updated_project_data
|
project_data = updated_project_data
|
||||||
# We also need to update the project data in the file
|
# We also need to update the project data in the file
|
||||||
|
|
|
||||||
|
|
@ -137,7 +137,7 @@ def update_params_with_load_from_db_fields(
|
||||||
except TypeError as exc:
|
except TypeError as exc:
|
||||||
raise exc
|
raise exc
|
||||||
|
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(f"Failed to get value for {field} from custom component. Setting it to None.")
|
logger.exception(f"Failed to get value for {field} from custom component. Setting it to None.")
|
||||||
params[field] = None
|
params[field] = None
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -96,7 +96,7 @@ def setup_llm_caching():
|
||||||
set_langchain_cache(settings_service.settings)
|
set_langchain_cache(settings_service.settings)
|
||||||
except ImportError:
|
except ImportError:
|
||||||
logger.warning(f"Could not import {settings_service.settings.cache_type}. ")
|
logger.warning(f"Could not import {settings_service.settings.cache_type}. ")
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.opt(exception=True).warning("Could not setup LLM caching.")
|
logger.opt(exception=True).warning("Could not setup LLM caching.")
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -111,7 +111,7 @@ def run_flow_from_json(
|
||||||
import nest_asyncio
|
import nest_asyncio
|
||||||
|
|
||||||
nest_asyncio.apply()
|
nest_asyncio.apply()
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.opt(exception=True).warning("Could not apply nest_asyncio")
|
logger.opt(exception=True).warning("Could not apply nest_asyncio")
|
||||||
if tweaks is None:
|
if tweaks is None:
|
||||||
tweaks = {}
|
tweaks = {}
|
||||||
|
|
|
||||||
|
|
@ -192,7 +192,7 @@ def configure(
|
||||||
rotation="10 MB", # Log rotation based on file size
|
rotation="10 MB", # Log rotation based on file size
|
||||||
serialize=True,
|
serialize=True,
|
||||||
)
|
)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Error setting up log file")
|
logger.exception("Error setting up log file")
|
||||||
|
|
||||||
if log_buffer.enabled():
|
if log_buffer.enabled():
|
||||||
|
|
|
||||||
|
|
@ -30,7 +30,7 @@ def get_langfuse_callback(trace_id):
|
||||||
try:
|
try:
|
||||||
trace = langfuse.trace(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:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Error initializing langfuse callback")
|
logger.exception("Error initializing langfuse callback")
|
||||||
|
|
||||||
return None
|
return None
|
||||||
|
|
|
||||||
|
|
@ -2,6 +2,7 @@ from collections.abc import Generator
|
||||||
from enum import Enum
|
from enum import Enum
|
||||||
|
|
||||||
from fastapi.encoders import jsonable_encoder
|
from fastapi.encoders import jsonable_encoder
|
||||||
|
from loguru import logger
|
||||||
from pydantic import BaseModel
|
from pydantic import BaseModel
|
||||||
|
|
||||||
from langflow.schema import Data
|
from langflow.schema import Data
|
||||||
|
|
@ -66,7 +67,8 @@ def post_process_raw(raw, artifact_type: str):
|
||||||
try:
|
try:
|
||||||
raw = jsonable_encoder(raw)
|
raw = jsonable_encoder(raw)
|
||||||
artifact_type = ArtifactType.OBJECT.value
|
artifact_type = ArtifactType.OBJECT.value
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error converting to json")
|
||||||
raw = "Built Successfully ✨"
|
raw = "Built Successfully ✨"
|
||||||
else:
|
else:
|
||||||
raw = "Built Successfully ✨"
|
raw = "Built Successfully ✨"
|
||||||
|
|
|
||||||
|
|
@ -5,6 +5,7 @@ from typing import TYPE_CHECKING, cast
|
||||||
from langchain_core.documents import Document
|
from langchain_core.documents import Document
|
||||||
from langchain_core.messages import AIMessage, BaseMessage, HumanMessage
|
from langchain_core.messages import AIMessage, BaseMessage, HumanMessage
|
||||||
from langchain_core.prompts.image import ImagePromptTemplate
|
from langchain_core.prompts.image import ImagePromptTemplate
|
||||||
|
from loguru import logger
|
||||||
from pydantic import BaseModel, model_serializer, model_validator
|
from pydantic import BaseModel, model_serializer, model_validator
|
||||||
|
|
||||||
from langflow.utils.constants import MESSAGE_SENDER_AI, MESSAGE_SENDER_USER
|
from langflow.utils.constants import MESSAGE_SENDER_AI, MESSAGE_SENDER_USER
|
||||||
|
|
@ -209,7 +210,8 @@ class Data(BaseModel):
|
||||||
try:
|
try:
|
||||||
data = {k: v.to_json() if hasattr(v, "to_json") else v for k, v in self.data.items()}
|
data = {k: v.to_json() if hasattr(v, "to_json") else v for k, v in self.data.items()}
|
||||||
return json.dumps(data, indent=4)
|
return json.dumps(data, indent=4)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error converting Data to JSON")
|
||||||
return str(self.data)
|
return str(self.data)
|
||||||
|
|
||||||
def __contains__(self, key):
|
def __contains__(self, key):
|
||||||
|
|
|
||||||
|
|
@ -2,6 +2,7 @@ from collections.abc import AsyncIterator, Generator, Iterator
|
||||||
from enum import Enum
|
from enum import Enum
|
||||||
from typing import Literal
|
from typing import Literal
|
||||||
|
|
||||||
|
from loguru import logger
|
||||||
from pydantic import BaseModel
|
from pydantic import BaseModel
|
||||||
from typing_extensions import TypedDict
|
from typing_extensions import TypedDict
|
||||||
|
|
||||||
|
|
@ -141,5 +142,6 @@ def recursive_serialize_or_str(obj):
|
||||||
# This a type BaseModel and not an instance of it
|
# This a type BaseModel and not an instance of it
|
||||||
return repr(obj)
|
return repr(obj)
|
||||||
return str(obj)
|
return str(obj)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug(f"Cannot serialize object {obj}")
|
||||||
return str(obj)
|
return str(obj)
|
||||||
|
|
|
||||||
|
|
@ -379,6 +379,7 @@ def decrypt_api_key(encrypted_api_key: str, settings_service=Depends(get_setting
|
||||||
if isinstance(encrypted_api_key, str):
|
if isinstance(encrypted_api_key, str):
|
||||||
try:
|
try:
|
||||||
decrypted_key = fernet.decrypt(encrypted_api_key.encode()).decode()
|
decrypted_key = fernet.decrypt(encrypted_api_key.encode()).decode()
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Failed to decrypt API key")
|
||||||
decrypted_key = fernet.decrypt(encrypted_api_key).decode()
|
decrypted_key = fernet.decrypt(encrypted_api_key).decode()
|
||||||
return decrypted_key
|
return decrypted_key
|
||||||
|
|
|
||||||
|
|
@ -2,6 +2,7 @@ from datetime import datetime, timezone
|
||||||
from uuid import UUID
|
from uuid import UUID
|
||||||
|
|
||||||
from fastapi import Depends, HTTPException, status
|
from fastapi import Depends, HTTPException, status
|
||||||
|
from loguru import logger
|
||||||
from sqlalchemy.exc import IntegrityError
|
from sqlalchemy.exc import IntegrityError
|
||||||
from sqlalchemy.orm.attributes import flag_modified
|
from sqlalchemy.orm.attributes import flag_modified
|
||||||
from sqlmodel import Session, select
|
from sqlmodel import Session, select
|
||||||
|
|
@ -53,5 +54,5 @@ def update_user_last_login_at(user_id: UUID, db: Session = Depends(get_session))
|
||||||
user_data = UserUpdate(last_login_at=datetime.now(timezone.utc))
|
user_data = UserUpdate(last_login_at=datetime.now(timezone.utc))
|
||||||
user = get_user_by_id(db, user_id)
|
user = get_user_by_id(db, user_id)
|
||||||
return update_user(user, user_data, db)
|
return update_user(user, user_data, db)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
pass
|
logger.opt(exception=True).debug("Error updating user last login at")
|
||||||
|
|
|
||||||
|
|
@ -187,8 +187,8 @@ class DatabaseService(Service):
|
||||||
# so we need to catch it
|
# so we need to catch it
|
||||||
try:
|
try:
|
||||||
session.exec(text("SELECT * FROM alembic_version"))
|
session.exec(text("SELECT * FROM alembic_version"))
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.info("Alembic not initialized")
|
logger.opt(exception=True).info("Alembic not initialized")
|
||||||
should_initialize_alembic = True
|
should_initialize_alembic = True
|
||||||
|
|
||||||
if should_initialize_alembic:
|
if should_initialize_alembic:
|
||||||
|
|
@ -206,7 +206,8 @@ class DatabaseService(Service):
|
||||||
try:
|
try:
|
||||||
buffer.write(f"{datetime.now(tz=timezone.utc).astimezone().isoformat()}: Checking migrations\n")
|
buffer.write(f"{datetime.now(tz=timezone.utc).astimezone().isoformat()}: Checking migrations\n")
|
||||||
command.check(alembic_cfg)
|
command.check(alembic_cfg)
|
||||||
except Exception as exc:
|
except Exception as exc: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error checking migrations")
|
||||||
if isinstance(exc, util.exc.CommandError | util.exc.AutogenerateDiffsDetected):
|
if isinstance(exc, util.exc.CommandError | util.exc.AutogenerateDiffsDetected):
|
||||||
command.upgrade(alembic_cfg, "head")
|
command.upgrade(alembic_cfg, "head")
|
||||||
time.sleep(3)
|
time.sleep(3)
|
||||||
|
|
@ -215,7 +216,7 @@ class DatabaseService(Service):
|
||||||
buffer.write(f"{datetime.now(tz=timezone.utc).astimezone()}: Checking migrations\n")
|
buffer.write(f"{datetime.now(tz=timezone.utc).astimezone()}: Checking migrations\n")
|
||||||
command.check(alembic_cfg)
|
command.check(alembic_cfg)
|
||||||
except util.exc.AutogenerateDiffsDetected as exc:
|
except util.exc.AutogenerateDiffsDetected as exc:
|
||||||
logger.exception("AutogenerateDiffsDetected")
|
logger.exception("Error checking migrations")
|
||||||
if not fix:
|
if not fix:
|
||||||
msg = f"There's a mismatch between the models and the database.\n{exc}"
|
msg = f"There's a mismatch between the models and the database.\n{exc}"
|
||||||
raise RuntimeError(msg) from exc
|
raise RuntimeError(msg) from exc
|
||||||
|
|
@ -313,7 +314,7 @@ class DatabaseService(Service):
|
||||||
with self.with_session() as session:
|
with self.with_session() as session:
|
||||||
teardown_superuser(settings_service, session)
|
teardown_superuser(settings_service, session)
|
||||||
|
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Error tearing down database")
|
logger.exception("Error tearing down database")
|
||||||
|
|
||||||
self.engine.dispose()
|
self.engine.dispose()
|
||||||
|
|
|
||||||
|
|
@ -34,7 +34,7 @@ class ServiceManager:
|
||||||
for factory in self.get_factories():
|
for factory in self.get_factories():
|
||||||
try:
|
try:
|
||||||
self.register_factory(factory)
|
self.register_factory(factory)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(f"Error initializing {factory}")
|
logger.exception(f"Error initializing {factory}")
|
||||||
|
|
||||||
def register_factory(
|
def register_factory(
|
||||||
|
|
@ -111,7 +111,7 @@ class ServiceManager:
|
||||||
result = service.teardown()
|
result = service.teardown()
|
||||||
if asyncio.iscoroutine(result):
|
if asyncio.iscoroutine(result):
|
||||||
await result
|
await result
|
||||||
except Exception as exc:
|
except Exception as exc: # noqa: BLE001
|
||||||
logger.exception(exc)
|
logger.exception(exc)
|
||||||
self.services = {}
|
self.services = {}
|
||||||
self.factories = {}
|
self.factories = {}
|
||||||
|
|
|
||||||
|
|
@ -80,7 +80,7 @@ class LangfusePlugin(CallbackPlugin):
|
||||||
if trace:
|
if trace:
|
||||||
return trace.getNewHandler()
|
return trace.getNewHandler()
|
||||||
|
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Error initializing langfuse callback")
|
logger.exception("Error initializing langfuse callback")
|
||||||
|
|
||||||
return None
|
return None
|
||||||
|
|
|
||||||
|
|
@ -40,7 +40,7 @@ class PluginService(Service):
|
||||||
and attr not in [CallbackPlugin, BasePlugin]
|
and attr not in [CallbackPlugin, BasePlugin]
|
||||||
):
|
):
|
||||||
self.register_plugin(plugin_name, attr())
|
self.register_plugin(plugin_name, attr())
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(f"Error loading plugin {plugin_name}")
|
logger.exception(f"Error loading plugin {plugin_name}")
|
||||||
|
|
||||||
def register_plugin(self, plugin_name, plugin_instance):
|
def register_plugin(self, plugin_name, plugin_instance):
|
||||||
|
|
|
||||||
|
|
@ -277,7 +277,7 @@ class Settings(BaseSettings):
|
||||||
logger.debug("Copying existing database to new location")
|
logger.debug("Copying existing database to new location")
|
||||||
copy2(f"./{db_file_name}", new_path)
|
copy2(f"./{db_file_name}", new_path)
|
||||||
logger.debug(f"Copied existing database to {new_path}")
|
logger.debug(f"Copied existing database to {new_path}")
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Failed to copy database, using default path")
|
logger.exception("Failed to copy database, using default path")
|
||||||
new_path = f"./{db_file_name}"
|
new_path = f"./{db_file_name}"
|
||||||
else:
|
else:
|
||||||
|
|
|
||||||
|
|
@ -33,7 +33,7 @@ def write_secret_to_file(path: Path, value: str) -> None:
|
||||||
f.write(value.encode("utf-8"))
|
f.write(value.encode("utf-8"))
|
||||||
try:
|
try:
|
||||||
set_secure_permissions(path)
|
set_secure_permissions(path)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Failed to set secure permissions on secret key")
|
logger.exception("Failed to set secure permissions on secret key")
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -2,6 +2,7 @@ import time
|
||||||
from collections.abc import Callable
|
from collections.abc import Callable
|
||||||
|
|
||||||
import socketio
|
import socketio
|
||||||
|
from loguru import logger
|
||||||
from sqlmodel import select
|
from sqlmodel import select
|
||||||
|
|
||||||
from langflow.api.utils import format_elapsed_time
|
from langflow.api.utils import format_elapsed_time
|
||||||
|
|
@ -35,7 +36,8 @@ async def get_vertices(sio, sid, flow_id, chat_service):
|
||||||
# Emit the vertices to the client
|
# Emit the vertices to the client
|
||||||
await sio.emit("vertices_order", data=vertices, to=sid)
|
await sio.emit("vertices_order", data=vertices, to=sid)
|
||||||
|
|
||||||
except Exception as exc:
|
except Exception as exc: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error getting vertices")
|
||||||
await sio.emit("error", data=str(exc), to=sid)
|
await sio.emit("error", data=str(exc), to=sid)
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -80,7 +82,8 @@ async def build_vertex(
|
||||||
duration=duration,
|
duration=duration,
|
||||||
timedelta=timedelta,
|
timedelta=timedelta,
|
||||||
)
|
)
|
||||||
except Exception as exc:
|
except Exception as exc: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error building vertex")
|
||||||
params = str(exc)
|
params = str(exc)
|
||||||
valid = False
|
valid = False
|
||||||
result_dict = ResultDataResponse(results={})
|
result_dict = ResultDataResponse(results={})
|
||||||
|
|
@ -99,5 +102,6 @@ async def build_vertex(
|
||||||
response = VertexBuildResponse(valid=valid, params=params, id=vertex.id, data=result_dict)
|
response = VertexBuildResponse(valid=valid, params=params, id=vertex.id, data=result_dict)
|
||||||
await sio.emit("vertex_build", data=response.model_dump(), to=sid)
|
await sio.emit("vertex_build", data=response.model_dump(), to=sid)
|
||||||
|
|
||||||
except Exception as exc:
|
except Exception as exc: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error building vertex")
|
||||||
await sio.emit("error", data=str(exc), to=sid)
|
await sio.emit("error", data=str(exc), to=sid)
|
||||||
|
|
|
||||||
|
|
@ -69,6 +69,6 @@ class InMemoryStateService(StateService):
|
||||||
for callback in self.observers[key]:
|
for callback in self.observers[key]:
|
||||||
try:
|
try:
|
||||||
callback(key, new_state, append=True)
|
callback(key, new_state, append=True)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(f"Error in observer {callback} for key {key}")
|
logger.exception(f"Error in observer {callback} for key {key}")
|
||||||
logger.warning("Callbacks not implemented yet")
|
logger.warning("Callbacks not implemented yet")
|
||||||
|
|
|
||||||
|
|
@ -161,7 +161,7 @@ class StoreService(Service):
|
||||||
return response.json()
|
return response.json()
|
||||||
except HTTPError as exc:
|
except HTTPError as exc:
|
||||||
raise exc
|
raise exc
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.opt(exception=True).debug("Webhook failed")
|
logger.opt(exception=True).debug("Webhook failed")
|
||||||
|
|
||||||
def build_tags_filter(self, tags: list[str]):
|
def build_tags_filter(self, tags: list[str]):
|
||||||
|
|
@ -587,7 +587,8 @@ class StoreService(Service):
|
||||||
)
|
)
|
||||||
authorized = True
|
authorized = True
|
||||||
result = updated_result
|
result = updated_result
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error updating components with user data")
|
||||||
# If we get an error here, it means the user is not authorized
|
# If we get an error here, it means the user is not authorized
|
||||||
authorized = False
|
authorized = False
|
||||||
return ListComponentResponseModel(results=result, authorized=authorized, count=comp_count)
|
return ListComponentResponseModel(results=result, authorized=authorized, count=comp_count)
|
||||||
|
|
|
||||||
|
|
@ -1,6 +1,7 @@
|
||||||
from typing import TYPE_CHECKING
|
from typing import TYPE_CHECKING
|
||||||
|
|
||||||
import httpx
|
import httpx
|
||||||
|
from loguru import logger
|
||||||
|
|
||||||
if TYPE_CHECKING:
|
if TYPE_CHECKING:
|
||||||
from langflow.services.store.schema import ListComponentResponse
|
from langflow.services.store.schema import ListComponentResponse
|
||||||
|
|
@ -47,7 +48,8 @@ def get_lf_version_from_pypi():
|
||||||
if response.status_code != httpx.codes.OK:
|
if response.status_code != httpx.codes.OK:
|
||||||
return None
|
return None
|
||||||
return response.json()["info"]["version"]
|
return response.json()["info"]["version"]
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error getting the latest version of langflow from PyPI")
|
||||||
return None
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -37,7 +37,7 @@ class AnyIOTaskResult:
|
||||||
async def run(self, func, *args, **kwargs):
|
async def run(self, func, *args, **kwargs):
|
||||||
try:
|
try:
|
||||||
self._result = await func(*args, **kwargs)
|
self._result = await func(*args, **kwargs)
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
self._exception = e
|
self._exception = e
|
||||||
self._traceback = e.__traceback__
|
self._traceback = e.__traceback__
|
||||||
finally:
|
finally:
|
||||||
|
|
@ -72,7 +72,7 @@ class AnyIOBackend(TaskBackend):
|
||||||
self.tasks[task_id] = task_result
|
self.tasks[task_id] = task_result
|
||||||
logger.info(f"Task {task_id} started.")
|
logger.info(f"Task {task_id} started.")
|
||||||
return task_id, task_result
|
return task_id, task_result
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("An error occurred while launching the task")
|
logger.exception("An error occurred while launching the task")
|
||||||
return None, None
|
return None, None
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -20,7 +20,7 @@ def check_celery_availability():
|
||||||
|
|
||||||
status = get_celery_worker_status(celery_app)
|
status = get_celery_worker_status(celery_app)
|
||||||
logger.debug(f"Celery status: {status}")
|
logger.debug(f"Celery status: {status}")
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.opt(exception=True).debug("Celery not available")
|
logger.opt(exception=True).debug("Celery not available")
|
||||||
status = {"availability": None}
|
status = {"availability": None}
|
||||||
return status
|
return status
|
||||||
|
|
|
||||||
|
|
@ -51,7 +51,7 @@ class TelemetryService(Service):
|
||||||
func, payload, path = await self.telemetry_queue.get()
|
func, payload, path = await self.telemetry_queue.get()
|
||||||
try:
|
try:
|
||||||
await func(payload, path)
|
await func(payload, path)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Error sending telemetry data")
|
logger.exception("Error sending telemetry data")
|
||||||
finally:
|
finally:
|
||||||
self.telemetry_queue.task_done()
|
self.telemetry_queue.task_done()
|
||||||
|
|
@ -75,7 +75,7 @@ class TelemetryService(Service):
|
||||||
logger.exception("HTTP error occurred")
|
logger.exception("HTTP error occurred")
|
||||||
except httpx.RequestError:
|
except httpx.RequestError:
|
||||||
logger.exception("Request error occurred")
|
logger.exception("Request error occurred")
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Unexpected error occurred")
|
logger.exception("Unexpected error occurred")
|
||||||
|
|
||||||
async def log_package_run(self, payload: RunPayload):
|
async def log_package_run(self, payload: RunPayload):
|
||||||
|
|
@ -120,7 +120,7 @@ class TelemetryService(Service):
|
||||||
self._start_time = datetime.now(timezone.utc)
|
self._start_time = datetime.now(timezone.utc)
|
||||||
self.worker_task = asyncio.create_task(self.telemetry_worker())
|
self.worker_task = asyncio.create_task(self.telemetry_worker())
|
||||||
asyncio.create_task(self.log_package_version())
|
asyncio.create_task(self.log_package_version())
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Error starting telemetry service")
|
logger.exception("Error starting telemetry service")
|
||||||
|
|
||||||
async def flush(self):
|
async def flush(self):
|
||||||
|
|
@ -128,7 +128,7 @@ class TelemetryService(Service):
|
||||||
return
|
return
|
||||||
try:
|
try:
|
||||||
await self.telemetry_queue.join()
|
await self.telemetry_queue.join()
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Error flushing logs")
|
logger.exception("Error flushing logs")
|
||||||
|
|
||||||
async def stop(self):
|
async def stop(self):
|
||||||
|
|
@ -144,7 +144,7 @@ class TelemetryService(Service):
|
||||||
with contextlib.suppress(asyncio.CancelledError):
|
with contextlib.suppress(asyncio.CancelledError):
|
||||||
await self.worker_task
|
await self.worker_task
|
||||||
await self.client.aclose()
|
await self.client.aclose()
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Error stopping tracing service")
|
logger.exception("Error stopping tracing service")
|
||||||
|
|
||||||
async def teardown(self):
|
async def teardown(self):
|
||||||
|
|
|
||||||
|
|
@ -57,7 +57,7 @@ class LangFuseTracer(BaseTracer):
|
||||||
logger.exception("Could not import langfuse. Please install it with `pip install langfuse`.")
|
logger.exception("Could not import langfuse. Please install it with `pip install langfuse`.")
|
||||||
return False
|
return False
|
||||||
|
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.opt(exception=True).debug("Error setting up LangSmith tracer")
|
logger.opt(exception=True).debug("Error setting up LangSmith tracer")
|
||||||
return False
|
return False
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -41,7 +41,7 @@ class LangSmithTracer(BaseTracer):
|
||||||
)
|
)
|
||||||
self._run_tree.add_event({"name": "Start", "time": datetime.now(timezone.utc).isoformat()})
|
self._run_tree.add_event({"name": "Start", "time": datetime.now(timezone.utc).isoformat()})
|
||||||
self._children: dict[str, RunTree] = {}
|
self._children: dict[str, RunTree] = {}
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.opt(exception=True).debug("Error setting up LangSmith tracer")
|
logger.opt(exception=True).debug("Error setting up LangSmith tracer")
|
||||||
self._ready = False
|
self._ready = False
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -46,7 +46,7 @@ class LangWatchTracer(BaseTracer):
|
||||||
name=name_without_id,
|
name=name_without_id,
|
||||||
type="workflow",
|
type="workflow",
|
||||||
)
|
)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.opt(exception=True).debug("Error setting up LangWatch tracer")
|
logger.opt(exception=True).debug("Error setting up LangWatch tracer")
|
||||||
self._ready = False
|
self._ready = False
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -63,7 +63,7 @@ class TracingService(Service):
|
||||||
log_func, args = await self.logs_queue.get()
|
log_func, args = await self.logs_queue.get()
|
||||||
try:
|
try:
|
||||||
await log_func(*args)
|
await log_func(*args)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Error processing log")
|
logger.exception("Error processing log")
|
||||||
finally:
|
finally:
|
||||||
self.logs_queue.task_done()
|
self.logs_queue.task_done()
|
||||||
|
|
@ -74,13 +74,13 @@ class TracingService(Service):
|
||||||
try:
|
try:
|
||||||
self.running = True
|
self.running = True
|
||||||
self.worker_task = asyncio.create_task(self.log_worker())
|
self.worker_task = asyncio.create_task(self.log_worker())
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Error starting tracing service")
|
logger.exception("Error starting tracing service")
|
||||||
|
|
||||||
async def flush(self):
|
async def flush(self):
|
||||||
try:
|
try:
|
||||||
await self.logs_queue.join()
|
await self.logs_queue.join()
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Error flushing logs")
|
logger.exception("Error flushing logs")
|
||||||
|
|
||||||
async def stop(self):
|
async def stop(self):
|
||||||
|
|
@ -94,7 +94,7 @@ class TracingService(Service):
|
||||||
self.worker_task.cancel()
|
self.worker_task.cancel()
|
||||||
self.worker_task = None
|
self.worker_task = None
|
||||||
|
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Error stopping tracing service")
|
logger.exception("Error stopping tracing service")
|
||||||
|
|
||||||
def _reset_io(self):
|
def _reset_io(self):
|
||||||
|
|
@ -109,7 +109,7 @@ class TracingService(Service):
|
||||||
self._initialize_langsmith_tracer()
|
self._initialize_langsmith_tracer()
|
||||||
self._initialize_langwatch_tracer()
|
self._initialize_langwatch_tracer()
|
||||||
self._initialize_langfuse_tracer()
|
self._initialize_langfuse_tracer()
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.opt(exception=True).debug("Error initializing tracers")
|
logger.opt(exception=True).debug("Error initializing tracers")
|
||||||
|
|
||||||
def _initialize_langsmith_tracer(self):
|
def _initialize_langsmith_tracer(self):
|
||||||
|
|
@ -166,7 +166,7 @@ class TracingService(Service):
|
||||||
continue
|
continue
|
||||||
try:
|
try:
|
||||||
tracer.add_trace(trace_id, trace_name, trace_type, inputs, metadata, vertex)
|
tracer.add_trace(trace_id, trace_name, trace_type, inputs, metadata, vertex)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(f"Error starting trace {trace_name}")
|
logger.exception(f"Error starting trace {trace_name}")
|
||||||
|
|
||||||
def _end_traces(self, trace_id: str, trace_name: str, error: Exception | None = None):
|
def _end_traces(self, trace_id: str, trace_name: str, error: Exception | None = None):
|
||||||
|
|
@ -181,7 +181,7 @@ class TracingService(Service):
|
||||||
error=error,
|
error=error,
|
||||||
logs=self._logs[trace_name],
|
logs=self._logs[trace_name],
|
||||||
)
|
)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(f"Error ending trace {trace_name}")
|
logger.exception(f"Error ending trace {trace_name}")
|
||||||
|
|
||||||
def _end_all_traces(self, outputs: dict, error: Exception | None = None):
|
def _end_all_traces(self, outputs: dict, error: Exception | None = None):
|
||||||
|
|
@ -190,7 +190,7 @@ class TracingService(Service):
|
||||||
continue
|
continue
|
||||||
try:
|
try:
|
||||||
tracer.end(self.inputs, outputs=self.outputs, error=error, metadata=outputs)
|
tracer.end(self.inputs, outputs=self.outputs, error=error, metadata=outputs)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception("Error ending all traces")
|
logger.exception("Error ending all traces")
|
||||||
|
|
||||||
async def end(self, outputs: dict, error: Exception | None = None):
|
async def end(self, outputs: dict, error: Exception | None = None):
|
||||||
|
|
|
||||||
|
|
@ -50,13 +50,14 @@ def get_or_create_super_user(session: Session, username, password, is_default):
|
||||||
logger.debug("Creating superuser.")
|
logger.debug("Creating superuser.")
|
||||||
try:
|
try:
|
||||||
return create_super_user(username, password, db=session)
|
return create_super_user(username, password, db=session)
|
||||||
except Exception as exc:
|
except Exception as exc: # noqa: BLE001
|
||||||
if "UNIQUE constraint failed: user.username" in str(exc):
|
if "UNIQUE constraint failed: user.username" in str(exc):
|
||||||
# This is to deal with workers running this
|
# This is to deal with workers running this
|
||||||
# at startup and trying to create the superuser
|
# at startup and trying to create the superuser
|
||||||
# at the same time.
|
# at the same time.
|
||||||
logger.debug("Superuser already exists.")
|
logger.opt(exception=True).debug("Superuser already exists.")
|
||||||
return None
|
return None
|
||||||
|
logger.opt(exception=True).debug("Error creating superuser.")
|
||||||
|
|
||||||
|
|
||||||
def setup_superuser(settings_service, session: Session):
|
def setup_superuser(settings_service, session: Session):
|
||||||
|
|
@ -117,13 +118,13 @@ async def teardown_services():
|
||||||
"""
|
"""
|
||||||
try:
|
try:
|
||||||
teardown_superuser(get_settings_service(), next(get_session()))
|
teardown_superuser(get_settings_service(), next(get_session()))
|
||||||
except Exception as exc:
|
except Exception as exc: # noqa: BLE001
|
||||||
logger.exception(exc)
|
logger.exception(exc)
|
||||||
try:
|
try:
|
||||||
from langflow.services.manager import service_manager
|
from langflow.services.manager import service_manager
|
||||||
|
|
||||||
await service_manager.teardown()
|
await service_manager.teardown()
|
||||||
except Exception as exc:
|
except Exception as exc: # noqa: BLE001
|
||||||
logger.exception(exc)
|
logger.exception(exc)
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -47,7 +47,7 @@ class KubernetesSecretService(VariableService, Service):
|
||||||
name=secret_name,
|
name=secret_name,
|
||||||
data=variables,
|
data=variables,
|
||||||
)
|
)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(f"Error creating {var} variable")
|
logger.exception(f"Error creating {var} variable")
|
||||||
|
|
||||||
else:
|
else:
|
||||||
|
|
|
||||||
|
|
@ -62,7 +62,7 @@ class DatabaseVariableService(VariableService, Service):
|
||||||
_type=CREDENTIAL_TYPE,
|
_type=CREDENTIAL_TYPE,
|
||||||
session=session,
|
session=session,
|
||||||
)
|
)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(f"Error creating {var} variable")
|
logger.exception(f"Error creating {var} variable")
|
||||||
|
|
||||||
else:
|
else:
|
||||||
|
|
|
||||||
|
|
@ -22,7 +22,7 @@ def run_in_thread(coro):
|
||||||
nonlocal result, exception
|
nonlocal result, exception
|
||||||
try:
|
try:
|
||||||
result = asyncio.run(coro)
|
result = asyncio.run(coro)
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
exception = e
|
exception = e
|
||||||
|
|
||||||
thread = threading.Thread(target=target)
|
thread = threading.Thread(target=target)
|
||||||
|
|
|
||||||
|
|
@ -51,7 +51,8 @@ def build_template_from_function(name: str, type_to_loader_dict: dict, add_funct
|
||||||
variables[class_field_items]["default"] = get_default_factory(
|
variables[class_field_items]["default"] = get_default_factory(
|
||||||
module=_class.__base__.__module__, function=value_
|
module=_class.__base__.__module__, function=value_
|
||||||
)
|
)
|
||||||
except Exception:
|
except Exception: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug(f"Error getting default factory for {value_}")
|
||||||
variables[class_field_items]["default"] = None
|
variables[class_field_items]["default"] = None
|
||||||
elif name_ != "name":
|
elif name_ != "name":
|
||||||
variables[class_field_items][name_] = value_
|
variables[class_field_items][name_] = value_
|
||||||
|
|
|
||||||
|
|
@ -4,6 +4,7 @@ import importlib
|
||||||
from types import FunctionType
|
from types import FunctionType
|
||||||
from typing import Optional, Union
|
from typing import Optional, Union
|
||||||
|
|
||||||
|
from loguru import logger
|
||||||
from pydantic import ValidationError
|
from pydantic import ValidationError
|
||||||
|
|
||||||
from langflow.field_typing.constants import CUSTOM_COMPONENT_SUPPORTED_TYPES
|
from langflow.field_typing.constants import CUSTOM_COMPONENT_SUPPORTED_TYPES
|
||||||
|
|
@ -25,7 +26,8 @@ def validate_code(code):
|
||||||
# Parse the code string into an abstract syntax tree (AST)
|
# Parse the code string into an abstract syntax tree (AST)
|
||||||
try:
|
try:
|
||||||
tree = ast.parse(code)
|
tree = ast.parse(code)
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error parsing code")
|
||||||
errors["function"]["errors"].append(str(e))
|
errors["function"]["errors"].append(str(e))
|
||||||
return errors
|
return errors
|
||||||
|
|
||||||
|
|
@ -48,7 +50,8 @@ def validate_code(code):
|
||||||
code_obj = compile(ast.Module(body=[node], type_ignores=[]), "<string>", "exec")
|
code_obj = compile(ast.Module(body=[node], type_ignores=[]), "<string>", "exec")
|
||||||
try:
|
try:
|
||||||
exec(code_obj)
|
exec(code_obj)
|
||||||
except Exception as e:
|
except Exception as e: # noqa: BLE001
|
||||||
|
logger.opt(exception=True).debug("Error executing function code")
|
||||||
errors["function"]["errors"].append(str(e))
|
errors["function"]["errors"].append(str(e))
|
||||||
|
|
||||||
# Return the errors dictionary
|
# Return the errors dictionary
|
||||||
|
|
|
||||||
|
|
@ -82,8 +82,8 @@ def fetch_latest_version(package_name: str, include_prerelease: bool) -> str | N
|
||||||
return None # Handle case where no valid versions are found
|
return None # Handle case where no valid versions are found
|
||||||
return max(valid_versions, key=lambda v: pkg_version.parse(v))
|
return max(valid_versions, key=lambda v: pkg_version.parse(v))
|
||||||
|
|
||||||
except Exception as e:
|
except Exception: # noqa: BLE001
|
||||||
logger.exception(e)
|
logger.exception("Error fetching latest version")
|
||||||
return None
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -56,7 +56,6 @@ ignore = [
|
||||||
# Rules that are TODOs
|
# Rules that are TODOs
|
||||||
"ANN",
|
"ANN",
|
||||||
"ARG",
|
"ARG",
|
||||||
"BLE",
|
|
||||||
"D",
|
"D",
|
||||||
"DOC",
|
"DOC",
|
||||||
"EXE",
|
"EXE",
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue