fix: flow_runner better run_id and support stream (#7892)
Co-authored-by: Gabriel Luiz Freitas Almeida <gabriel@langflow.org>
This commit is contained in:
parent
36c9f6fa15
commit
21f9182a82
3 changed files with 14 additions and 4 deletions
|
|
@ -641,7 +641,7 @@ class Graph:
|
||||||
raise ValueError(msg)
|
raise ValueError(msg)
|
||||||
return self._run_id
|
return self._run_id
|
||||||
|
|
||||||
def set_run_id(self, run_id: uuid.UUID | None = None) -> None:
|
def set_run_id(self, run_id: uuid.UUID | str | None = None) -> None:
|
||||||
"""Sets the ID of the current run.
|
"""Sets the ID of the current run.
|
||||||
|
|
||||||
Args:
|
Args:
|
||||||
|
|
|
||||||
|
|
@ -71,6 +71,7 @@ async def run_graph(
|
||||||
session_id: str | None = None,
|
session_id: str | None = None,
|
||||||
fallback_to_env_vars: bool = False,
|
fallback_to_env_vars: bool = False,
|
||||||
output_component: str | None = None,
|
output_component: str | None = None,
|
||||||
|
stream: bool = False,
|
||||||
) -> list[RunOutputs]:
|
) -> list[RunOutputs]:
|
||||||
"""Runs the given Langflow Graph with the specified input and returns the outputs.
|
"""Runs the given Langflow Graph with the specified input and returns the outputs.
|
||||||
|
|
||||||
|
|
@ -83,6 +84,7 @@ async def run_graph(
|
||||||
fallback_to_env_vars (bool, optional): Whether to fallback to environment variables.
|
fallback_to_env_vars (bool, optional): Whether to fallback to environment variables.
|
||||||
Defaults to False.
|
Defaults to False.
|
||||||
output_component (Optional[str], optional): The specific output component to retrieve. Defaults to None.
|
output_component (Optional[str], optional): The specific output component to retrieve. Defaults to None.
|
||||||
|
stream (bool, optional): Whether to stream the results or not. Defaults to False.
|
||||||
|
|
||||||
Returns:
|
Returns:
|
||||||
List[RunOutputs]: A list of RunOutputs objects representing the outputs of the graph.
|
List[RunOutputs]: A list of RunOutputs objects representing the outputs of the graph.
|
||||||
|
|
@ -113,7 +115,7 @@ async def run_graph(
|
||||||
inputs_components=components,
|
inputs_components=components,
|
||||||
types=types,
|
types=types,
|
||||||
outputs=outputs or [],
|
outputs=outputs or [],
|
||||||
stream=False,
|
stream=stream,
|
||||||
session_id=session_id,
|
session_id=session_id,
|
||||||
fallback_to_env_vars=fallback_to_env_vars,
|
fallback_to_env_vars=fallback_to_env_vars,
|
||||||
)
|
)
|
||||||
|
|
|
||||||
|
|
@ -43,8 +43,10 @@ class LangflowRunnerExperimental:
|
||||||
session_id: str, # UUID required currently
|
session_id: str, # UUID required currently
|
||||||
flow: Path | str | dict,
|
flow: Path | str | dict,
|
||||||
input_value: str,
|
input_value: str,
|
||||||
|
*,
|
||||||
input_type: str = "chat",
|
input_type: str = "chat",
|
||||||
output_type: str = "chat",
|
output_type: str = "chat",
|
||||||
|
stream: bool = False,
|
||||||
):
|
):
|
||||||
logger.info(f"Start Handling {session_id=}")
|
logger.info(f"Start Handling {session_id=}")
|
||||||
await self.init_db_if_needed()
|
await self.init_db_if_needed()
|
||||||
|
|
@ -56,7 +58,7 @@ class LangflowRunnerExperimental:
|
||||||
await self.add_flow_to_db(session_id, flow_dict)
|
await self.add_flow_to_db(session_id, flow_dict)
|
||||||
graph = await self.create_graph_from_flow(session_id, flow_dict)
|
graph = await self.create_graph_from_flow(session_id, flow_dict)
|
||||||
try:
|
try:
|
||||||
result = await self.run_graph(input_value, input_type, output_type, session_id, graph)
|
result = await self.run_graph(input_value, input_type, output_type, session_id, graph, stream=stream)
|
||||||
finally:
|
finally:
|
||||||
await self.clear_flow_state(session_id, flow_dict)
|
await self.clear_flow_state(session_id, flow_dict)
|
||||||
logger.info(f"Finish Handling {session_id=}")
|
logger.info(f"Finish Handling {session_id=}")
|
||||||
|
|
@ -74,7 +76,9 @@ class LangflowRunnerExperimental:
|
||||||
await session.commit()
|
await session.commit()
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
async def run_graph(input_value: str, input_type: str, output_type: str, session_id: str, graph: Graph):
|
async def run_graph(
|
||||||
|
input_value: str, input_type: str, output_type: str, session_id: str, graph: Graph, *, stream: bool
|
||||||
|
):
|
||||||
return await run_graph(
|
return await run_graph(
|
||||||
graph=graph,
|
graph=graph,
|
||||||
session_id=session_id,
|
session_id=session_id,
|
||||||
|
|
@ -82,13 +86,17 @@ class LangflowRunnerExperimental:
|
||||||
fallback_to_env_vars=True,
|
fallback_to_env_vars=True,
|
||||||
input_type=input_type,
|
input_type=input_type,
|
||||||
output_type=output_type,
|
output_type=output_type,
|
||||||
|
stream=stream,
|
||||||
)
|
)
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
async def create_graph_from_flow(session_id: str, flow_dict: dict):
|
async def create_graph_from_flow(session_id: str, flow_dict: dict):
|
||||||
graph = await aload_flow_from_json(flow=flow_dict, disable_logs=False)
|
graph = await aload_flow_from_json(flow=flow_dict, disable_logs=False)
|
||||||
graph.flow_id = flow_dict["id"]
|
graph.flow_id = flow_dict["id"]
|
||||||
|
graph.flow_name = flow_dict.get("name")
|
||||||
graph.session_id = session_id
|
graph.session_id = session_id
|
||||||
|
graph.set_run_id(session_id)
|
||||||
|
await graph.initialize_run()
|
||||||
return graph
|
return graph
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue