fix: update condition to run end_all_traces (#2707)
This commit is contained in:
parent
a3959650fc
commit
537e358b65
3 changed files with 6 additions and 1 deletions
|
|
@ -260,7 +260,7 @@ async def build_vertex(
|
||||||
if graph.stop_vertex and graph.stop_vertex in next_runnable_vertices:
|
if graph.stop_vertex and graph.stop_vertex in next_runnable_vertices:
|
||||||
next_runnable_vertices = [graph.stop_vertex]
|
next_runnable_vertices = [graph.stop_vertex]
|
||||||
|
|
||||||
if not graph.run_manager.vertices_to_run and not next_runnable_vertices:
|
if graph.run_manager.all_predecessors_are_fulfilled() and not next_runnable_vertices:
|
||||||
background_tasks.add_task(graph.end_all_traces)
|
background_tasks.add_task(graph.end_all_traces)
|
||||||
|
|
||||||
build_response = VertexBuildResponse(
|
build_response = VertexBuildResponse(
|
||||||
|
|
|
||||||
|
|
@ -539,6 +539,8 @@ class Graph:
|
||||||
"""Marks a vertex in the graph."""
|
"""Marks a vertex in the graph."""
|
||||||
vertex = self.get_vertex(vertex_id)
|
vertex = self.get_vertex(vertex_id)
|
||||||
vertex.set_state(state)
|
vertex.set_state(state)
|
||||||
|
if state == VertexStates.INACTIVE:
|
||||||
|
self.run_manager.remove_from_predecessors(vertex_id)
|
||||||
|
|
||||||
def mark_branch(self, vertex_id: str, state: str, visited: Optional[set] = None, output_name: Optional[str] = None):
|
def mark_branch(self, vertex_id: str, state: str, visited: Optional[set] = None, output_name: Optional[str] = None):
|
||||||
"""Marks a branch of the graph."""
|
"""Marks a branch of the graph."""
|
||||||
|
|
|
||||||
|
|
@ -41,6 +41,9 @@ 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 all_predecessors_are_fulfilled(self) -> bool:
|
||||||
|
return all(not value for value in self.run_predecessors.values())
|
||||||
|
|
||||||
def update_run_state(self, run_predecessors: dict, vertices_to_run: set):
|
def update_run_state(self, run_predecessors: dict, vertices_to_run: set):
|
||||||
self.run_predecessors.update(run_predecessors)
|
self.run_predecessors.update(run_predecessors)
|
||||||
self.vertices_to_run.update(vertices_to_run)
|
self.vertices_to_run.update(vertices_to_run)
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue