refactor: remove state manager from Graph class (#8258)
* refactor: remove state management from Graph class * Remove unused state mgmt methods * [autofix.ci] apply automated fixes * refactor: remove GraphStateManager class and its methods This commit deletes the GraphStateManager class, which was responsible for managing the state within the graph. The removal is part of a cleanup effort to streamline the codebase and eliminate unused components. * refactor: remove ListenComponent and NotifyComponent classes This commit deletes the ListenComponent and NotifyComponent classes, which were previously responsible for handling notifications within the system. The removal is part of a codebase cleanup to eliminate unused components and streamline the architecture. --------- Co-authored-by: Jordan Frazier <jordan.frazier@datastax.com> Co-authored-by: autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com> Co-authored-by: Jordan Frazier <122494242+jordanrfrazier@users.noreply.github.com>
This commit is contained in:
parent
9954f7fa0c
commit
092c12f7b6
7 changed files with 260 additions and 568 deletions
|
|
@ -1,29 +0,0 @@
|
|||
from langflow.custom.custom_component.custom_component import CustomComponent
|
||||
from langflow.schema.data import Data
|
||||
|
||||
|
||||
class ListenComponent(CustomComponent):
|
||||
display_name = "Listen"
|
||||
description = "A component to listen for a notification."
|
||||
name = "Listen"
|
||||
beta: bool = True
|
||||
icon = "Radio"
|
||||
|
||||
def build_config(self):
|
||||
return {
|
||||
"name": {
|
||||
"display_name": "Name",
|
||||
"info": "The name of the notification to listen for.",
|
||||
},
|
||||
}
|
||||
|
||||
def build(self, name: str) -> Data:
|
||||
state = self.get_state(name)
|
||||
self._set_successors_ids()
|
||||
self.status = state
|
||||
return state
|
||||
|
||||
def _set_successors_ids(self):
|
||||
self._vertex.is_state = True
|
||||
successors = self._vertex.graph.successor_map.get(self._vertex.id, [])
|
||||
return successors + self._vertex.graph.activated_vertices
|
||||
|
|
@ -1,46 +0,0 @@
|
|||
from langflow.custom.custom_component.custom_component import CustomComponent
|
||||
from langflow.schema.data import Data
|
||||
|
||||
|
||||
class NotifyComponent(CustomComponent):
|
||||
display_name = "Notify"
|
||||
description = "A component to generate a notification to Get Notified component."
|
||||
icon = "Notify"
|
||||
name = "Notify"
|
||||
beta: bool = True
|
||||
|
||||
def build_config(self):
|
||||
return {
|
||||
"name": {"display_name": "Name", "info": "The name of the notification."},
|
||||
"data": {"display_name": "Data", "info": "The data to store."},
|
||||
"append": {
|
||||
"display_name": "Append",
|
||||
"info": "If True, the record will be appended to the notification.",
|
||||
},
|
||||
}
|
||||
|
||||
def build(self, name: str, *, data: Data | None = None, append: bool = False) -> Data:
|
||||
if data and not isinstance(data, Data):
|
||||
if isinstance(data, str):
|
||||
data = Data(text=data)
|
||||
elif isinstance(data, dict):
|
||||
data = Data(data=data)
|
||||
else:
|
||||
data = Data(text=str(data))
|
||||
elif not data:
|
||||
data = Data(text="")
|
||||
if data:
|
||||
if append:
|
||||
self.append_state(name, data)
|
||||
else:
|
||||
self.update_state(name, data)
|
||||
else:
|
||||
self.status = "No record provided."
|
||||
self.status = data
|
||||
self._set_successors_ids()
|
||||
return data
|
||||
|
||||
def _set_successors_ids(self):
|
||||
self._vertex.is_state = True
|
||||
successors = self._vertex.graph.successor_map.get(self._vertex.id, [])
|
||||
return successors + self._vertex.graph.activated_vertices
|
||||
|
|
@ -116,16 +116,6 @@ class CustomComponent(BaseComponent):
|
|||
return f"{self.display_name} ({self._id})"
|
||||
return f"{self.display_name}"
|
||||
|
||||
def update_state(self, name: str, value: Any) -> None:
|
||||
if not self._vertex:
|
||||
msg = "Vertex is not set"
|
||||
raise ValueError(msg)
|
||||
try:
|
||||
self._vertex.graph.update_state(name=name, record=value, caller=self._vertex.id)
|
||||
except Exception as e:
|
||||
msg = f"Error updating state: {e}"
|
||||
raise ValueError(msg) from e
|
||||
|
||||
def stop(self, output_name: str | None = None) -> None:
|
||||
if not output_name and self._vertex and len(self._vertex.outputs) == 1:
|
||||
output_name = self._vertex.outputs[0]["name"]
|
||||
|
|
@ -156,26 +146,6 @@ class CustomComponent(BaseComponent):
|
|||
msg = f"Error starting {self.display_name}: {e}"
|
||||
raise ValueError(msg) from e
|
||||
|
||||
def append_state(self, name: str, value: Any) -> None:
|
||||
if not self._vertex:
|
||||
msg = "Vertex is not set"
|
||||
raise ValueError(msg)
|
||||
try:
|
||||
self._vertex.graph.append_state(name=name, record=value, caller=self._vertex.id)
|
||||
except Exception as e:
|
||||
msg = f"Error appending state: {e}"
|
||||
raise ValueError(msg) from e
|
||||
|
||||
def get_state(self, name: str):
|
||||
if not self._vertex:
|
||||
msg = "Vertex is not set"
|
||||
raise ValueError(msg)
|
||||
try:
|
||||
return self._vertex.graph.get_state(name=name)
|
||||
except Exception as e:
|
||||
msg = f"Error getting state: {e}"
|
||||
raise ValueError(msg) from e
|
||||
|
||||
@staticmethod
|
||||
def resolve_path(path: str) -> str:
|
||||
"""Resolves the path to an absolute path."""
|
||||
|
|
|
|||
|
|
@ -22,7 +22,6 @@ from langflow.graph.edge.base import CycleEdge, Edge
|
|||
from langflow.graph.graph.constants import Finish, lazy_load_vertex_dict
|
||||
from langflow.graph.graph.runnable_vertices_manager import RunnableVerticesManager
|
||||
from langflow.graph.graph.schema import GraphData, GraphDump, StartConfigDict, VertexBuildResult
|
||||
from langflow.graph.graph.state_manager import GraphStateManager
|
||||
from langflow.graph.graph.state_model import create_state_model_from_graph
|
||||
from langflow.graph.graph.utils import (
|
||||
find_all_cycle_edges,
|
||||
|
|
@ -52,7 +51,6 @@ if TYPE_CHECKING:
|
|||
from langflow.events.event_manager import EventManager
|
||||
from langflow.graph.edge.schema import EdgeData
|
||||
from langflow.graph.schema import ResultData
|
||||
from langflow.schema.data import Data
|
||||
from langflow.services.chat.schema import GetCache, SetCache
|
||||
from langflow.services.tracing.service import TracingService
|
||||
|
||||
|
|
@ -113,7 +111,6 @@ class Graph:
|
|||
self.edges: list[CycleEdge] = []
|
||||
self.vertices: list[Vertex] = []
|
||||
self.run_manager = RunnableVerticesManager()
|
||||
self.state_manager = GraphStateManager()
|
||||
self._vertices: list[NodeData] = []
|
||||
self._edges: list[EdgeData] = []
|
||||
|
||||
|
|
@ -494,34 +491,6 @@ class Graph:
|
|||
self.build_graph_maps(self.edges)
|
||||
self.define_vertices_lists()
|
||||
|
||||
def get_state(self, name: str) -> Data | None:
|
||||
"""Returns the state of the graph with the given name.
|
||||
|
||||
Args:
|
||||
name (str): The name of the state.
|
||||
|
||||
Returns:
|
||||
Optional[Data]: The state record, or None if the state does not exist.
|
||||
"""
|
||||
return self.state_manager.get_state(name, run_id=self._run_id)
|
||||
|
||||
def update_state(self, name: str, record: str | Data, caller: str | None = None) -> None:
|
||||
"""Updates the state of the graph with the given name.
|
||||
|
||||
Args:
|
||||
name (str): The name of the state.
|
||||
record (Union[str, Data]): The new state record.
|
||||
caller (Optional[str], optional): The ID of the vertex that is updating the state. Defaults to None.
|
||||
"""
|
||||
if caller:
|
||||
# If there is a caller which is a vertex_id, I want to activate
|
||||
# all StateVertex in self.vertices that are not the caller
|
||||
# essentially notifying all the other vertices that the state has changed
|
||||
# This also has to activate their successors
|
||||
self.activate_state_vertices(name, caller)
|
||||
|
||||
self.state_manager.update_state(name, record, run_id=self._run_id)
|
||||
|
||||
def activate_state_vertices(self, name: str, caller: str) -> None:
|
||||
"""Activates the state vertices in the graph with the given name and caller.
|
||||
|
||||
|
|
@ -575,19 +544,6 @@ class Graph:
|
|||
"""Resets the activated vertices in the graph."""
|
||||
self.activated_vertices = []
|
||||
|
||||
def append_state(self, name: str, record: str | Data, caller: str | None = None) -> None:
|
||||
"""Appends the state of the graph with the given name.
|
||||
|
||||
Args:
|
||||
name (str): The name of the state.
|
||||
record (Union[str, Data]): The state record to append.
|
||||
caller (Optional[str], optional): The ID of the vertex that is updating the state. Defaults to None.
|
||||
"""
|
||||
if caller:
|
||||
self.activate_state_vertices(name, caller)
|
||||
|
||||
self.state_manager.append_state(name, record, run_id=self._run_id)
|
||||
|
||||
def validate_stream(self) -> None:
|
||||
"""Validates the stream configuration of the graph.
|
||||
|
||||
|
|
@ -1056,7 +1012,6 @@ class Graph:
|
|||
state["run_manager"] = RunnableVerticesManager.from_dict(run_manager)
|
||||
self.__dict__.update(state)
|
||||
self.vertex_map = {vertex.id: vertex for vertex in self.vertices}
|
||||
self.state_manager = GraphStateManager()
|
||||
self.tracing_service = get_tracing_service()
|
||||
self.set_run_id(self._run_id)
|
||||
|
||||
|
|
|
|||
|
|
@ -1,38 +0,0 @@
|
|||
from __future__ import annotations
|
||||
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
from loguru import logger
|
||||
|
||||
from langflow.services.deps import get_settings_service, get_state_service
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from collections.abc import Callable
|
||||
|
||||
from langflow.services.state.service import StateService
|
||||
|
||||
|
||||
class GraphStateManager:
|
||||
def __init__(self) -> None:
|
||||
try:
|
||||
self.state_service: StateService = get_state_service()
|
||||
except Exception: # noqa: BLE001
|
||||
logger.opt(exception=True).debug("Error getting state service. Defaulting to InMemoryStateService")
|
||||
from langflow.services.state.service import InMemoryStateService
|
||||
|
||||
self.state_service = InMemoryStateService(get_settings_service())
|
||||
|
||||
def append_state(self, key, new_state, run_id: str) -> None:
|
||||
self.state_service.append_state(key, new_state, run_id)
|
||||
|
||||
def update_state(self, key, new_state, run_id: str) -> None:
|
||||
self.state_service.update_state(key, new_state, run_id)
|
||||
|
||||
def get_state(self, key, run_id: str):
|
||||
return self.state_service.get_state(key, run_id)
|
||||
|
||||
def subscribe(self, key, observer: Callable) -> None:
|
||||
self.state_service.subscribe(key, observer)
|
||||
|
||||
def unsubscribe(self, key, observer: Callable) -> None:
|
||||
self.state_service.unsubscribe(key, observer)
|
||||
|
|
@ -132,12 +132,6 @@ class Vertex:
|
|||
def add_result(self, name: str, result: Any) -> None:
|
||||
self.results[name] = result
|
||||
|
||||
def update_graph_state(self, key, new_state, *, append: bool) -> None:
|
||||
if append:
|
||||
self.graph.append_state(key, new_state, caller=self.id)
|
||||
else:
|
||||
self.graph.update_state(key, new_state, caller=self.id)
|
||||
|
||||
def set_state(self, state: str) -> None:
|
||||
self.state = VertexStates[state]
|
||||
if self.state == VertexStates.INACTIVE and self.graph.in_degree_map[self.id] <= 1:
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue