Underground-Digital/Workflow-Engine
0
1from collections.abc import Generator, Mapping2from typing import Any, Union3 4from openai._exceptions import RateLimitError5 6from configs import dify_config7from core.app.apps.advanced_chat.app_generator import AdvancedChatAppGenerator8from core.app.apps.agent_chat.app_generator import AgentChatAppGenerator9from core.app.apps.chat.app_generator import ChatAppGenerator10from core.app.apps.completion.app_generator import CompletionAppGenerator11from core.app.apps.workflow.app_generator import WorkflowAppGenerator12from core.app.entities.app_invoke_entities import InvokeFrom13from core.app.features.rate_limiting import RateLimit14from models.model import Account, App, AppMode, EndUser15from models.workflow import Workflow16from services.errors.llm import InvokeRateLimitError17from services.workflow_service import WorkflowService18 19 20class AppGenerateService:21 @classmethod22 def generate(23 cls,24 app_model: App,25 user: Union[Account, EndUser],26 args: Mapping[str, Any],27 invoke_from: InvokeFrom,28 streaming: bool = True,29 ):30 """31 App Content Generate32 :param app_model: app model33 :param user: user34 :param args: args35 :param invoke_from: invoke from36 :param streaming: streaming37 :return:38 """39 max_active_request = AppGenerateService._get_max_active_requests(app_model)40 rate_limit = RateLimit(app_model.id, max_active_request)41 request_id = RateLimit.gen_request_key()42 try:43 request_id = rate_limit.enter(request_id)44 if app_model.mode == AppMode.COMPLETION.value:45 return rate_limit.generate(46 CompletionAppGenerator().generate(47 app_model=app_model, user=user, args=args, invoke_from=invoke_from, stream=streaming48 ),49 request_id,50 )51 elif app_model.mode == AppMode.AGENT_CHAT.value or app_model.is_agent:52 return rate_limit.generate(53 AgentChatAppGenerator().generate(54 app_model=app_model, user=user, args=args, invoke_from=invoke_from, stream=streaming55 ),56 request_id,57 )58 elif app_model.mode == AppMode.CHAT.value:59 return rate_limit.generate(60 ChatAppGenerator().generate(61 app_model=app_model, user=user, args=args, invoke_from=invoke_from, stream=streaming62 ),63 request_id,64 )65 elif app_model.mode == AppMode.ADVANCED_CHAT.value:66 workflow = cls._get_workflow(app_model, invoke_from)67 return rate_limit.generate(68 AdvancedChatAppGenerator().generate(69 app_model=app_model,70 workflow=workflow,71 user=user,72 args=args,73 invoke_from=invoke_from,74 stream=streaming,75 ),76 request_id,77 )78 elif app_model.mode == AppMode.WORKFLOW.value:79 workflow = cls._get_workflow(app_model, invoke_from)80 return rate_limit.generate(81 WorkflowAppGenerator().generate(82 app_model=app_model,83 workflow=workflow,84 user=user,85 args=args,86 invoke_from=invoke_from,87 stream=streaming,88 ),89 request_id,90 )91 else:92 raise ValueError(f"Invalid app mode {app_model.mode}")93 except RateLimitError as e:94 raise InvokeRateLimitError(str(e))95 finally:96 if not streaming:97 rate_limit.exit(request_id)98 99 @staticmethod100 def _get_max_active_requests(app_model: App) -> int:101 max_active_requests = app_model.max_active_requests102 if app_model.max_active_requests is None:103 max_active_requests = int(dify_config.APP_MAX_ACTIVE_REQUESTS)104 return max_active_requests105 106 @classmethod107 def generate_single_iteration(cls, app_model: App, user: Account, node_id: str, args: Any, streaming: bool = True):108 if app_model.mode == AppMode.ADVANCED_CHAT.value:109 workflow = cls._get_workflow(app_model, InvokeFrom.DEBUGGER)110 return AdvancedChatAppGenerator().single_iteration_generate(111 app_model=app_model, workflow=workflow, node_id=node_id, user=user, args=args, stream=streaming112 )113 elif app_model.mode == AppMode.WORKFLOW.value:114 workflow = cls._get_workflow(app_model, InvokeFrom.DEBUGGER)115 return WorkflowAppGenerator().single_iteration_generate(116 app_model=app_model, workflow=workflow, node_id=node_id, user=user, args=args, stream=streaming117 )118 else:119 raise ValueError(f"Invalid app mode {app_model.mode}")120 121 @classmethod122 def generate_more_like_this(123 cls,124 app_model: App,125 user: Union[Account, EndUser],126 message_id: str,127 invoke_from: InvokeFrom,128 streaming: bool = True,129 ) -> Union[dict, Generator]:130 """131 Generate more like this132 :param app_model: app model133 :param user: user134 :param message_id: message id135 :param invoke_from: invoke from136 :param streaming: streaming137 :return:138 """139 return CompletionAppGenerator().generate_more_like_this(140 app_model=app_model, message_id=message_id, user=user, invoke_from=invoke_from, stream=streaming141 )142 143 @classmethod144 def _get_workflow(cls, app_model: App, invoke_from: InvokeFrom) -> Workflow:145 """146 Get workflow147 :param app_model: app model148 :param invoke_from: invoke from149 :return:150 """151 workflow_service = WorkflowService()152 if invoke_from == InvokeFrom.DEBUGGER:153 # fetch draft workflow by app_model154 workflow = workflow_service.get_draft_workflow(app_model=app_model)155 156 if not workflow:157 raise ValueError("Workflow not initialized")158 else:159 # fetch published workflow by app_model160 workflow = workflow_service.get_published_workflow(app_model=app_model)161 162 if not workflow:163 raise ValueError("Workflow not published")164 165 return workflow166 