Team Ai
Apppublic

Underground-Digital/Workflow-Engine

sourceHugging Faceupdated 2y agoView on Hugging Face
0likes
graph_engine.py742 linesDownload Raw Back to graph_engine
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 
Underground-Digital/Workflow-Engine · Team Ai