Underground-Digital/Workflow-Engine
0
1import logging2 3from flask_restful import reqparse4from werkzeug.exceptions import InternalServerError, NotFound5 6import services7from controllers.web import api8from controllers.web.error import (9 AppUnavailableError,10 CompletionRequestError,11 ConversationCompletedError,12 NotChatAppError,13 NotCompletionAppError,14 ProviderModelCurrentlyNotSupportError,15 ProviderNotInitializeError,16 ProviderQuotaExceededError,17)18from controllers.web.error import InvokeRateLimitError as InvokeRateLimitHttpError19from controllers.web.wraps import WebApiResource20from core.app.apps.base_app_queue_manager import AppQueueManager21from core.app.entities.app_invoke_entities import InvokeFrom22from core.errors.error import ModelCurrentlyNotSupportError, ProviderTokenNotInitError, QuotaExceededError23from core.model_runtime.errors.invoke import InvokeError24from libs import helper25from libs.helper import uuid_value26from models.model import AppMode27from services.app_generate_service import AppGenerateService28from services.errors.llm import InvokeRateLimitError29 30 31# define completion api for user32class CompletionApi(WebApiResource):33 def post(self, app_model, end_user):34 if app_model.mode != "completion":35 raise NotCompletionAppError()36 37 parser = reqparse.RequestParser()38 parser.add_argument("inputs", type=dict, required=True, location="json")39 parser.add_argument("query", type=str, location="json", default="")40 parser.add_argument("files", type=list, required=False, location="json")41 parser.add_argument("response_mode", type=str, choices=["blocking", "streaming"], location="json")42 parser.add_argument("retriever_from", type=str, required=False, default="web_app", location="json")43 44 args = parser.parse_args()45 46 streaming = args["response_mode"] == "streaming"47 args["auto_generate_name"] = False48 49 try:50 response = AppGenerateService.generate(51 app_model=app_model, user=end_user, args=args, invoke_from=InvokeFrom.WEB_APP, streaming=streaming52 )53 54 return helper.compact_generate_response(response)55 except services.errors.conversation.ConversationNotExistsError:56 raise NotFound("Conversation Not Exists.")57 except services.errors.conversation.ConversationCompletedError:58 raise ConversationCompletedError()59 except services.errors.app_model_config.AppModelConfigBrokenError:60 logging.exception("App model config broken.")61 raise AppUnavailableError()62 except ProviderTokenNotInitError as ex:63 raise ProviderNotInitializeError(ex.description)64 except QuotaExceededError:65 raise ProviderQuotaExceededError()66 except ModelCurrentlyNotSupportError:67 raise ProviderModelCurrentlyNotSupportError()68 except InvokeError as e:69 raise CompletionRequestError(e.description)70 except ValueError as e:71 raise e72 except Exception as e:73 logging.exception("internal server error.")74 raise InternalServerError()75 76 77class CompletionStopApi(WebApiResource):78 def post(self, app_model, end_user, task_id):79 if app_model.mode != "completion":80 raise NotCompletionAppError()81 82 AppQueueManager.set_stop_flag(task_id, InvokeFrom.WEB_APP, end_user.id)83 84 return {"result": "success"}, 20085 86 87class ChatApi(WebApiResource):88 def post(self, app_model, end_user):89 app_mode = AppMode.value_of(app_model.mode)90 if app_mode not in {AppMode.CHAT, AppMode.AGENT_CHAT, AppMode.ADVANCED_CHAT}:91 raise NotChatAppError()92 93 parser = reqparse.RequestParser()94 parser.add_argument("inputs", type=dict, required=True, location="json")95 parser.add_argument("query", type=str, required=True, location="json")96 parser.add_argument("files", type=list, required=False, location="json")97 parser.add_argument("response_mode", type=str, choices=["blocking", "streaming"], location="json")98 parser.add_argument("conversation_id", type=uuid_value, location="json")99 parser.add_argument("parent_message_id", type=uuid_value, required=False, location="json")100 parser.add_argument("retriever_from", type=str, required=False, default="web_app", location="json")101 102 args = parser.parse_args()103 104 streaming = args["response_mode"] == "streaming"105 args["auto_generate_name"] = False106 107 try:108 response = AppGenerateService.generate(109 app_model=app_model, user=end_user, args=args, invoke_from=InvokeFrom.WEB_APP, streaming=streaming110 )111 112 return helper.compact_generate_response(response)113 except services.errors.conversation.ConversationNotExistsError:114 raise NotFound("Conversation Not Exists.")115 except services.errors.conversation.ConversationCompletedError:116 raise ConversationCompletedError()117 except services.errors.app_model_config.AppModelConfigBrokenError:118 logging.exception("App model config broken.")119 raise AppUnavailableError()120 except ProviderTokenNotInitError as ex:121 raise ProviderNotInitializeError(ex.description)122 except QuotaExceededError:123 raise ProviderQuotaExceededError()124 except ModelCurrentlyNotSupportError:125 raise ProviderModelCurrentlyNotSupportError()126 except InvokeRateLimitError as ex:127 raise InvokeRateLimitHttpError(ex.description)128 except InvokeError as e:129 raise CompletionRequestError(e.description)130 except ValueError as e:131 raise e132 except Exception as e:133 logging.exception("internal server error.")134 raise InternalServerError()135 136 137class ChatStopApi(WebApiResource):138 def post(self, app_model, end_user, task_id):139 app_mode = AppMode.value_of(app_model.mode)140 if app_mode not in {AppMode.CHAT, AppMode.AGENT_CHAT, AppMode.ADVANCED_CHAT}:141 raise NotChatAppError()142 143 AppQueueManager.set_stop_flag(task_id, InvokeFrom.WEB_APP, end_user.id)144 145 return {"result": "success"}, 200146 147 148api.add_resource(CompletionApi, "/completion-messages")149api.add_resource(CompletionStopApi, "/completion-messages/<string:task_id>/stop")150api.add_resource(ChatApi, "/chat-messages")151api.add_resource(ChatStopApi, "/chat-messages/<string:task_id>/stop")152 