Merge branch 'main' into feat/firecrawl-integration
This commit is contained in:
commit
9373749163
44 changed files with 773 additions and 518 deletions
|
|
@ -4,12 +4,12 @@ import sys
|
|||
import time
|
||||
import warnings
|
||||
from pathlib import Path
|
||||
from typing import Optional
|
||||
from typing import Any, Callable, Optional
|
||||
|
||||
import click
|
||||
import httpx
|
||||
import typer
|
||||
from dotenv import load_dotenv
|
||||
from dotenv import dotenv_values, load_dotenv
|
||||
from multiprocess import Process, cpu_count # type: ignore
|
||||
from packaging import version as pkg_version
|
||||
from rich import box
|
||||
|
|
@ -130,6 +130,29 @@ def run(
|
|||
|
||||
if env_file:
|
||||
load_dotenv(env_file, override=True)
|
||||
env_vars = dotenv_values(env_file)
|
||||
|
||||
# Define a mapping of environment variables to their corresponding variables and types
|
||||
env_var_mapping: dict[str, tuple[str, type | Callable[[Any], bool]]] = {
|
||||
"LANGFLOW_HOST": ("host", str),
|
||||
"LANGFLOW_PORT": ("port", int),
|
||||
"LANGFLOW_WORKERS": ("workers", int),
|
||||
"LANGFLOW_WORKER_TIMEOUT": ("timeout", int),
|
||||
"LANGFLOW_COMPONENTS_PATH": ("components_path", Path),
|
||||
"LANGFLOW_LOG_LEVEL": ("log_level", str),
|
||||
"LANGFLOW_LOG_FILE": ("log_file", Path),
|
||||
"LANGFLOW_LANGCHAIN_CACHE": ("cache", str),
|
||||
"LANGFLOW_FRONTEND_PATH": ("path", str),
|
||||
"LANGFLOW_OPEN_BROWSER": ("open_browser", lambda x: x.lower() == "true"),
|
||||
"LANGFLOW_REMOVE_API_KEYS": ("remove_api_keys", lambda x: x.lower() == "true"),
|
||||
"LANGFLOW_BACKEND_ONLY": ("backend_only", lambda x: x.lower() == "true"),
|
||||
"LANGFLOW_STORE": ("store", lambda x: x.lower() == "true"),
|
||||
}
|
||||
|
||||
# Update variables based on environment variables
|
||||
for env_var, (var_name, var_type) in env_var_mapping.items():
|
||||
if env_var in env_vars:
|
||||
locals()[var_name] = var_type(env_vars[env_var])
|
||||
|
||||
update_settings(
|
||||
dev=dev,
|
||||
|
|
|
|||
|
|
@ -121,9 +121,9 @@ async def retrieve_vertices_order(
|
|||
background_tasks.add_task(
|
||||
telemetry_service.log_package_playground,
|
||||
PlaygroundPayload(
|
||||
seconds=int(time.perf_counter() - start_time),
|
||||
componentCount=components_count,
|
||||
success=True,
|
||||
playgroundSeconds=int(time.perf_counter() - start_time),
|
||||
playgroundComponentCount=components_count,
|
||||
playgroundSuccess=True,
|
||||
),
|
||||
)
|
||||
return VerticesOrderResponse(ids=first_layer, run_id=graph._run_id, vertices_to_run=vertices_to_run)
|
||||
|
|
@ -131,10 +131,10 @@ async def retrieve_vertices_order(
|
|||
background_tasks.add_task(
|
||||
telemetry_service.log_package_playground,
|
||||
PlaygroundPayload(
|
||||
seconds=int(time.perf_counter() - start_time),
|
||||
componentCount=components_count,
|
||||
success=False,
|
||||
errorMessage=str(exc),
|
||||
playgroundSeconds=int(time.perf_counter() - start_time),
|
||||
playgroundComponentCount=components_count,
|
||||
playgroundSuccess=False,
|
||||
playgroundErrorMessage=str(exc),
|
||||
),
|
||||
)
|
||||
if "stream or streaming set to True" in str(exc):
|
||||
|
|
@ -280,10 +280,10 @@ async def build_vertex(
|
|||
background_tasks.add_task(
|
||||
telemetry_service.log_package_component,
|
||||
ComponentPayload(
|
||||
name=vertex_id,
|
||||
seconds=int(time.perf_counter() - start_time),
|
||||
success=valid,
|
||||
errorMessage=params,
|
||||
componentName=vertex_id,
|
||||
componentSeconds=int(time.perf_counter() - start_time),
|
||||
componentSuccess=valid,
|
||||
componentErrorMessage=params,
|
||||
),
|
||||
)
|
||||
return build_response
|
||||
|
|
@ -291,10 +291,10 @@ async def build_vertex(
|
|||
background_tasks.add_task(
|
||||
telemetry_service.log_package_component,
|
||||
ComponentPayload(
|
||||
name=vertex_id,
|
||||
seconds=int(time.perf_counter() - start_time),
|
||||
success=False,
|
||||
errorMessage=str(exc),
|
||||
componentName=vertex_id,
|
||||
componentSeconds=int(time.perf_counter() - start_time),
|
||||
componentSuccess=False,
|
||||
componentErrorMessage=str(exc),
|
||||
),
|
||||
)
|
||||
logger.error(f"Error building Component:\n\n{exc}")
|
||||
|
|
|
|||
|
|
@ -116,11 +116,29 @@ async def simple_run_flow(
|
|||
return RunResponse(outputs=task_result, session_id=session_id)
|
||||
|
||||
except sa.exc.StatementError as exc:
|
||||
# StatementError('(builtins.ValueError) badly formed hexadecimal UUID string')
|
||||
if "badly formed hexadecimal UUID string" in str(exc):
|
||||
logger.error(f"Flow ID {flow_id_str} is not a valid UUID")
|
||||
# This means the Flow ID is not a valid UUID which means it can't find the flow
|
||||
raise ValueError(str(exc)) from exc
|
||||
raise ValueError(str(exc)) from exc
|
||||
|
||||
|
||||
async def simple_run_flow_task(
|
||||
flow: Flow,
|
||||
input_request: SimplifiedAPIRequest,
|
||||
stream: bool = False,
|
||||
api_key_user: Optional[User] = None,
|
||||
):
|
||||
"""
|
||||
Run a flow task as a BackgroundTask, therefore it should not throw exceptions.
|
||||
"""
|
||||
try:
|
||||
result = await simple_run_flow(
|
||||
flow=flow,
|
||||
input_request=input_request,
|
||||
stream=stream,
|
||||
api_key_user=api_key_user,
|
||||
)
|
||||
return result
|
||||
|
||||
except Exception as exc:
|
||||
logger.exception(f"Error running flow {flow.id} task: {exc}")
|
||||
|
||||
|
||||
@router.post("/run/{flow_id_or_name}", response_model=RunResponse, response_model_exclude_none=True)
|
||||
|
|
@ -191,7 +209,7 @@ async def simplified_run_flow(
|
|||
end_time = time.perf_counter()
|
||||
background_tasks.add_task(
|
||||
telemetry_service.log_package_run,
|
||||
RunPayload(IsWebhook=False, seconds=int(end_time - start_time), success=True, errorMessage=""),
|
||||
RunPayload(runIsWebhook=False, runSeconds=int(end_time - start_time), runSuccess=True, runErrorMessage=""),
|
||||
)
|
||||
return result
|
||||
|
||||
|
|
@ -199,7 +217,9 @@ async def simplified_run_flow(
|
|||
end_time = time.perf_counter()
|
||||
background_tasks.add_task(
|
||||
telemetry_service.log_package_run,
|
||||
RunPayload(IsWebhook=False, seconds=int(end_time - start_time), success=False, errorMessage=str(exc)),
|
||||
RunPayload(
|
||||
runIsWebhook=False, runSeconds=int(end_time - start_time), runSuccess=False, runErrorMessage=str(exc)
|
||||
),
|
||||
)
|
||||
if "badly formed hexadecimal UUID string" in str(exc):
|
||||
# This means the Flow ID is not a valid UUID which means it can't find the flow
|
||||
|
|
@ -213,7 +233,9 @@ async def simplified_run_flow(
|
|||
logger.exception(exc)
|
||||
background_tasks.add_task(
|
||||
telemetry_service.log_package_run,
|
||||
RunPayload(IsWebhook=False, seconds=int(end_time - start_time), success=False, errorMessage=str(exc)),
|
||||
RunPayload(
|
||||
runIsWebhook=False, runSeconds=int(end_time - start_time), runSuccess=False, runErrorMessage=str(exc)
|
||||
),
|
||||
)
|
||||
raise HTTPException(status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail=str(exc)) from exc
|
||||
|
||||
|
|
@ -266,20 +288,25 @@ async def webhook_run_flow(
|
|||
)
|
||||
logger.debug("Starting background task")
|
||||
background_tasks.add_task( # type: ignore
|
||||
simple_run_flow,
|
||||
simple_run_flow_task,
|
||||
flow=flow,
|
||||
input_request=input_request,
|
||||
)
|
||||
background_tasks.add_task(
|
||||
telemetry_service.log_package_run,
|
||||
RunPayload(IsWebhook=True, seconds=int(time.perf_counter() - start_time), success=True, errorMessage=""),
|
||||
RunPayload(
|
||||
runIsWebhook=True, runSeconds=int(time.perf_counter() - start_time), runSuccess=True, runErrorMessage=""
|
||||
),
|
||||
)
|
||||
return {"message": "Task started in the background", "status": "in progress"}
|
||||
except Exception as exc:
|
||||
background_tasks.add_task(
|
||||
telemetry_service.log_package_run,
|
||||
RunPayload(
|
||||
IsWebhook=True, seconds=int(time.perf_counter() - start_time), success=False, errorMessage=str(exc)
|
||||
runIsWebhook=True,
|
||||
runSeconds=int(time.perf_counter() - start_time),
|
||||
runSuccess=False,
|
||||
runErrorMessage=str(exc),
|
||||
),
|
||||
)
|
||||
if "Flow ID is required" in str(exc) or "Request body is empty" in str(exc):
|
||||
|
|
|
|||
|
|
@ -209,6 +209,23 @@ def update_flow(
|
|||
webhook_component = get_webhook_component_in_flow(db_flow.data)
|
||||
db_flow.webhook = webhook_component is not None
|
||||
db_flow.updated_at = datetime.now(timezone.utc)
|
||||
|
||||
# First check if the flow.name is unique
|
||||
# there might be flows with name like: "MyFlow", "MyFlow (1)", "MyFlow (2)"
|
||||
# so we need to check if the name is unique with `like` operator
|
||||
# if we find a flow with the same name, we add a number to the end of the name
|
||||
# based on the highest number found
|
||||
flow_from_db = session.exec(select(Flow).where(Flow.id == flow_id, Flow.user_id == current_user.id)).first()
|
||||
if flow_from_db:
|
||||
flows = session.exec(
|
||||
select(Flow).where(Flow.name.like(f"{flow.name} (%")).where(Flow.user_id == current_user.id) # type: ignore
|
||||
).all()
|
||||
if flows:
|
||||
numbers = [int(flow.name.split("(")[1].split(")")[0]) for flow in flows]
|
||||
flow.name = f"{flow.name} ({max(numbers) + 1})"
|
||||
else:
|
||||
flow.name = f"{flow.name} (1)"
|
||||
|
||||
if db_flow.folder_id is None:
|
||||
default_folder = session.exec(select(Folder).where(Folder.name == DEFAULT_FOLDER_NAME)).first()
|
||||
if default_folder:
|
||||
|
|
|
|||
|
|
@ -3,44 +3,40 @@ from pydantic.v1 import SecretStr
|
|||
from langflow.base.constants import STREAM_INFO_TEXT
|
||||
from langflow.base.models.model import LCModelComponent
|
||||
from langflow.field_typing import LanguageModel
|
||||
from langflow.io import BoolInput, DropdownInput, FloatInput, IntInput, MessageInput, Output, SecretStrInput, StrInput
|
||||
from langflow.inputs import (
|
||||
BoolInput,
|
||||
DropdownInput,
|
||||
FloatInput,
|
||||
IntInput,
|
||||
MessageInput,
|
||||
SecretStrInput,
|
||||
StrInput,
|
||||
)
|
||||
|
||||
|
||||
class GoogleGenerativeAIComponent(LCModelComponent):
|
||||
display_name: str = "Google Generative AI"
|
||||
description: str = "Generate text using Google Generative AI."
|
||||
display_name = "Google Generative AI"
|
||||
description = "Generate text using Google Generative AI."
|
||||
icon = "GoogleGenerativeAI"
|
||||
|
||||
inputs = [
|
||||
SecretStrInput(
|
||||
name="google_api_key",
|
||||
display_name="Google API Key",
|
||||
info="The Google API Key to use for the Google Generative AI.",
|
||||
MessageInput(name="input_value", display_name="Input"),
|
||||
IntInput(
|
||||
name="max_output_tokens",
|
||||
display_name="Max Output Tokens",
|
||||
info="The maximum number of tokens to generate.",
|
||||
),
|
||||
DropdownInput(
|
||||
name="model",
|
||||
display_name="Model",
|
||||
info="The name of the model to use.",
|
||||
options=["gemini-1.5-pro", "gemini-1.5-flash"],
|
||||
options=["gemini-1.5-pro", "gemini-1.5-flash", "gemini-1.0-pro", "gemini-1.0-pro-vision"],
|
||||
value="gemini-1.5-pro",
|
||||
),
|
||||
IntInput(
|
||||
name="max_output_tokens",
|
||||
display_name="Max Output Tokens",
|
||||
info="The maximum number of tokens to generate.",
|
||||
advanced=True,
|
||||
),
|
||||
FloatInput(
|
||||
name="temperature",
|
||||
display_name="Temperature",
|
||||
info="Run inference with this temperature. Must by in the closed interval [0.0, 1.0].",
|
||||
value=0.1,
|
||||
),
|
||||
IntInput(
|
||||
name="top_k",
|
||||
display_name="Top K",
|
||||
info="Decode using top-k sampling: consider the set of top_k most probable tokens. Must be positive.",
|
||||
advanced=True,
|
||||
SecretStrInput(
|
||||
name="google_api_key",
|
||||
display_name="Google API Key",
|
||||
info="The Google API Key to use for the Google Generative AI.",
|
||||
),
|
||||
FloatInput(
|
||||
name="top_p",
|
||||
|
|
@ -48,29 +44,26 @@ class GoogleGenerativeAIComponent(LCModelComponent):
|
|||
info="The maximum cumulative probability of tokens to consider when sampling.",
|
||||
advanced=True,
|
||||
),
|
||||
FloatInput(name="temperature", display_name="Temperature", value=0.1),
|
||||
BoolInput(name="stream", display_name="Stream", info=STREAM_INFO_TEXT, advanced=True),
|
||||
IntInput(
|
||||
name="n",
|
||||
display_name="N",
|
||||
info="Number of chat completions to generate for each prompt. Note that the API may not return the full n completions if duplicates are generated.",
|
||||
advanced=True,
|
||||
),
|
||||
MessageInput(
|
||||
name="input_value",
|
||||
display_name="Input",
|
||||
info="The input to the model.",
|
||||
input_types=["Text", "Data", "Prompt"],
|
||||
),
|
||||
BoolInput(name="stream", display_name="Stream", info=STREAM_INFO_TEXT, advanced=True),
|
||||
StrInput(
|
||||
name="system_message",
|
||||
display_name="System Message",
|
||||
info="System message to pass to the model.",
|
||||
advanced=True,
|
||||
),
|
||||
]
|
||||
outputs = [
|
||||
Output(display_name="Text", name="text_output", method="text_response"),
|
||||
Output(display_name="Language Model", name="model_output", method="build_model"),
|
||||
IntInput(
|
||||
name="top_k",
|
||||
display_name="Top K",
|
||||
info="Decode using top-k sampling: consider the set of top_k most probable tokens. Must be positive.",
|
||||
advanced=True,
|
||||
),
|
||||
]
|
||||
|
||||
def build_model(self) -> LanguageModel:
|
||||
|
|
|
|||
|
|
@ -164,15 +164,15 @@ async def get_current_user_for_websocket(
|
|||
|
||||
def get_current_active_user(current_user: Annotated[User, Depends(get_current_user)]):
|
||||
if not current_user.is_active:
|
||||
raise HTTPException(status_code=400, detail="Inactive user")
|
||||
raise HTTPException(status_code=status.HTTP_401_UNAUTHORIZED, detail="Inactive user")
|
||||
return current_user
|
||||
|
||||
|
||||
def get_current_active_superuser(current_user: Annotated[User, Depends(get_current_user)]) -> User:
|
||||
if not current_user.is_active:
|
||||
raise HTTPException(status_code=401, detail="Inactive user")
|
||||
raise HTTPException(status_code=status.HTTP_401_UNAUTHORIZED, detail="Inactive user")
|
||||
if not current_user.is_superuser:
|
||||
raise HTTPException(status_code=400, detail="The user doesn't have enough privileges")
|
||||
raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="The user doesn't have enough privileges")
|
||||
return current_user
|
||||
|
||||
|
||||
|
|
@ -324,8 +324,8 @@ def authenticate_user(username: str, password: str, db: Session = Depends(get_se
|
|||
|
||||
if not user.is_active:
|
||||
if not user.last_login_at:
|
||||
raise HTTPException(status_code=400, detail="Waiting for approval")
|
||||
raise HTTPException(status_code=400, detail="Inactive user")
|
||||
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="Waiting for approval")
|
||||
raise HTTPException(status_code=status.HTTP_401_UNAUTHORIZED, detail="Inactive user")
|
||||
|
||||
return user if verify_password(password, user.password) else None
|
||||
|
||||
|
|
|
|||
|
|
@ -2,10 +2,10 @@ from pydantic import BaseModel
|
|||
|
||||
|
||||
class RunPayload(BaseModel):
|
||||
IsWebhook: bool = False
|
||||
seconds: int
|
||||
success: bool
|
||||
errorMessage: str = ""
|
||||
runIsWebhook: bool = False
|
||||
runSeconds: int
|
||||
runSuccess: bool
|
||||
runErrorMessage: str = ""
|
||||
|
||||
|
||||
class ShutdownPayload(BaseModel):
|
||||
|
|
@ -23,14 +23,14 @@ class VersionPayload(BaseModel):
|
|||
|
||||
|
||||
class PlaygroundPayload(BaseModel):
|
||||
seconds: int
|
||||
componentCount: int | None = None
|
||||
success: bool
|
||||
errorMessage: str = ""
|
||||
playgroundSeconds: int
|
||||
playgroundComponentCount: int | None = None
|
||||
playgroundSuccess: bool
|
||||
playgroundErrorMessage: str = ""
|
||||
|
||||
|
||||
class ComponentPayload(BaseModel):
|
||||
name: str
|
||||
seconds: int
|
||||
success: bool
|
||||
errorMessage: str
|
||||
componentName: str
|
||||
componentSeconds: int
|
||||
componentSuccess: bool
|
||||
componentErrorMessage: str
|
||||
|
|
|
|||
11
src/backend/base/poetry.lock
generated
11
src/backend/base/poetry.lock
generated
|
|
@ -1295,18 +1295,21 @@ types-requests = ">=2.31.0.2,<3.0.0.0"
|
|||
|
||||
[[package]]
|
||||
name = "langsmith"
|
||||
version = "0.1.81"
|
||||
version = "0.1.82"
|
||||
description = "Client library to connect to the LangSmith LLM Tracing and Evaluation Platform."
|
||||
optional = false
|
||||
python-versions = "<4.0,>=3.8.1"
|
||||
files = [
|
||||
{file = "langsmith-0.1.81-py3-none-any.whl", hash = "sha256:3251d823225eef23ee541980b9d9e506367eabbb7f985a086b5d09e8f78ba7e9"},
|
||||
{file = "langsmith-0.1.81.tar.gz", hash = "sha256:585ef3a2251380bd2843a664c9a28da4a7d28432e3ee8bcebf291ffb8e1f0af0"},
|
||||
{file = "langsmith-0.1.82-py3-none-any.whl", hash = "sha256:9b3653e7d316036b0c60bf0bc3e280662d660f485a4ebd8e5c9d84f9831ae79c"},
|
||||
{file = "langsmith-0.1.82.tar.gz", hash = "sha256:c02e2bbc488c10c13b52c69d271eb40bd38da078d37b6ae7ae04a18bd48140be"},
|
||||
]
|
||||
|
||||
[package.dependencies]
|
||||
orjson = ">=3.9.14,<4.0.0"
|
||||
pydantic = ">=1,<3"
|
||||
pydantic = [
|
||||
{version = ">=1,<3", markers = "python_full_version < \"3.12.4\""},
|
||||
{version = ">=2.7.4,<3.0.0", markers = "python_full_version >= \"3.12.4\""},
|
||||
]
|
||||
requests = ">=2,<3"
|
||||
|
||||
[[package]]
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
[tool.poetry]
|
||||
name = "langflow-base"
|
||||
version = "0.0.79"
|
||||
version = "0.0.81"
|
||||
description = "A Python package with a built-in web application"
|
||||
authors = ["Langflow <contact@langflow.org>"]
|
||||
maintainers = [
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue