Underground-Digital/Workflow-Engine
0
1import logging2import queue3import time4import uuid5from collections.abc import Generator, Mapping6from concurrent.futures import ThreadPoolExecutor, wait7from copy import copy, deepcopy8from typing import Any, Optional9 10from flask import Flask, current_app11 12from core.app.apps.base_app_queue_manager import GenerateTaskStoppedError13from core.app.entities.app_invoke_entities import InvokeFrom14from core.workflow.entities.node_entities import NodeRunMetadataKey15from core.workflow.entities.variable_pool import VariablePool, VariableValue16from core.workflow.graph_engine.condition_handlers.condition_manager import ConditionManager17from core.workflow.graph_engine.entities.event import (18 BaseIterationEvent,19 GraphEngineEvent,20 GraphRunFailedEvent,21 GraphRunStartedEvent,22 GraphRunSucceededEvent,23 NodeRunFailedEvent,24 NodeRunRetrieverResourceEvent,25 NodeRunStartedEvent,26 NodeRunStreamChunkEvent,27 NodeRunSucceededEvent,28 ParallelBranchRunFailedEvent,29 ParallelBranchRunStartedEvent,30 ParallelBranchRunSucceededEvent,31)32from core.workflow.graph_engine.entities.graph import Graph, GraphEdge33from core.workflow.graph_engine.entities.graph_init_params import GraphInitParams34from core.workflow.graph_engine.entities.graph_runtime_state import GraphRuntimeState35from core.workflow.graph_engine.entities.runtime_route_state import RouteNodeState36from core.workflow.nodes import NodeType37from core.workflow.nodes.answer.answer_stream_processor import AnswerStreamProcessor38from core.workflow.nodes.base import BaseNode39from core.workflow.nodes.end.end_stream_processor import EndStreamProcessor40from core.workflow.nodes.event import RunCompletedEvent, RunRetrieverResourceEvent, RunStreamChunkEvent41from core.workflow.nodes.node_mapping import node_type_classes_mapping42from extensions.ext_database import db43from models.enums import UserFrom44from models.workflow import WorkflowNodeExecutionStatus, WorkflowType45 46logger = logging.getLogger(__name__)47 48 49class GraphEngineThreadPool(ThreadPoolExecutor):50 def __init__(51 self, max_workers=None, thread_name_prefix="", initializer=None, initargs=(), max_submit_count=10052 ) -> None:53 super().__init__(max_workers, thread_name_prefix, initializer, initargs)54 self.max_submit_count = max_submit_count55 self.submit_count = 056 57 def submit(self, fn, *args, **kwargs):58 self.submit_count += 159 self.check_is_full()60 61 return super().submit(fn, *args, **kwargs)62 63 def task_done_callback(self, future):64 self.submit_count -= 165 66 def check_is_full(self) -> None:67 print(f"submit_count: {self.submit_count}, max_submit_count: {self.max_submit_count}")68 if self.submit_count > self.max_submit_count:69 raise ValueError(f"Max submit count {self.max_submit_count} of workflow thread pool reached.")70 71 72class GraphEngine:73 workflow_thread_pool_mapping: dict[str, GraphEngineThreadPool] = {}74 75 def __init__(76 self,77 tenant_id: str,78 app_id: str,79 workflow_type: WorkflowType,80 workflow_id: str,81 user_id: str,82 user_from: UserFrom,83 invoke_from: InvokeFrom,84 call_depth: int,85 graph: Graph,86 graph_config: Mapping[str, Any],87 variable_pool: VariablePool,88 max_execution_steps: int,89 max_execution_time: int,90 thread_pool_id: Optional[str] = None,91 ) -> None:92 thread_pool_max_submit_count = 10093 thread_pool_max_workers = 1094 95 # init thread pool96 if thread_pool_id:97 if thread_pool_id not in GraphEngine.workflow_thread_pool_mapping:98 raise ValueError(f"Max submit count {thread_pool_max_submit_count} of workflow thread pool reached.")99 100 self.thread_pool_id = thread_pool_id101 self.thread_pool = GraphEngine.workflow_thread_pool_mapping[thread_pool_id]102 self.is_main_thread_pool = False103 else:104 self.thread_pool = GraphEngineThreadPool(105 max_workers=thread_pool_max_workers, max_submit_count=thread_pool_max_submit_count106 )107 self.thread_pool_id = str(uuid.uuid4())108 self.is_main_thread_pool = True109 GraphEngine.workflow_thread_pool_mapping[self.thread_pool_id] = self.thread_pool110 111 self.graph = graph112 self.init_params = GraphInitParams(113 tenant_id=tenant_id,114 app_id=app_id,115 workflow_type=workflow_type,116 workflow_id=workflow_id,117 graph_config=graph_config,118 user_id=user_id,119 user_from=user_from,120 invoke_from=invoke_from,121 call_depth=call_depth,122 )123 124 self.graph_runtime_state = GraphRuntimeState(variable_pool=variable_pool, start_at=time.perf_counter())125 126 self.max_execution_steps = max_execution_steps127 self.max_execution_time = max_execution_time128 129 def run(self) -> Generator[GraphEngineEvent, None, None]:130 # trigger graph run start event131 yield GraphRunStartedEvent()132 133 try:134 if self.init_params.workflow_type == WorkflowType.CHAT:135 stream_processor = AnswerStreamProcessor(136 graph=self.graph, variable_pool=self.graph_runtime_state.variable_pool137 )138 else:139 stream_processor = EndStreamProcessor(140 graph=self.graph, variable_pool=self.graph_runtime_state.variable_pool141 )142 143 # run graph144 generator = stream_processor.process(self._run(start_node_id=self.graph.root_node_id))145 146 for item in generator:147 try:148 yield item149 if isinstance(item, NodeRunFailedEvent):150 yield GraphRunFailedEvent(error=item.route_node_state.failed_reason or "Unknown error.")151 return152 elif isinstance(item, NodeRunSucceededEvent):153 if item.node_type == NodeType.END:154 self.graph_runtime_state.outputs = (155 item.route_node_state.node_run_result.outputs156 if item.route_node_state.node_run_result157 and item.route_node_state.node_run_result.outputs158 else {}159 )160 elif item.node_type == NodeType.ANSWER:161 if "answer" not in self.graph_runtime_state.outputs:162 self.graph_runtime_state.outputs["answer"] = ""163 164 self.graph_runtime_state.outputs["answer"] += "\n" + (165 item.route_node_state.node_run_result.outputs.get("answer", "")166 if item.route_node_state.node_run_result167 and item.route_node_state.node_run_result.outputs168 else ""169 )170 171 self.graph_runtime_state.outputs["answer"] = self.graph_runtime_state.outputs[172 "answer"173 ].strip()174 except Exception as e:175 logger.exception(f"Graph run failed: {str(e)}")176 yield GraphRunFailedEvent(error=str(e))177 return178 179 # trigger graph run success event180 yield GraphRunSucceededEvent(outputs=self.graph_runtime_state.outputs)181 self._release_thread()182 except GraphRunFailedError as e:183 yield GraphRunFailedEvent(error=e.error)184 self._release_thread()185 return186 except Exception as e:187 logger.exception("Unknown Error when graph running")188 yield GraphRunFailedEvent(error=str(e))189 self._release_thread()190 raise e191 192 def _release_thread(self):193 if self.is_main_thread_pool and self.thread_pool_id in GraphEngine.workflow_thread_pool_mapping:194 del GraphEngine.workflow_thread_pool_mapping[self.thread_pool_id]195 196 def _run(197 self,198 start_node_id: str,199 in_parallel_id: Optional[str] = None,200 parent_parallel_id: Optional[str] = None,201 parent_parallel_start_node_id: Optional[str] = None,202 ) -> Generator[GraphEngineEvent, None, None]:203 parallel_start_node_id = None204 if in_parallel_id:205 parallel_start_node_id = start_node_id206 207 next_node_id = start_node_id208 previous_route_node_state: Optional[RouteNodeState] = None209 while True:210 # max steps reached211 if self.graph_runtime_state.node_run_steps > self.max_execution_steps:212 raise GraphRunFailedError("Max steps {} reached.".format(self.max_execution_steps))213 214 # or max execution time reached215 if self._is_timed_out(216 start_at=self.graph_runtime_state.start_at, max_execution_time=self.max_execution_time217 ):218 raise GraphRunFailedError("Max execution time {}s reached.".format(self.max_execution_time))219 220 # init route node state221 route_node_state = self.graph_runtime_state.node_run_state.create_node_state(node_id=next_node_id)222 223 # get node config224 node_id = route_node_state.node_id225 node_config = self.graph.node_id_config_mapping.get(node_id)226 if not node_config:227 raise GraphRunFailedError(f"Node {node_id} config not found.")228 229 # convert to specific node230 node_type = NodeType(node_config.get("data", {}).get("type"))231 node_cls = node_type_classes_mapping[node_type]232 233 previous_node_id = previous_route_node_state.node_id if previous_route_node_state else None234 235 # init workflow run state236 node_instance = node_cls( # type: ignore237 id=route_node_state.id,238 config=node_config,239 graph_init_params=self.init_params,240 graph=self.graph,241 graph_runtime_state=self.graph_runtime_state,242 previous_node_id=previous_node_id,243 thread_pool_id=self.thread_pool_id,244 )245 246 try:247 # run node248 generator = self._run_node(249 node_instance=node_instance,250 route_node_state=route_node_state,251 parallel_id=in_parallel_id,252 parallel_start_node_id=parallel_start_node_id,253 parent_parallel_id=parent_parallel_id,254 parent_parallel_start_node_id=parent_parallel_start_node_id,255 )256 257 for item in generator:258 if isinstance(item, NodeRunStartedEvent):259 self.graph_runtime_state.node_run_steps += 1260 item.route_node_state.index = self.graph_runtime_state.node_run_steps261 262 yield item263 264 self.graph_runtime_state.node_run_state.node_state_mapping[route_node_state.id] = route_node_state265 266 # append route267 if previous_route_node_state:268 self.graph_runtime_state.node_run_state.add_route(269 source_node_state_id=previous_route_node_state.id, target_node_state_id=route_node_state.id270 )271 except Exception as e:272 route_node_state.status = RouteNodeState.Status.FAILED273 route_node_state.failed_reason = str(e)274 yield NodeRunFailedEvent(275 error=str(e),276 id=node_instance.id,277 node_id=next_node_id,278 node_type=node_type,279 node_data=node_instance.node_data,280 route_node_state=route_node_state,281 parallel_id=in_parallel_id,282 parallel_start_node_id=parallel_start_node_id,283 parent_parallel_id=parent_parallel_id,284 parent_parallel_start_node_id=parent_parallel_start_node_id,285 )286 raise e287 288 # It may not be necessary, but it is necessary. :)289 if (290 self.graph.node_id_config_mapping[next_node_id].get("data", {}).get("type", "").lower()291 == NodeType.END.value292 ):293 break294 295 previous_route_node_state = route_node_state296 297 # get next node ids298 edge_mappings = self.graph.edge_mapping.get(next_node_id)299 if not edge_mappings:300 break301 302 if len(edge_mappings) == 1:303 edge = edge_mappings[0]304 305 if edge.run_condition:306 result = ConditionManager.get_condition_handler(307 init_params=self.init_params,308 graph=self.graph,309 run_condition=edge.run_condition,310 ).check(311 graph_runtime_state=self.graph_runtime_state,312 previous_route_node_state=previous_route_node_state,313 )314 315 if not result:316 break317 318 next_node_id = edge.target_node_id319 else:320 final_node_id = None321 322 if any(edge.run_condition for edge in edge_mappings):323 # if nodes has run conditions, get node id which branch to take based on the run condition results324 condition_edge_mappings = {}325 for edge in edge_mappings:326 if edge.run_condition:327 run_condition_hash = edge.run_condition.hash328 if run_condition_hash not in condition_edge_mappings:329 condition_edge_mappings[run_condition_hash] = []330 331 condition_edge_mappings[run_condition_hash].append(edge)332 333 for _, sub_edge_mappings in condition_edge_mappings.items():334 if len(sub_edge_mappings) == 0:335 continue336 337 edge = sub_edge_mappings[0]338 339 result = ConditionManager.get_condition_handler(340 init_params=self.init_params,341 graph=self.graph,342 run_condition=edge.run_condition,343 ).check(344 graph_runtime_state=self.graph_runtime_state,345 previous_route_node_state=previous_route_node_state,346 )347 348 if not result:349 continue350 351 if len(sub_edge_mappings) == 1:352 final_node_id = edge.target_node_id353 else:354 parallel_generator = self._run_parallel_branches(355 edge_mappings=sub_edge_mappings,356 in_parallel_id=in_parallel_id,357 parallel_start_node_id=parallel_start_node_id,358 )359 360 for item in parallel_generator:361 if isinstance(item, str):362 final_node_id = item363 else:364 yield item365 366 break367 368 if not final_node_id:369 break370 371 next_node_id = final_node_id372 else:373 parallel_generator = self._run_parallel_branches(374 edge_mappings=edge_mappings,375 in_parallel_id=in_parallel_id,376 parallel_start_node_id=parallel_start_node_id,377 )378 379 for item in parallel_generator:380 if isinstance(item, str):381 final_node_id = item382 else:383 yield item384 385 if not final_node_id:386 break387 388 next_node_id = final_node_id389 390 if in_parallel_id and self.graph.node_parallel_mapping.get(next_node_id, "") != in_parallel_id:391 break392 393 def _run_parallel_branches(394 self,395 edge_mappings: list[GraphEdge],396 in_parallel_id: Optional[str] = None,397 parallel_start_node_id: Optional[str] = None,398 ) -> Generator[GraphEngineEvent | str, None, None]:399 # if nodes has no run conditions, parallel run all nodes400 parallel_id = self.graph.node_parallel_mapping.get(edge_mappings[0].target_node_id)401 if not parallel_id:402 node_id = edge_mappings[0].target_node_id403 node_config = self.graph.node_id_config_mapping.get(node_id)404 if not node_config:405 raise GraphRunFailedError(406 f"Node {node_id} related parallel not found or incorrectly connected to multiple parallel branches."407 )408 409 node_title = node_config.get("data", {}).get("title")410 raise GraphRunFailedError(411 f"Node {node_title} related parallel not found or incorrectly connected to multiple parallel branches."412 )413 414 parallel = self.graph.parallel_mapping.get(parallel_id)415 if not parallel:416 raise GraphRunFailedError(f"Parallel {parallel_id} not found.")417 418 # run parallel nodes, run in new thread and use queue to get results419 q: queue.Queue = queue.Queue()420 421 # Create a list to store the threads422 futures = []423 424 # new thread425 for edge in edge_mappings:426 if (427 edge.target_node_id not in self.graph.node_parallel_mapping428 or self.graph.node_parallel_mapping.get(edge.target_node_id, "") != parallel_id429 ):430 continue431 432 future = self.thread_pool.submit(433 self._run_parallel_node,434 **{435 "flask_app": current_app._get_current_object(), # type: ignore[attr-defined]436 "q": q,437 "parallel_id": parallel_id,438 "parallel_start_node_id": edge.target_node_id,439 "parent_parallel_id": in_parallel_id,440 "parent_parallel_start_node_id": parallel_start_node_id,441 },442 )443 444 future.add_done_callback(self.thread_pool.task_done_callback)445 446 futures.append(future)447 448 succeeded_count = 0449 while True:450 try:451 event = q.get(timeout=1)452 if event is None:453 break454 455 yield event456 if event.parallel_id == parallel_id:457 if isinstance(event, ParallelBranchRunSucceededEvent):458 succeeded_count += 1459 if succeeded_count == len(futures):460 q.put(None)461 462 continue463 elif isinstance(event, ParallelBranchRunFailedEvent):464 raise GraphRunFailedError(event.error)465 except queue.Empty:466 continue467 468 # wait all threads469 wait(futures)470 471 # get final node id472 final_node_id = parallel.end_to_node_id473 if final_node_id:474 yield final_node_id475 476 def _run_parallel_node(477 self,478 flask_app: Flask,479 q: queue.Queue,480 parallel_id: str,481 parallel_start_node_id: str,482 parent_parallel_id: Optional[str] = None,483 parent_parallel_start_node_id: Optional[str] = None,484 ) -> None:485 """486 Run parallel nodes487 """488 with flask_app.app_context():489 try:490 q.put(491 ParallelBranchRunStartedEvent(492 parallel_id=parallel_id,493 parallel_start_node_id=parallel_start_node_id,494 parent_parallel_id=parent_parallel_id,495 parent_parallel_start_node_id=parent_parallel_start_node_id,496 )497 )498 499 # run node500 generator = self._run(501 start_node_id=parallel_start_node_id,502 in_parallel_id=parallel_id,503 parent_parallel_id=parent_parallel_id,504 parent_parallel_start_node_id=parent_parallel_start_node_id,505 )506 507 for item in generator:508 q.put(item)509 510 # trigger graph run success event511 q.put(512 ParallelBranchRunSucceededEvent(513 parallel_id=parallel_id,514 parallel_start_node_id=parallel_start_node_id,515 parent_parallel_id=parent_parallel_id,516 parent_parallel_start_node_id=parent_parallel_start_node_id,517 )518 )519 except GraphRunFailedError as e:520 q.put(521 ParallelBranchRunFailedEvent(522 parallel_id=parallel_id,523 parallel_start_node_id=parallel_start_node_id,524 parent_parallel_id=parent_parallel_id,525 parent_parallel_start_node_id=parent_parallel_start_node_id,526 error=e.error,527 )528 )529 except Exception as e:530 logger.exception("Unknown Error when generating in parallel")531 q.put(532 ParallelBranchRunFailedEvent(533 parallel_id=parallel_id,534 parallel_start_node_id=parallel_start_node_id,535 parent_parallel_id=parent_parallel_id,536 parent_parallel_start_node_id=parent_parallel_start_node_id,537 error=str(e),538 )539 )540 finally:541 db.session.remove()542 543 def _run_node(544 self,545 node_instance: BaseNode,546 route_node_state: RouteNodeState,547 parallel_id: Optional[str] = None,548 parallel_start_node_id: Optional[str] = None,549 parent_parallel_id: Optional[str] = None,550 parent_parallel_start_node_id: Optional[str] = None,551 ) -> Generator[GraphEngineEvent, None, None]:552 """553 Run node554 """555 # trigger node run start event556 yield NodeRunStartedEvent(557 id=node_instance.id,558 node_id=node_instance.node_id,559 node_type=node_instance.node_type,560 node_data=node_instance.node_data,561 route_node_state=route_node_state,562 predecessor_node_id=node_instance.previous_node_id,563 parallel_id=parallel_id,564 parallel_start_node_id=parallel_start_node_id,565 parent_parallel_id=parent_parallel_id,566 parent_parallel_start_node_id=parent_parallel_start_node_id,567 )568 569 db.session.close()570 571 try:572 # run node573 generator = node_instance.run()574 for item in generator:575 if isinstance(item, GraphEngineEvent):576 if isinstance(item, BaseIterationEvent):577 # add parallel info to iteration event578 item.parallel_id = parallel_id579 item.parallel_start_node_id = parallel_start_node_id580 item.parent_parallel_id = parent_parallel_id581 item.parent_parallel_start_node_id = parent_parallel_start_node_id582 583 yield item584 else:585 if isinstance(item, RunCompletedEvent):586 run_result = item.run_result587 route_node_state.set_finished(run_result=run_result)588 589 if run_result.status == WorkflowNodeExecutionStatus.FAILED:590 yield NodeRunFailedEvent(591 error=route_node_state.failed_reason or "Unknown error.",592 id=node_instance.id,593 node_id=node_instance.node_id,594 node_type=node_instance.node_type,595 node_data=node_instance.node_data,596 route_node_state=route_node_state,597 parallel_id=parallel_id,598 parallel_start_node_id=parallel_start_node_id,599 parent_parallel_id=parent_parallel_id,600 parent_parallel_start_node_id=parent_parallel_start_node_id,601 )602 elif run_result.status == WorkflowNodeExecutionStatus.SUCCEEDED:603 if run_result.metadata and run_result.metadata.get(NodeRunMetadataKey.TOTAL_TOKENS):604 # plus state total_tokens605 self.graph_runtime_state.total_tokens += int(606 run_result.metadata.get(NodeRunMetadataKey.TOTAL_TOKENS) # type: ignore[arg-type]607 )608 609 if run_result.llm_usage:610 # use the latest usage611 self.graph_runtime_state.llm_usage += run_result.llm_usage612 613 # append node output variables to variable pool614 if run_result.outputs:615 for variable_key, variable_value in run_result.outputs.items():616 # append variables to variable pool recursively617 self._append_variables_recursively(618 node_id=node_instance.node_id,619 variable_key_list=[variable_key],620 variable_value=variable_value,621 )622 623 # add parallel info to run result metadata624 if parallel_id and parallel_start_node_id:625 if not run_result.metadata:626 run_result.metadata = {}627 628 run_result.metadata[NodeRunMetadataKey.PARALLEL_ID] = parallel_id629 run_result.metadata[NodeRunMetadataKey.PARALLEL_START_NODE_ID] = parallel_start_node_id630 if parent_parallel_id and parent_parallel_start_node_id:631 run_result.metadata[NodeRunMetadataKey.PARENT_PARALLEL_ID] = parent_parallel_id632 run_result.metadata[NodeRunMetadataKey.PARENT_PARALLEL_START_NODE_ID] = (633 parent_parallel_start_node_id634 )635 636 yield NodeRunSucceededEvent(637 id=node_instance.id,638 node_id=node_instance.node_id,639 node_type=node_instance.node_type,640 node_data=node_instance.node_data,641 route_node_state=route_node_state,642 parallel_id=parallel_id,643 parallel_start_node_id=parallel_start_node_id,644 parent_parallel_id=parent_parallel_id,645 parent_parallel_start_node_id=parent_parallel_start_node_id,646 )647 648 break649 elif isinstance(item, RunStreamChunkEvent):650 yield NodeRunStreamChunkEvent(651 id=node_instance.id,652 node_id=node_instance.node_id,653 node_type=node_instance.node_type,654 node_data=node_instance.node_data,655 chunk_content=item.chunk_content,656 from_variable_selector=item.from_variable_selector,657 route_node_state=route_node_state,658 parallel_id=parallel_id,659 parallel_start_node_id=parallel_start_node_id,660 parent_parallel_id=parent_parallel_id,661 parent_parallel_start_node_id=parent_parallel_start_node_id,662 )663 elif isinstance(item, RunRetrieverResourceEvent):664 yield NodeRunRetrieverResourceEvent(665 id=node_instance.id,666 node_id=node_instance.node_id,667 node_type=node_instance.node_type,668 node_data=node_instance.node_data,669 retriever_resources=item.retriever_resources,670 context=item.context,671 route_node_state=route_node_state,672 parallel_id=parallel_id,673 parallel_start_node_id=parallel_start_node_id,674 parent_parallel_id=parent_parallel_id,675 parent_parallel_start_node_id=parent_parallel_start_node_id,676 )677 except GenerateTaskStoppedError:678 # trigger node run failed event679 route_node_state.status = RouteNodeState.Status.FAILED680 route_node_state.failed_reason = "Workflow stopped."681 yield NodeRunFailedEvent(682 error="Workflow stopped.",683 id=node_instance.id,684 node_id=node_instance.node_id,685 node_type=node_instance.node_type,686 node_data=node_instance.node_data,687 route_node_state=route_node_state,688 parallel_id=parallel_id,689 parallel_start_node_id=parallel_start_node_id,690 parent_parallel_id=parent_parallel_id,691 parent_parallel_start_node_id=parent_parallel_start_node_id,692 )693 return694 except Exception as e:695 logger.exception(f"Node {node_instance.node_data.title} run failed: {str(e)}")696 raise e697 finally:698 db.session.close()699 700 def _append_variables_recursively(self, node_id: str, variable_key_list: list[str], variable_value: VariableValue):701 """702 Append variables recursively703 :param node_id: node id704 :param variable_key_list: variable key list705 :param variable_value: variable value706 :return:707 """708 self.graph_runtime_state.variable_pool.add([node_id] + variable_key_list, variable_value)709 710 # if variable_value is a dict, then recursively append variables711 if isinstance(variable_value, dict):712 for key, value in variable_value.items():713 # construct new key list714 new_key_list = variable_key_list + [key]715 self._append_variables_recursively(716 node_id=node_id, variable_key_list=new_key_list, variable_value=value717 )718 719 def _is_timed_out(self, start_at: float, max_execution_time: int) -> bool:720 """721 Check timeout722 :param start_at: start time723 :param max_execution_time: max execution time724 :return:725 """726 return time.perf_counter() - start_at > max_execution_time727 728 def create_copy(self):729 """730 create a graph engine copy731 :return: with a new variable pool instance of graph engine732 """733 new_instance = copy(self)734 new_instance.graph_runtime_state = copy(self.graph_runtime_state)735 new_instance.graph_runtime_state.variable_pool = deepcopy(self.graph_runtime_state.variable_pool)736 return new_instance737 738 739class GraphRunFailedError(Exception):740 def __init__(self, error: str):741 self.error = error742 