Merge branch 'dev' into fixGroupLogs
This commit is contained in:
commit
4e28b2a09c
2 changed files with 17 additions and 9 deletions
|
|
@ -1444,7 +1444,7 @@ class Graph:
|
||||||
|
|
||||||
def is_vertex_runnable(self, vertex_id: str) -> bool:
|
def is_vertex_runnable(self, vertex_id: str) -> bool:
|
||||||
"""Returns whether a vertex is runnable."""
|
"""Returns whether a vertex is runnable."""
|
||||||
return self.run_manager.is_vertex_runnable(vertex_id)
|
return self.run_manager.is_vertex_runnable(vertex_id, self.inactivated_vertices)
|
||||||
|
|
||||||
def build_run_map(self):
|
def build_run_map(self):
|
||||||
"""
|
"""
|
||||||
|
|
@ -1464,7 +1464,7 @@ class Graph:
|
||||||
This checks the direct predecessors of each successor to identify any that are
|
This checks the direct predecessors of each successor to identify any that are
|
||||||
immediately runnable, expanding the search to ensure progress can be made.
|
immediately runnable, expanding the search to ensure progress can be made.
|
||||||
"""
|
"""
|
||||||
return self.run_manager.find_runnable_predecessors_for_successors(vertex_id)
|
return self.run_manager.find_runnable_predecessors_for_successors(vertex_id, self.inactivated_vertices)
|
||||||
|
|
||||||
def remove_from_predecessors(self, vertex_id: str):
|
def remove_from_predecessors(self, vertex_id: str):
|
||||||
self.run_manager.remove_from_predecessors(vertex_id)
|
self.run_manager.remove_from_predecessors(vertex_id)
|
||||||
|
|
|
||||||
|
|
@ -1,6 +1,6 @@
|
||||||
import asyncio
|
import asyncio
|
||||||
from collections import defaultdict
|
from collections import defaultdict
|
||||||
from typing import TYPE_CHECKING, Callable, List, Coroutine
|
from typing import TYPE_CHECKING, Callable, Coroutine, List
|
||||||
|
|
||||||
if TYPE_CHECKING:
|
if TYPE_CHECKING:
|
||||||
from langflow.graph.graph.base import Graph
|
from langflow.graph.graph.base import Graph
|
||||||
|
|
@ -40,19 +40,23 @@ class RunnableVerticesManager:
|
||||||
self.run_predecessors = state["run_predecessors"]
|
self.run_predecessors = state["run_predecessors"]
|
||||||
self.vertices_to_run = state["vertices_to_run"]
|
self.vertices_to_run = state["vertices_to_run"]
|
||||||
|
|
||||||
def is_vertex_runnable(self, vertex_id: str) -> bool:
|
def is_vertex_runnable(self, vertex_id: str, inactivated_vertices: set[str]) -> bool:
|
||||||
"""Determines if a vertex is runnable."""
|
"""Determines if a vertex is runnable."""
|
||||||
|
|
||||||
return vertex_id in self.vertices_to_run and not self.run_predecessors.get(vertex_id)
|
return (
|
||||||
|
vertex_id in self.vertices_to_run
|
||||||
|
and not self.run_predecessors.get(vertex_id)
|
||||||
|
and vertex_id not in inactivated_vertices
|
||||||
|
)
|
||||||
|
|
||||||
def find_runnable_predecessors_for_successors(self, vertex_id: str) -> List[str]:
|
def find_runnable_predecessors_for_successors(self, vertex_id: str, inactivated_vertices: set[str]) -> List[str]:
|
||||||
"""Finds runnable predecessors for the successors of a given vertex."""
|
"""Finds runnable predecessors for the successors of a given vertex."""
|
||||||
runnable_vertices = []
|
runnable_vertices = []
|
||||||
visited = set()
|
visited = set()
|
||||||
|
|
||||||
for successor_id in self.run_map.get(vertex_id, []):
|
for successor_id in self.run_map.get(vertex_id, []):
|
||||||
for predecessor_id in self.run_predecessors.get(successor_id, []):
|
for predecessor_id in self.run_predecessors.get(successor_id, []):
|
||||||
if predecessor_id not in visited and self.is_vertex_runnable(predecessor_id):
|
if predecessor_id not in visited and self.is_vertex_runnable(predecessor_id, inactivated_vertices):
|
||||||
runnable_vertices.append(predecessor_id)
|
runnable_vertices.append(predecessor_id)
|
||||||
visited.add(predecessor_id)
|
visited.add(predecessor_id)
|
||||||
return runnable_vertices
|
return runnable_vertices
|
||||||
|
|
@ -104,10 +108,14 @@ class RunnableVerticesManager:
|
||||||
"""
|
"""
|
||||||
async with lock:
|
async with lock:
|
||||||
self.remove_from_predecessors(vertex.id)
|
self.remove_from_predecessors(vertex.id)
|
||||||
direct_successors_ready = [v for v in vertex.successors_ids if self.is_vertex_runnable(v)]
|
direct_successors_ready = [
|
||||||
|
v for v in vertex.successors_ids if self.is_vertex_runnable(v, graph.inactivated_vertices)
|
||||||
|
]
|
||||||
if not direct_successors_ready:
|
if not direct_successors_ready:
|
||||||
# No direct successors ready, look for runnable predecessors of successors
|
# No direct successors ready, look for runnable predecessors of successors
|
||||||
next_runnable_vertices = self.find_runnable_predecessors_for_successors(vertex.id)
|
next_runnable_vertices = self.find_runnable_predecessors_for_successors(
|
||||||
|
vertex.id, graph.inactivated_vertices
|
||||||
|
)
|
||||||
else:
|
else:
|
||||||
next_runnable_vertices = direct_successors_ready
|
next_runnable_vertices = direct_successors_ready
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue