feat: migrate vertex_builds to sql database (#2978)
* feat: migrate vertex_builds to sql database * [autofix.ci] apply automated fixes * name --------- Co-authored-by: autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com>
This commit is contained in:
parent
d71356bc16
commit
c7575b18df
17 changed files with 219 additions and 193 deletions
|
|
@ -0,0 +1,51 @@
|
||||||
|
"""create vertex_builds table
|
||||||
|
|
||||||
|
Revision ID: 0d60fcbd4e8e
|
||||||
|
Revises: 90be8e2ed91e
|
||||||
|
Create Date: 2024-07-26 11:41:31.274271
|
||||||
|
|
||||||
|
"""
|
||||||
|
|
||||||
|
from typing import Sequence, Union
|
||||||
|
|
||||||
|
from alembic import op
|
||||||
|
import sqlalchemy as sa
|
||||||
|
import sqlmodel
|
||||||
|
from langflow.utils import migration
|
||||||
|
|
||||||
|
|
||||||
|
# revision identifiers, used by Alembic.
|
||||||
|
revision: str = "0d60fcbd4e8e"
|
||||||
|
down_revision: Union[str, None] = "90be8e2ed91e"
|
||||||
|
branch_labels: Union[str, Sequence[str], None] = None
|
||||||
|
depends_on: Union[str, Sequence[str], None] = None
|
||||||
|
|
||||||
|
|
||||||
|
def upgrade() -> None:
|
||||||
|
conn = op.get_bind()
|
||||||
|
if not migration.table_exists("vertex_build", conn):
|
||||||
|
op.create_table(
|
||||||
|
"vertex_build",
|
||||||
|
sa.Column("timestamp", sa.DateTime(), nullable=False),
|
||||||
|
sa.Column("id", sqlmodel.sql.sqltypes.AutoString(), nullable=True),
|
||||||
|
sa.Column("data", sa.JSON(), nullable=True),
|
||||||
|
sa.Column("artifacts", sa.JSON(), nullable=True),
|
||||||
|
sa.Column("params", sqlmodel.sql.sqltypes.AutoString(), nullable=True),
|
||||||
|
sa.Column("build_id", sqlmodel.sql.sqltypes.GUID(), nullable=False),
|
||||||
|
sa.Column("flow_id", sqlmodel.sql.sqltypes.GUID(), nullable=False),
|
||||||
|
sa.Column("valid", sa.BOOLEAN(), nullable=False),
|
||||||
|
sa.ForeignKeyConstraint(
|
||||||
|
["flow_id"],
|
||||||
|
["flow.id"],
|
||||||
|
"fk_vertex_build_flow_id",
|
||||||
|
),
|
||||||
|
sa.PrimaryKeyConstraint("build_id"),
|
||||||
|
)
|
||||||
|
pass
|
||||||
|
|
||||||
|
|
||||||
|
def downgrade() -> None:
|
||||||
|
conn = op.get_bind()
|
||||||
|
if migration.table_exists("vertex_build", conn):
|
||||||
|
op.drop_table("vertex_build")
|
||||||
|
pass
|
||||||
|
|
@ -25,11 +25,11 @@ from langflow.api.v1.schemas import (
|
||||||
)
|
)
|
||||||
from langflow.exceptions.component import ComponentBuildException
|
from langflow.exceptions.component import ComponentBuildException
|
||||||
from langflow.graph.graph.base import Graph
|
from langflow.graph.graph.base import Graph
|
||||||
|
from langflow.graph.utils import log_vertex_build
|
||||||
from langflow.schema.schema import OutputValue
|
from langflow.schema.schema import OutputValue
|
||||||
from langflow.services.auth.utils import get_current_active_user
|
from langflow.services.auth.utils import get_current_active_user
|
||||||
from langflow.services.chat.service import ChatService
|
from langflow.services.chat.service import ChatService
|
||||||
from langflow.services.deps import get_chat_service, get_session, get_session_service, get_telemetry_service
|
from langflow.services.deps import get_chat_service, get_session, get_session_service, get_telemetry_service
|
||||||
from langflow.services.monitor.utils import log_vertex_build
|
|
||||||
from langflow.services.telemetry.schema import ComponentPayload, PlaygroundPayload
|
from langflow.services.telemetry.schema import ComponentPayload, PlaygroundPayload
|
||||||
from langflow.services.telemetry.service import TelemetryService
|
from langflow.services.telemetry.service import TelemetryService
|
||||||
|
|
||||||
|
|
@ -233,7 +233,7 @@ async def build_vertex(
|
||||||
background_tasks.add_task(
|
background_tasks.add_task(
|
||||||
log_vertex_build,
|
log_vertex_build,
|
||||||
flow_id=flow_id_str,
|
flow_id=flow_id_str,
|
||||||
vertex_id=vertex_id.split("-")[0],
|
vertex_id=vertex_id,
|
||||||
valid=valid,
|
valid=valid,
|
||||||
params=params,
|
params=params,
|
||||||
data=result_data_response,
|
data=result_data_response,
|
||||||
|
|
|
||||||
|
|
@ -10,39 +10,36 @@ from langflow.services.database.models.message.model import MessageRead, Message
|
||||||
from langflow.services.database.models.transactions.crud import get_transactions_by_flow_id
|
from langflow.services.database.models.transactions.crud import get_transactions_by_flow_id
|
||||||
from langflow.services.database.models.transactions.model import TransactionReadResponse
|
from langflow.services.database.models.transactions.model import TransactionReadResponse
|
||||||
from langflow.services.database.models.user.model import User
|
from langflow.services.database.models.user.model import User
|
||||||
from langflow.services.deps import get_monitor_service, get_session
|
from langflow.services.database.models.vertex_builds.crud import (
|
||||||
from langflow.services.monitor.schema import MessageModelResponse, VertexBuildMapModel
|
get_vertex_builds_by_flow_id,
|
||||||
from langflow.services.monitor.service import MonitorService
|
delete_vertex_builds_by_flow_id,
|
||||||
|
)
|
||||||
|
from langflow.services.database.models.vertex_builds.model import VertexBuildMapModel
|
||||||
|
from langflow.services.deps import get_session
|
||||||
|
from langflow.services.monitor.schema import MessageModelResponse
|
||||||
|
|
||||||
router = APIRouter(prefix="/monitor", tags=["Monitor"])
|
router = APIRouter(prefix="/monitor", tags=["Monitor"])
|
||||||
|
|
||||||
|
|
||||||
# Get vertex_builds data from the monitor service
|
|
||||||
@router.get("/builds", response_model=VertexBuildMapModel)
|
@router.get("/builds", response_model=VertexBuildMapModel)
|
||||||
async def get_vertex_builds(
|
async def get_vertex_builds(
|
||||||
flow_id: Optional[str] = Query(None),
|
flow_id: UUID = Query(),
|
||||||
vertex_id: Optional[str] = Query(None),
|
session: Session = Depends(get_session),
|
||||||
valid: Optional[bool] = Query(None),
|
|
||||||
order_by: Optional[str] = Query("timestamp"),
|
|
||||||
monitor_service: MonitorService = Depends(get_monitor_service),
|
|
||||||
):
|
):
|
||||||
try:
|
try:
|
||||||
vertex_build_dicts = monitor_service.get_vertex_builds(
|
vertex_builds = get_vertex_builds_by_flow_id(session, flow_id)
|
||||||
flow_id=flow_id, vertex_id=vertex_id, valid=valid, order_by=order_by
|
return VertexBuildMapModel.from_list_of_dicts(vertex_builds)
|
||||||
)
|
|
||||||
vertex_build_map = VertexBuildMapModel.from_list_of_dicts(vertex_build_dicts)
|
|
||||||
return vertex_build_map
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
raise HTTPException(status_code=500, detail=str(e))
|
raise HTTPException(status_code=500, detail=str(e))
|
||||||
|
|
||||||
|
|
||||||
@router.delete("/builds", status_code=204)
|
@router.delete("/builds", status_code=204)
|
||||||
async def delete_vertex_builds(
|
async def delete_vertex_builds(
|
||||||
flow_id: Optional[str] = Query(None),
|
flow_id: UUID = Query(),
|
||||||
monitor_service: MonitorService = Depends(get_monitor_service),
|
session: Session = Depends(get_session),
|
||||||
):
|
):
|
||||||
try:
|
try:
|
||||||
monitor_service.delete_vertex_builds(flow_id=flow_id)
|
delete_vertex_builds_by_flow_id(session, flow_id)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
raise HTTPException(status_code=500, detail=str(e))
|
raise HTTPException(status_code=500, detail=str(e))
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -4,7 +4,6 @@ from loguru import logger
|
||||||
from pydantic import BaseModel, Field, field_validator
|
from pydantic import BaseModel, Field, field_validator
|
||||||
|
|
||||||
from langflow.schema.schema import INPUT_FIELD_NAME
|
from langflow.schema.schema import INPUT_FIELD_NAME
|
||||||
from langflow.services.monitor.utils import log_message
|
|
||||||
|
|
||||||
if TYPE_CHECKING:
|
if TYPE_CHECKING:
|
||||||
from langflow.graph.vertex.base import Vertex
|
from langflow.graph.vertex.base import Vertex
|
||||||
|
|
@ -224,13 +223,6 @@ class ContractEdge(Edge):
|
||||||
):
|
):
|
||||||
if target.params.get("message") == "":
|
if target.params.get("message") == "":
|
||||||
return self.result
|
return self.result
|
||||||
await log_message(
|
|
||||||
sender=target.params.get("sender", ""),
|
|
||||||
sender_name=target.params.get("sender_name", ""),
|
|
||||||
message=target.params.get(INPUT_FIELD_NAME, {}),
|
|
||||||
session_id=target.params.get("session_id", ""),
|
|
||||||
flow_id=target.graph.flow_id,
|
|
||||||
)
|
|
||||||
return self.result
|
return self.result
|
||||||
|
|
||||||
def __repr__(self) -> str:
|
def __repr__(self) -> str:
|
||||||
|
|
|
||||||
|
|
@ -11,12 +11,15 @@ from langflow.schema.data import Data
|
||||||
from langflow.schema.message import Message
|
from langflow.schema.message import Message
|
||||||
from langflow.services.database.models.transactions.model import TransactionBase
|
from langflow.services.database.models.transactions.model import TransactionBase
|
||||||
from langflow.services.database.models.transactions.crud import log_transaction as crud_log_transaction
|
from langflow.services.database.models.transactions.crud import log_transaction as crud_log_transaction
|
||||||
|
from langflow.services.database.models.vertex_builds.crud import log_vertex_build as crud_log_vertex_build
|
||||||
|
from langflow.services.database.models.vertex_builds.model import VertexBuildBase
|
||||||
from langflow.services.database.utils import session_getter
|
from langflow.services.database.utils import session_getter
|
||||||
from langflow.services.deps import get_db_service
|
from langflow.services.deps import get_db_service
|
||||||
from loguru import logger
|
from loguru import logger
|
||||||
|
|
||||||
if TYPE_CHECKING:
|
if TYPE_CHECKING:
|
||||||
from langflow.graph.vertex.base import Vertex
|
from langflow.graph.vertex.base import Vertex
|
||||||
|
from langflow.api.v1.schemas import ResultDataResponse
|
||||||
|
|
||||||
|
|
||||||
class UnbuiltObject:
|
class UnbuiltObject:
|
||||||
|
|
@ -145,3 +148,28 @@ async def log_transaction(
|
||||||
logger.debug(f"Logged transaction: {inserted.id}")
|
logger.debug(f"Logged transaction: {inserted.id}")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"Error logging transaction: {e}")
|
logger.error(f"Error logging transaction: {e}")
|
||||||
|
|
||||||
|
|
||||||
|
def log_vertex_build(
|
||||||
|
flow_id: str,
|
||||||
|
vertex_id: str,
|
||||||
|
valid: bool,
|
||||||
|
params: Any,
|
||||||
|
data: "ResultDataResponse",
|
||||||
|
artifacts: Optional[dict] = None,
|
||||||
|
):
|
||||||
|
try:
|
||||||
|
vertex_build = VertexBuildBase(
|
||||||
|
flow_id=flow_id,
|
||||||
|
id=vertex_id,
|
||||||
|
valid=valid,
|
||||||
|
params=str(params) if params else None,
|
||||||
|
# ugly hack to get the model dump with weird datatypes
|
||||||
|
data=json.loads(data.model_dump_json()),
|
||||||
|
artifacts=artifacts,
|
||||||
|
)
|
||||||
|
with session_getter(get_db_service()) as session:
|
||||||
|
inserted = crud_log_vertex_build(session, vertex_build)
|
||||||
|
logger.debug(f"Logged vertex build: {inserted.build_id}")
|
||||||
|
except Exception as e:
|
||||||
|
logger.exception(f"Error logging vertex build: {e}")
|
||||||
|
|
|
||||||
|
|
@ -7,13 +7,12 @@ from langchain_core.messages import AIMessage, AIMessageChunk
|
||||||
from loguru import logger
|
from loguru import logger
|
||||||
|
|
||||||
from langflow.graph.schema import CHAT_COMPONENTS, RECORDS_COMPONENTS, InterfaceComponentTypes, ResultData
|
from langflow.graph.schema import CHAT_COMPONENTS, RECORDS_COMPONENTS, InterfaceComponentTypes, ResultData
|
||||||
from langflow.graph.utils import UnbuiltObject, serialize_field, log_transaction
|
from langflow.graph.utils import UnbuiltObject, serialize_field, log_transaction, log_vertex_build
|
||||||
from langflow.graph.vertex.base import Vertex
|
from langflow.graph.vertex.base import Vertex
|
||||||
from langflow.schema import Data
|
from langflow.schema import Data
|
||||||
from langflow.schema.artifact import ArtifactType
|
from langflow.schema.artifact import ArtifactType
|
||||||
from langflow.schema.message import Message
|
from langflow.schema.message import Message
|
||||||
from langflow.schema.schema import INPUT_FIELD_NAME
|
from langflow.schema.schema import INPUT_FIELD_NAME
|
||||||
from langflow.services.monitor.utils import log_vertex_build
|
|
||||||
from langflow.template.field.base import UNDEFINED
|
from langflow.template.field.base import UNDEFINED
|
||||||
from langflow.utils.schemas import ChatOutputResponse, DataOutputResponse
|
from langflow.utils.schemas import ChatOutputResponse, DataOutputResponse
|
||||||
from langflow.utils.util import unescape_string
|
from langflow.utils.util import unescape_string
|
||||||
|
|
@ -389,13 +388,15 @@ class InterfaceVertex(ComponentVertex):
|
||||||
if isinstance(value, (AsyncIterator, Iterator)):
|
if isinstance(value, (AsyncIterator, Iterator)):
|
||||||
origin_vertex.results[key] = complete_message
|
origin_vertex.results[key] = complete_message
|
||||||
|
|
||||||
await log_vertex_build(
|
asyncio.create_task(
|
||||||
flow_id=self.graph.flow_id,
|
log_vertex_build(
|
||||||
vertex_id=self.id,
|
flow_id=self.graph.flow_id,
|
||||||
valid=True,
|
vertex_id=self.id,
|
||||||
params=self._built_object_repr(),
|
valid=True,
|
||||||
data=self.result,
|
params=self._built_object_repr(),
|
||||||
artifacts=self.artifacts,
|
data=self.result,
|
||||||
|
artifacts=self.artifacts,
|
||||||
|
)
|
||||||
)
|
)
|
||||||
|
|
||||||
self._validate_built_object()
|
self._validate_built_object()
|
||||||
|
|
|
||||||
|
|
@ -14,6 +14,7 @@ from sqlalchemy import UniqueConstraint
|
||||||
from sqlmodel import JSON, Column, Field, Relationship, SQLModel
|
from sqlmodel import JSON, Column, Field, Relationship, SQLModel
|
||||||
|
|
||||||
from langflow.schema import Data
|
from langflow.schema import Data
|
||||||
|
from langflow.services.database.models.vertex_builds.model import VertexBuildTable
|
||||||
|
|
||||||
if TYPE_CHECKING:
|
if TYPE_CHECKING:
|
||||||
from langflow.services.database.models.folder import Folder
|
from langflow.services.database.models.folder import Folder
|
||||||
|
|
@ -145,6 +146,7 @@ class Flow(FlowBase, table=True):
|
||||||
folder: Optional["Folder"] = Relationship(back_populates="flows")
|
folder: Optional["Folder"] = Relationship(back_populates="flows")
|
||||||
messages: List["MessageTable"] = Relationship(back_populates="flow")
|
messages: List["MessageTable"] = Relationship(back_populates="flow")
|
||||||
transactions: List["TransactionTable"] = Relationship(back_populates="flow")
|
transactions: List["TransactionTable"] = Relationship(back_populates="flow")
|
||||||
|
vertex_builds: List["VertexBuildTable"] = Relationship(back_populates="flow")
|
||||||
|
|
||||||
def to_data(self):
|
def to_data(self):
|
||||||
serialized = self.model_dump()
|
serialized = self.model_dump()
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,3 @@
|
||||||
|
from .model import VertexBuildTable
|
||||||
|
|
||||||
|
__all__ = ["VertexBuildTable"]
|
||||||
|
|
@ -0,0 +1,36 @@
|
||||||
|
from typing import Optional
|
||||||
|
from uuid import UUID
|
||||||
|
|
||||||
|
from sqlalchemy.exc import IntegrityError
|
||||||
|
from sqlmodel import Session, select, col
|
||||||
|
from sqlalchemy import delete
|
||||||
|
|
||||||
|
from langflow.services.database.models.vertex_builds.model import VertexBuildBase, VertexBuildTable
|
||||||
|
|
||||||
|
|
||||||
|
def get_vertex_builds_by_flow_id(db: Session, flow_id: UUID, limit: Optional[int] = 1000) -> list[VertexBuildTable]:
|
||||||
|
stmt = (
|
||||||
|
select(VertexBuildTable)
|
||||||
|
.where(VertexBuildTable.flow_id == flow_id)
|
||||||
|
.order_by(col(VertexBuildTable.timestamp))
|
||||||
|
.limit(limit)
|
||||||
|
)
|
||||||
|
|
||||||
|
builds = db.exec(stmt)
|
||||||
|
return [t for t in builds]
|
||||||
|
|
||||||
|
|
||||||
|
def log_vertex_build(db: Session, vertex_build: VertexBuildBase) -> VertexBuildTable:
|
||||||
|
table = VertexBuildTable(**vertex_build.model_dump())
|
||||||
|
db.add(table)
|
||||||
|
try:
|
||||||
|
db.commit()
|
||||||
|
return table
|
||||||
|
except IntegrityError as e:
|
||||||
|
db.rollback()
|
||||||
|
raise e
|
||||||
|
|
||||||
|
|
||||||
|
def delete_vertex_builds_by_flow_id(db: Session, flow_id: UUID) -> None:
|
||||||
|
delete(VertexBuildTable).where(col(VertexBuildTable.flow_id == flow_id))
|
||||||
|
db.commit()
|
||||||
|
|
@ -0,0 +1,52 @@
|
||||||
|
from datetime import datetime, timezone
|
||||||
|
from typing import TYPE_CHECKING, Optional
|
||||||
|
from uuid import UUID, uuid4
|
||||||
|
|
||||||
|
from pydantic import field_validator, BaseModel
|
||||||
|
from sqlmodel import JSON, Column, Field, Relationship, SQLModel
|
||||||
|
|
||||||
|
|
||||||
|
if TYPE_CHECKING:
|
||||||
|
from langflow.services.database.models.flow.model import Flow
|
||||||
|
|
||||||
|
|
||||||
|
class VertexBuildBase(SQLModel):
|
||||||
|
timestamp: datetime = Field(default_factory=lambda: datetime.now(timezone.utc))
|
||||||
|
id: str = Field(nullable=False)
|
||||||
|
data: Optional[dict] = Field(default=None, sa_column=Column(JSON))
|
||||||
|
artifacts: Optional[dict] = Field(default=None, sa_column=Column(JSON))
|
||||||
|
params: Optional[str] = Field(nullable=True)
|
||||||
|
valid: bool = Field(nullable=False)
|
||||||
|
flow_id: UUID = Field(foreign_key="flow.id")
|
||||||
|
|
||||||
|
# Needed for Column(JSON)
|
||||||
|
class Config:
|
||||||
|
arbitrary_types_allowed = True
|
||||||
|
|
||||||
|
@field_validator("flow_id", mode="before")
|
||||||
|
@classmethod
|
||||||
|
def validate_flow_id(cls, value):
|
||||||
|
if value is None:
|
||||||
|
return value
|
||||||
|
if isinstance(value, str):
|
||||||
|
value = UUID(value)
|
||||||
|
return value
|
||||||
|
|
||||||
|
|
||||||
|
class VertexBuildTable(VertexBuildBase, table=True):
|
||||||
|
__tablename__ = "vertex_build"
|
||||||
|
build_id: Optional[UUID] = Field(default_factory=uuid4, primary_key=True)
|
||||||
|
flow: "Flow" = Relationship(back_populates="vertex_builds")
|
||||||
|
|
||||||
|
|
||||||
|
class VertexBuildMapModel(BaseModel):
|
||||||
|
vertex_builds: dict[str, list[VertexBuildTable]]
|
||||||
|
|
||||||
|
@classmethod
|
||||||
|
def from_list_of_dicts(cls, vertex_build_dicts: list[VertexBuildTable]):
|
||||||
|
vertex_build_map: dict[str, list[VertexBuildTable]] = {}
|
||||||
|
for vertex_build in vertex_build_dicts:
|
||||||
|
if vertex_build.id not in vertex_build_map:
|
||||||
|
vertex_build_map[vertex_build.id] = []
|
||||||
|
vertex_build_map[vertex_build.id].append(vertex_build)
|
||||||
|
return cls(vertex_builds=vertex_build_map)
|
||||||
|
|
@ -274,7 +274,7 @@ class DatabaseService(Service):
|
||||||
|
|
||||||
inspector = inspect(self.engine)
|
inspector = inspect(self.engine)
|
||||||
table_names = inspector.get_table_names()
|
table_names = inspector.get_table_names()
|
||||||
current_tables = ["flow", "user", "apikey", "folder", "message", "variable", "transaction"]
|
current_tables = ["flow", "user", "apikey", "folder", "message", "variable", "transaction", "vertex_build"]
|
||||||
|
|
||||||
if table_names and all(table in table_names for table in current_tables):
|
if table_names and all(table in table_names for table in current_tables):
|
||||||
logger.debug("Database and tables already exist")
|
logger.debug("Database and tables already exist")
|
||||||
|
|
|
||||||
|
|
@ -244,12 +244,6 @@ class VertexBuildResponseModel(VertexBuildModel):
|
||||||
return v
|
return v
|
||||||
|
|
||||||
|
|
||||||
def to_map(value: dict):
|
|
||||||
keys = list(value.keys())
|
|
||||||
values = list(value.values())
|
|
||||||
return {"key": keys, "value": values}
|
|
||||||
|
|
||||||
|
|
||||||
class VertexBuildMapModel(BaseModel):
|
class VertexBuildMapModel(BaseModel):
|
||||||
vertex_builds: dict[str, list[VertexBuildResponseModel]]
|
vertex_builds: dict[str, list[VertexBuildResponseModel]]
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -1,71 +1,33 @@
|
||||||
from datetime import datetime
|
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from typing import TYPE_CHECKING, Union, List
|
from typing import TYPE_CHECKING, List
|
||||||
|
|
||||||
from loguru import logger
|
|
||||||
from platformdirs import user_cache_dir
|
from platformdirs import user_cache_dir
|
||||||
|
|
||||||
from langflow.services.base import Service
|
from langflow.services.base import Service
|
||||||
from langflow.services.monitor.utils import (
|
from langflow.services.monitor.utils import (
|
||||||
add_row_to_table,
|
|
||||||
drop_and_create_table_if_schema_mismatch,
|
|
||||||
new_duckdb_locked_connection,
|
new_duckdb_locked_connection,
|
||||||
)
|
)
|
||||||
|
|
||||||
if TYPE_CHECKING:
|
if TYPE_CHECKING:
|
||||||
from langflow.services.monitor.schema import VertexBuildModel
|
|
||||||
from langflow.services.settings.service import SettingsService
|
from langflow.services.settings.service import SettingsService
|
||||||
|
|
||||||
|
|
||||||
class MonitorService(Service):
|
class MonitorService(Service):
|
||||||
|
"""
|
||||||
|
Deprecated. Still connecting to duckdb to migrate old installations.
|
||||||
|
"""
|
||||||
|
|
||||||
name = "monitor_service"
|
name = "monitor_service"
|
||||||
|
|
||||||
def __init__(self, settings_service: "SettingsService"):
|
def __init__(self, settings_service: "SettingsService"):
|
||||||
from langflow.services.monitor.schema import VertexBuildModel
|
|
||||||
|
|
||||||
self.settings_service = settings_service
|
self.settings_service = settings_service
|
||||||
self.base_cache_dir = Path(user_cache_dir("langflow"), ensure_exists=True)
|
self.base_cache_dir = Path(user_cache_dir("langflow"), ensure_exists=True)
|
||||||
self.db_path = self.base_cache_dir / "monitor.duckdb"
|
self.db_path = self.base_cache_dir / "monitor.duckdb"
|
||||||
self.table_map: dict[str, type[VertexBuildModel]] = {
|
|
||||||
"vertex_builds": VertexBuildModel,
|
|
||||||
}
|
|
||||||
|
|
||||||
try:
|
|
||||||
self.ensure_tables_exist()
|
|
||||||
except Exception as e:
|
|
||||||
logger.exception(f"Error initializing monitor service: {e}")
|
|
||||||
|
|
||||||
def exec_query(self, query: str, read_only: bool = False):
|
def exec_query(self, query: str, read_only: bool = False):
|
||||||
with new_duckdb_locked_connection(self.db_path, read_only=read_only) as conn:
|
with new_duckdb_locked_connection(self.db_path, read_only=read_only) as conn:
|
||||||
return conn.execute(query).df()
|
return conn.execute(query).df()
|
||||||
|
|
||||||
def to_df(self, table_name):
|
|
||||||
return self.load_table_as_dataframe(table_name)
|
|
||||||
|
|
||||||
def ensure_tables_exist(self):
|
|
||||||
for table_name, model in self.table_map.items():
|
|
||||||
drop_and_create_table_if_schema_mismatch(str(self.db_path), table_name, model)
|
|
||||||
|
|
||||||
def add_row(
|
|
||||||
self,
|
|
||||||
table_name: str,
|
|
||||||
data: Union[dict, "VertexBuildModel"],
|
|
||||||
):
|
|
||||||
model = self.table_map.get(table_name)
|
|
||||||
if model is None:
|
|
||||||
raise ValueError(f"Unknown table name: {table_name}")
|
|
||||||
|
|
||||||
with new_duckdb_locked_connection(self.db_path, read_only=False) as conn:
|
|
||||||
add_row_to_table(conn, table_name, model, data)
|
|
||||||
|
|
||||||
def load_table_as_dataframe(self, table_name):
|
|
||||||
with new_duckdb_locked_connection(self.db_path, read_only=True) as conn:
|
|
||||||
return conn.table(table_name).df()
|
|
||||||
|
|
||||||
@staticmethod
|
|
||||||
def get_timestamp():
|
|
||||||
return datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
|
||||||
|
|
||||||
def get_messages(
|
def get_messages(
|
||||||
self,
|
self,
|
||||||
flow_id: str | None = None,
|
flow_id: str | None = None,
|
||||||
|
|
@ -102,47 +64,6 @@ class MonitorService(Service):
|
||||||
|
|
||||||
return df
|
return df
|
||||||
|
|
||||||
def get_vertex_builds(
|
|
||||||
self,
|
|
||||||
flow_id: str | None = None,
|
|
||||||
vertex_id: str | None = None,
|
|
||||||
valid: bool | None = None,
|
|
||||||
order_by: str | None = "timestamp",
|
|
||||||
):
|
|
||||||
query = "SELECT id, index,flow_id, valid, params, data, artifacts, timestamp FROM vertex_builds"
|
|
||||||
conditions = []
|
|
||||||
if flow_id:
|
|
||||||
conditions.append(f"flow_id = '{flow_id}'")
|
|
||||||
if vertex_id:
|
|
||||||
conditions.append(f"id = '{vertex_id}'")
|
|
||||||
if valid is not None: # Check for None because valid is a boolean
|
|
||||||
valid_str = "true" if valid else "false"
|
|
||||||
conditions.append(f"valid = {valid_str}")
|
|
||||||
|
|
||||||
if conditions:
|
|
||||||
query += " WHERE " + " AND ".join(conditions)
|
|
||||||
|
|
||||||
if order_by:
|
|
||||||
query += f" ORDER BY {order_by}"
|
|
||||||
|
|
||||||
with new_duckdb_locked_connection(self.db_path, read_only=True) as conn:
|
|
||||||
df = conn.execute(query).df()
|
|
||||||
|
|
||||||
return df.to_dict(orient="records")
|
|
||||||
|
|
||||||
def delete_vertex_builds(self, flow_id: str | None = None):
|
|
||||||
query = "DELETE FROM vertex_builds"
|
|
||||||
if flow_id:
|
|
||||||
query += f" WHERE flow_id = '{flow_id}'"
|
|
||||||
|
|
||||||
with new_duckdb_locked_connection(self.db_path, read_only=False) as conn:
|
|
||||||
conn.execute(query)
|
|
||||||
|
|
||||||
def delete_messages_session(self, session_id: str):
|
|
||||||
query = f"DELETE FROM messages WHERE session_id = '{session_id}'"
|
|
||||||
|
|
||||||
return self.exec_query(query, read_only=False)
|
|
||||||
|
|
||||||
def delete_messages(self, message_ids: list[int] | str):
|
def delete_messages(self, message_ids: list[int] | str):
|
||||||
if isinstance(message_ids, list):
|
if isinstance(message_ids, list):
|
||||||
# If message_ids is a list, join the string representations of the integers
|
# If message_ids is a list, join the string representations of the integers
|
||||||
|
|
@ -157,13 +78,6 @@ class MonitorService(Service):
|
||||||
|
|
||||||
return self.exec_query(query, read_only=False)
|
return self.exec_query(query, read_only=False)
|
||||||
|
|
||||||
def update_message(self, message_id: str, **kwargs):
|
|
||||||
query = (
|
|
||||||
f"""UPDATE messages SET {', '.join(f"{k} = '{v}'" for k, v in kwargs.items())} WHERE index = {message_id}"""
|
|
||||||
)
|
|
||||||
|
|
||||||
return self.exec_query(query, read_only=False)
|
|
||||||
|
|
||||||
def get_transactions(self, limit: int = 100):
|
def get_transactions(self, limit: int = 100):
|
||||||
query = f"SELECT index,flow_id, status, error, timestamp, vertex_id, inputs, outputs, target_id FROM transactions LIMIT {str(limit)}"
|
query = f"SELECT index,flow_id, status, error, timestamp, vertex_id, inputs, outputs, target_id FROM transactions LIMIT {str(limit)}"
|
||||||
with new_duckdb_locked_connection(self.db_path, read_only=True) as conn:
|
with new_duckdb_locked_connection(self.db_path, read_only=True) as conn:
|
||||||
|
|
|
||||||
|
|
@ -1,16 +1,15 @@
|
||||||
from contextlib import contextmanager
|
from contextlib import contextmanager
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from typing import TYPE_CHECKING, Any, Dict, Optional, Type, Union
|
from typing import TYPE_CHECKING, Any, Dict, Type, Union
|
||||||
|
|
||||||
import duckdb
|
import duckdb
|
||||||
from loguru import logger
|
from loguru import logger
|
||||||
from pydantic import BaseModel
|
from pydantic import BaseModel
|
||||||
|
|
||||||
from langflow.services.deps import get_monitor_service
|
|
||||||
from langflow.utils.concurrency import KeyedWorkerLockManager
|
from langflow.utils.concurrency import KeyedWorkerLockManager
|
||||||
|
|
||||||
if TYPE_CHECKING:
|
if TYPE_CHECKING:
|
||||||
from langflow.api.v1.schemas import ResultDataResponse
|
pass
|
||||||
|
|
||||||
|
|
||||||
INDEX_KEY = "index"
|
INDEX_KEY = "index"
|
||||||
|
|
@ -124,52 +123,3 @@ def add_row_to_table(
|
||||||
logger.error(f"Error adding row to {table_name}: {column_error_message}")
|
logger.error(f"Error adding row to {table_name}: {column_error_message}")
|
||||||
else:
|
else:
|
||||||
logger.error(f"Error adding row to {table_name}: {e}")
|
logger.error(f"Error adding row to {table_name}: {e}")
|
||||||
|
|
||||||
|
|
||||||
async def log_message(
|
|
||||||
sender: str,
|
|
||||||
sender_name: str,
|
|
||||||
message: str,
|
|
||||||
session_id: str,
|
|
||||||
files: Optional[list] = None,
|
|
||||||
flow_id: Optional[str] = None,
|
|
||||||
):
|
|
||||||
try:
|
|
||||||
monitor_service = get_monitor_service()
|
|
||||||
row = {
|
|
||||||
"sender": sender,
|
|
||||||
"sender_name": sender_name,
|
|
||||||
"message": message,
|
|
||||||
"files": files or [],
|
|
||||||
"session_id": session_id,
|
|
||||||
"timestamp": monitor_service.get_timestamp(),
|
|
||||||
"flow_id": flow_id,
|
|
||||||
}
|
|
||||||
monitor_service.add_row(table_name="messages", data=row)
|
|
||||||
except Exception as e:
|
|
||||||
logger.error(f"Error logging message: {e}")
|
|
||||||
|
|
||||||
|
|
||||||
async def log_vertex_build(
|
|
||||||
flow_id: str,
|
|
||||||
vertex_id: str,
|
|
||||||
valid: bool,
|
|
||||||
params: Any,
|
|
||||||
data: "ResultDataResponse",
|
|
||||||
artifacts: Optional[dict] = None,
|
|
||||||
):
|
|
||||||
try:
|
|
||||||
monitor_service = get_monitor_service()
|
|
||||||
|
|
||||||
row = {
|
|
||||||
"flow_id": flow_id,
|
|
||||||
"id": vertex_id,
|
|
||||||
"valid": valid,
|
|
||||||
"params": params,
|
|
||||||
"data": data.model_dump(),
|
|
||||||
"artifacts": artifacts or {},
|
|
||||||
"timestamp": monitor_service.get_timestamp(),
|
|
||||||
}
|
|
||||||
monitor_service.add_row(table_name="vertex_builds", data=row)
|
|
||||||
except Exception as e:
|
|
||||||
logger.exception(f"Error logging vertex build: {e}")
|
|
||||||
|
|
|
||||||
|
|
@ -7,10 +7,10 @@ from sqlmodel import select
|
||||||
from langflow.api.utils import format_elapsed_time
|
from langflow.api.utils import format_elapsed_time
|
||||||
from langflow.api.v1.schemas import ResultDataResponse, VertexBuildResponse
|
from langflow.api.v1.schemas import ResultDataResponse, VertexBuildResponse
|
||||||
from langflow.graph.graph.base import Graph
|
from langflow.graph.graph.base import Graph
|
||||||
|
from langflow.graph.utils import log_vertex_build
|
||||||
from langflow.graph.vertex.base import Vertex
|
from langflow.graph.vertex.base import Vertex
|
||||||
from langflow.services.database.models.flow.model import Flow
|
from langflow.services.database.models.flow.model import Flow
|
||||||
from langflow.services.deps import get_session
|
from langflow.services.deps import get_session
|
||||||
from langflow.services.monitor.utils import log_vertex_build
|
|
||||||
|
|
||||||
|
|
||||||
def set_socketio_server(socketio_server):
|
def set_socketio_server(socketio_server):
|
||||||
|
|
@ -86,7 +86,7 @@ async def build_vertex(
|
||||||
result_dict = ResultDataResponse(results={})
|
result_dict = ResultDataResponse(results={})
|
||||||
artifacts = {}
|
artifacts = {}
|
||||||
set_cache(flow_id, graph)
|
set_cache(flow_id, graph)
|
||||||
await log_vertex_build(
|
log_vertex_build(
|
||||||
flow_id=flow_id,
|
flow_id=flow_id,
|
||||||
vertex_id=vertex_id,
|
vertex_id=vertex_id,
|
||||||
valid=valid,
|
valid=valid,
|
||||||
|
|
|
||||||
|
|
@ -329,3 +329,14 @@ def test_migrate_transactions(client: TestClient):
|
||||||
with session_scope() as session:
|
with session_scope() as session:
|
||||||
new_trans = get_transactions_by_flow_id(session, UUID(flow_id))
|
new_trans = get_transactions_by_flow_id(session, UUID(flow_id))
|
||||||
assert 0 == len(new_trans)
|
assert 0 == len(new_trans)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.load_flows
|
||||||
|
def test_migrate_transactions_no_duckdb(client: TestClient):
|
||||||
|
flow_id = "c54f9130-f2fa-4a3e-b22a-3856d946351b"
|
||||||
|
get_monitor_service()
|
||||||
|
|
||||||
|
with session_scope() as session:
|
||||||
|
migrate_transactions_from_monitor_service_to_database(session)
|
||||||
|
new_trans = get_transactions_by_flow_id(session, UUID(flow_id))
|
||||||
|
assert 0 == len(new_trans)
|
||||||
|
|
|
||||||
|
|
@ -985,16 +985,11 @@ export async function downloadImage({ flowId, fileName }): Promise<any> {
|
||||||
|
|
||||||
export async function getFlowPool({
|
export async function getFlowPool({
|
||||||
flowId,
|
flowId,
|
||||||
nodeId,
|
|
||||||
}: {
|
}: {
|
||||||
flowId: string;
|
flowId: string;
|
||||||
nodeId?: string;
|
|
||||||
}): Promise<AxiosResponse<{ vertex_builds: FlowPoolType }>> {
|
}): Promise<AxiosResponse<{ vertex_builds: FlowPoolType }>> {
|
||||||
const config = {};
|
const config = {};
|
||||||
config["params"] = { flow_id: flowId };
|
config["params"] = { flow_id: flowId };
|
||||||
if (nodeId) {
|
|
||||||
config["params"] = { nodeId };
|
|
||||||
}
|
|
||||||
return await api.get(`${BASE_URL_API}monitor/builds`, config);
|
return await api.get(`${BASE_URL_API}monitor/builds`, config);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue