Team Ai
Apppublic

Underground-Digital/Workflow-Engine

sourceHugging Faceupdated 2y agoView on Hugging Face
0likes
workflow.py474 linesDownload Raw Back to app
1import json2import logging3 4from flask import abort, request5from flask_restful import Resource, marshal_with, reqparse6from werkzeug.exceptions import Forbidden, InternalServerError, NotFound7 8import services9from controllers.console import api10from controllers.console.app.error import ConversationCompletedError, DraftWorkflowNotExist, DraftWorkflowNotSync11from controllers.console.app.wraps import get_app_model12from controllers.console.wraps import account_initialization_required, setup_required13from core.app.apps.base_app_queue_manager import AppQueueManager14from core.app.entities.app_invoke_entities import InvokeFrom15from factories import variable_factory16from fields.workflow_fields import workflow_fields17from fields.workflow_run_fields import workflow_run_node_execution_fields18from libs import helper19from libs.helper import TimestampField, uuid_value20from libs.login import current_user, login_required21from models import App22from models.model import AppMode23from services.app_dsl_service import AppDslService24from services.app_generate_service import AppGenerateService25from services.errors.app import WorkflowHashNotEqualError26from services.workflow_service import WorkflowService27 28logger = logging.getLogger(__name__)29 30 31class DraftWorkflowApi(Resource):32    @setup_required33    @login_required34    @account_initialization_required35    @get_app_model(mode=[AppMode.ADVANCED_CHAT, AppMode.WORKFLOW])36    @marshal_with(workflow_fields)37    def get(self, app_model: App):38        """39        Get draft workflow40        """41        # The role of the current user in the ta table must be admin, owner, or editor42        if not current_user.is_editor:43            raise Forbidden()44 45        # fetch draft workflow by app_model46        workflow_service = WorkflowService()47        workflow = workflow_service.get_draft_workflow(app_model=app_model)48 49        if not workflow:50            raise DraftWorkflowNotExist()51 52        # return workflow, if not found, return None (initiate graph by frontend)53        return workflow54 55    @setup_required56    @login_required57    @account_initialization_required58    @get_app_model(mode=[AppMode.ADVANCED_CHAT, AppMode.WORKFLOW])59    def post(self, app_model: App):60        """61        Sync draft workflow62        """63        # The role of the current user in the ta table must be admin, owner, or editor64        if not current_user.is_editor:65            raise Forbidden()66 67        content_type = request.headers.get("Content-Type", "")68 69        if "application/json" in content_type:70            parser = reqparse.RequestParser()71            parser.add_argument("graph", type=dict, required=True, nullable=False, location="json")72            parser.add_argument("features", type=dict, required=True, nullable=False, location="json")73            parser.add_argument("hash", type=str, required=False, location="json")74            # TODO: set this to required=True after frontend is updated75            parser.add_argument("environment_variables", type=list, required=False, location="json")76            parser.add_argument("conversation_variables", type=list, required=False, location="json")77            args = parser.parse_args()78        elif "text/plain" in content_type:79            try:80                data = json.loads(request.data.decode("utf-8"))81                if "graph" not in data or "features" not in data:82                    raise ValueError("graph or features not found in data")83 84                if not isinstance(data.get("graph"), dict) or not isinstance(data.get("features"), dict):85                    raise ValueError("graph or features is not a dict")86 87                args = {88                    "graph": data.get("graph"),89                    "features": data.get("features"),90                    "hash": data.get("hash"),91                    "environment_variables": data.get("environment_variables"),92                    "conversation_variables": data.get("conversation_variables"),93                }94            except json.JSONDecodeError:95                return {"message": "Invalid JSON data"}, 40096        else:97            abort(415)98 99        workflow_service = WorkflowService()100 101        try:102            environment_variables_list = args.get("environment_variables") or []103            environment_variables = [104                variable_factory.build_variable_from_mapping(obj) for obj in environment_variables_list105            ]106            conversation_variables_list = args.get("conversation_variables") or []107            conversation_variables = [108                variable_factory.build_variable_from_mapping(obj) for obj in conversation_variables_list109            ]110            workflow = workflow_service.sync_draft_workflow(111                app_model=app_model,112                graph=args["graph"],113                features=args["features"],114                unique_hash=args.get("hash"),115                account=current_user,116                environment_variables=environment_variables,117                conversation_variables=conversation_variables,118            )119        except WorkflowHashNotEqualError:120            raise DraftWorkflowNotSync()121 122        return {123            "result": "success",124            "hash": workflow.unique_hash,125            "updated_at": TimestampField().format(workflow.updated_at or workflow.created_at),126        }127 128 129class DraftWorkflowImportApi(Resource):130    @setup_required131    @login_required132    @account_initialization_required133    @get_app_model(mode=[AppMode.ADVANCED_CHAT, AppMode.WORKFLOW])134    @marshal_with(workflow_fields)135    def post(self, app_model: App):136        """137        Import draft workflow138        """139        # The role of the current user in the ta table must be admin, owner, or editor140        if not current_user.is_editor:141            raise Forbidden()142 143        parser = reqparse.RequestParser()144        parser.add_argument("data", type=str, required=True, nullable=False, location="json")145        args = parser.parse_args()146 147        workflow = AppDslService.import_and_overwrite_workflow(148            app_model=app_model, data=args["data"], account=current_user149        )150 151        return workflow152 153 154class AdvancedChatDraftWorkflowRunApi(Resource):155    @setup_required156    @login_required157    @account_initialization_required158    @get_app_model(mode=[AppMode.ADVANCED_CHAT])159    def post(self, app_model: App):160        """161        Run draft workflow162        """163        # The role of the current user in the ta table must be admin, owner, or editor164        if not current_user.is_editor:165            raise Forbidden()166 167        parser = reqparse.RequestParser()168        parser.add_argument("inputs", type=dict, location="json")169        parser.add_argument("query", type=str, required=True, location="json", default="")170        parser.add_argument("files", type=list, location="json")171        parser.add_argument("conversation_id", type=uuid_value, location="json")172        parser.add_argument("parent_message_id", type=uuid_value, required=False, location="json")173 174        args = parser.parse_args()175 176        try:177            response = AppGenerateService.generate(178                app_model=app_model, user=current_user, args=args, invoke_from=InvokeFrom.DEBUGGER, streaming=True179            )180 181            return helper.compact_generate_response(response)182        except services.errors.conversation.ConversationNotExistsError:183            raise NotFound("Conversation Not Exists.")184        except services.errors.conversation.ConversationCompletedError:185            raise ConversationCompletedError()186        except ValueError as e:187            raise e188        except Exception as e:189            logging.exception("internal server error.")190            raise InternalServerError()191 192 193class AdvancedChatDraftRunIterationNodeApi(Resource):194    @setup_required195    @login_required196    @account_initialization_required197    @get_app_model(mode=[AppMode.ADVANCED_CHAT])198    def post(self, app_model: App, node_id: str):199        """200        Run draft workflow iteration node201        """202        # The role of the current user in the ta table must be admin, owner, or editor203        if not current_user.is_editor:204            raise Forbidden()205 206        parser = reqparse.RequestParser()207        parser.add_argument("inputs", type=dict, location="json")208        args = parser.parse_args()209 210        try:211            response = AppGenerateService.generate_single_iteration(212                app_model=app_model, user=current_user, node_id=node_id, args=args, streaming=True213            )214 215            return helper.compact_generate_response(response)216        except services.errors.conversation.ConversationNotExistsError:217            raise NotFound("Conversation Not Exists.")218        except services.errors.conversation.ConversationCompletedError:219            raise ConversationCompletedError()220        except ValueError as e:221            raise e222        except Exception as e:223            logging.exception("internal server error.")224            raise InternalServerError()225 226 227class WorkflowDraftRunIterationNodeApi(Resource):228    @setup_required229    @login_required230    @account_initialization_required231    @get_app_model(mode=[AppMode.WORKFLOW])232    def post(self, app_model: App, node_id: str):233        """234        Run draft workflow iteration node235        """236        # The role of the current user in the ta table must be admin, owner, or editor237        if not current_user.is_editor:238            raise Forbidden()239 240        parser = reqparse.RequestParser()241        parser.add_argument("inputs", type=dict, location="json")242        args = parser.parse_args()243 244        try:245            response = AppGenerateService.generate_single_iteration(246                app_model=app_model, user=current_user, node_id=node_id, args=args, streaming=True247            )248 249            return helper.compact_generate_response(response)250        except services.errors.conversation.ConversationNotExistsError:251            raise NotFound("Conversation Not Exists.")252        except services.errors.conversation.ConversationCompletedError:253            raise ConversationCompletedError()254        except ValueError as e:255            raise e256        except Exception as e:257            logging.exception("internal server error.")258            raise InternalServerError()259 260 261class DraftWorkflowRunApi(Resource):262    @setup_required263    @login_required264    @account_initialization_required265    @get_app_model(mode=[AppMode.WORKFLOW])266    def post(self, app_model: App):267        """268        Run draft workflow269        """270        # The role of the current user in the ta table must be admin, owner, or editor271        if not current_user.is_editor:272            raise Forbidden()273 274        parser = reqparse.RequestParser()275        parser.add_argument("inputs", type=dict, required=True, nullable=False, location="json")276        parser.add_argument("files", type=list, required=False, location="json")277        args = parser.parse_args()278 279        response = AppGenerateService.generate(280            app_model=app_model,281            user=current_user,282            args=args,283            invoke_from=InvokeFrom.DEBUGGER,284            streaming=True,285        )286 287        return helper.compact_generate_response(response)288 289 290class WorkflowTaskStopApi(Resource):291    @setup_required292    @login_required293    @account_initialization_required294    @get_app_model(mode=[AppMode.ADVANCED_CHAT, AppMode.WORKFLOW])295    def post(self, app_model: App, task_id: str):296        """297        Stop workflow task298        """299        # The role of the current user in the ta table must be admin, owner, or editor300        if not current_user.is_editor:301            raise Forbidden()302 303        AppQueueManager.set_stop_flag(task_id, InvokeFrom.DEBUGGER, current_user.id)304 305        return {"result": "success"}306 307 308class DraftWorkflowNodeRunApi(Resource):309    @setup_required310    @login_required311    @account_initialization_required312    @get_app_model(mode=[AppMode.ADVANCED_CHAT, AppMode.WORKFLOW])313    @marshal_with(workflow_run_node_execution_fields)314    def post(self, app_model: App, node_id: str):315        """316        Run draft workflow node317        """318        # The role of the current user in the ta table must be admin, owner, or editor319        if not current_user.is_editor:320            raise Forbidden()321 322        parser = reqparse.RequestParser()323        parser.add_argument("inputs", type=dict, required=True, nullable=False, location="json")324        args = parser.parse_args()325 326        workflow_service = WorkflowService()327        workflow_node_execution = workflow_service.run_draft_workflow_node(328            app_model=app_model, node_id=node_id, user_inputs=args.get("inputs"), account=current_user329        )330 331        return workflow_node_execution332 333 334class PublishedWorkflowApi(Resource):335    @setup_required336    @login_required337    @account_initialization_required338    @get_app_model(mode=[AppMode.ADVANCED_CHAT, AppMode.WORKFLOW])339    @marshal_with(workflow_fields)340    def get(self, app_model: App):341        """342        Get published workflow343        """344        # The role of the current user in the ta table must be admin, owner, or editor345        if not current_user.is_editor:346            raise Forbidden()347 348        # fetch published workflow by app_model349        workflow_service = WorkflowService()350        workflow = workflow_service.get_published_workflow(app_model=app_model)351 352        # return workflow, if not found, return None353        return workflow354 355    @setup_required356    @login_required357    @account_initialization_required358    @get_app_model(mode=[AppMode.ADVANCED_CHAT, AppMode.WORKFLOW])359    def post(self, app_model: App):360        """361        Publish workflow362        """363        # The role of the current user in the ta table must be admin, owner, or editor364        if not current_user.is_editor:365            raise Forbidden()366 367        workflow_service = WorkflowService()368        workflow = workflow_service.publish_workflow(app_model=app_model, account=current_user)369 370        return {"result": "success", "created_at": TimestampField().format(workflow.created_at)}371 372 373class DefaultBlockConfigsApi(Resource):374    @setup_required375    @login_required376    @account_initialization_required377    @get_app_model(mode=[AppMode.ADVANCED_CHAT, AppMode.WORKFLOW])378    def get(self, app_model: App):379        """380        Get default block config381        """382        # The role of the current user in the ta table must be admin, owner, or editor383        if not current_user.is_editor:384            raise Forbidden()385 386        # Get default block configs387        workflow_service = WorkflowService()388        return workflow_service.get_default_block_configs()389 390 391class DefaultBlockConfigApi(Resource):392    @setup_required393    @login_required394    @account_initialization_required395    @get_app_model(mode=[AppMode.ADVANCED_CHAT, AppMode.WORKFLOW])396    def get(self, app_model: App, block_type: str):397        """398        Get default block config399        """400        # The role of the current user in the ta table must be admin, owner, or editor401        if not current_user.is_editor:402            raise Forbidden()403 404        parser = reqparse.RequestParser()405        parser.add_argument("q", type=str, location="args")406        args = parser.parse_args()407 408        filters = None409        if args.get("q"):410            try:411                filters = json.loads(args.get("q"))412            except json.JSONDecodeError:413                raise ValueError("Invalid filters")414 415        # Get default block configs416        workflow_service = WorkflowService()417        return workflow_service.get_default_block_config(node_type=block_type, filters=filters)418 419 420class ConvertToWorkflowApi(Resource):421    @setup_required422    @login_required423    @account_initialization_required424    @get_app_model(mode=[AppMode.CHAT, AppMode.COMPLETION])425    def post(self, app_model: App):426        """427        Convert basic mode of chatbot app to workflow mode428        Convert expert mode of chatbot app to workflow mode429        Convert Completion App to Workflow App430        """431        # The role of the current user in the ta table must be admin, owner, or editor432        if not current_user.is_editor:433            raise Forbidden()434 435        if request.data:436            parser = reqparse.RequestParser()437            parser.add_argument("name", type=str, required=False, nullable=True, location="json")438            parser.add_argument("icon_type", type=str, required=False, nullable=True, location="json")439            parser.add_argument("icon", type=str, required=False, nullable=True, location="json")440            parser.add_argument("icon_background", type=str, required=False, nullable=True, location="json")441            args = parser.parse_args()442        else:443            args = {}444 445        # convert to workflow mode446        workflow_service = WorkflowService()447        new_app_model = workflow_service.convert_to_workflow(app_model=app_model, account=current_user, args=args)448 449        # return app id450        return {451            "new_app_id": new_app_model.id,452        }453 454 455api.add_resource(DraftWorkflowApi, "/apps/<uuid:app_id>/workflows/draft")456api.add_resource(DraftWorkflowImportApi, "/apps/<uuid:app_id>/workflows/draft/import")457api.add_resource(AdvancedChatDraftWorkflowRunApi, "/apps/<uuid:app_id>/advanced-chat/workflows/draft/run")458api.add_resource(DraftWorkflowRunApi, "/apps/<uuid:app_id>/workflows/draft/run")459api.add_resource(WorkflowTaskStopApi, "/apps/<uuid:app_id>/workflow-runs/tasks/<string:task_id>/stop")460api.add_resource(DraftWorkflowNodeRunApi, "/apps/<uuid:app_id>/workflows/draft/nodes/<string:node_id>/run")461api.add_resource(462    AdvancedChatDraftRunIterationNodeApi,463    "/apps/<uuid:app_id>/advanced-chat/workflows/draft/iteration/nodes/<string:node_id>/run",464)465api.add_resource(466    WorkflowDraftRunIterationNodeApi, "/apps/<uuid:app_id>/workflows/draft/iteration/nodes/<string:node_id>/run"467)468api.add_resource(PublishedWorkflowApi, "/apps/<uuid:app_id>/workflows/publish")469api.add_resource(DefaultBlockConfigsApi, "/apps/<uuid:app_id>/workflows/default-workflow-block-configs")470api.add_resource(471    DefaultBlockConfigApi, "/apps/<uuid:app_id>/workflows/default-workflow-block-configs/<string:block_type>"472)473api.add_resource(ConvertToWorkflowApi, "/apps/<uuid:app_id>/convert-to-workflow")474