Underground-Digital/Workflow-Engine
0
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 