Underground-Digital/Workflow-Engine
0
1import json2import time3from collections.abc import Sequence4from datetime import datetime, timezone5from typing import Optional6 7from core.app.apps.advanced_chat.app_config_manager import AdvancedChatAppConfigManager8from core.app.apps.workflow.app_config_manager import WorkflowAppConfigManager9from core.model_runtime.utils.encoders import jsonable_encoder10from core.variables import Variable11from core.workflow.entities.node_entities import NodeRunResult12from core.workflow.errors import WorkflowNodeRunFailedError13from core.workflow.nodes import NodeType14from core.workflow.nodes.event import RunCompletedEvent15from core.workflow.nodes.node_mapping import node_type_classes_mapping16from core.workflow.workflow_entry import WorkflowEntry17from events.app_event import app_draft_workflow_was_synced, app_published_workflow_was_updated18from extensions.ext_database import db19from models.account import Account20from models.enums import CreatedByRole21from models.model import App, AppMode22from models.workflow import (23 Workflow,24 WorkflowNodeExecution,25 WorkflowNodeExecutionStatus,26 WorkflowNodeExecutionTriggeredFrom,27 WorkflowType,28)29from services.errors.app import WorkflowHashNotEqualError30from services.workflow.workflow_converter import WorkflowConverter31 32 33class WorkflowService:34 """35 Workflow Service36 """37 38 def get_draft_workflow(self, app_model: App) -> Optional[Workflow]:39 """40 Get draft workflow41 """42 # fetch draft workflow by app_model43 workflow = (44 db.session.query(Workflow)45 .filter(46 Workflow.tenant_id == app_model.tenant_id, Workflow.app_id == app_model.id, Workflow.version == "draft"47 )48 .first()49 )50 51 # return draft workflow52 return workflow53 54 def get_published_workflow(self, app_model: App) -> Optional[Workflow]:55 """56 Get published workflow57 """58 59 if not app_model.workflow_id:60 return None61 62 # fetch published workflow by workflow_id63 workflow = (64 db.session.query(Workflow)65 .filter(66 Workflow.tenant_id == app_model.tenant_id,67 Workflow.app_id == app_model.id,68 Workflow.id == app_model.workflow_id,69 )70 .first()71 )72 73 return workflow74 75 def sync_draft_workflow(76 self,77 *,78 app_model: App,79 graph: dict,80 features: dict,81 unique_hash: Optional[str],82 account: Account,83 environment_variables: Sequence[Variable],84 conversation_variables: Sequence[Variable],85 ) -> Workflow:86 """87 Sync draft workflow88 :raises WorkflowHashNotEqualError89 """90 # fetch draft workflow by app_model91 workflow = self.get_draft_workflow(app_model=app_model)92 93 if workflow and workflow.unique_hash != unique_hash:94 raise WorkflowHashNotEqualError()95 96 # validate features structure97 self.validate_features_structure(app_model=app_model, features=features)98 99 # create draft workflow if not found100 if not workflow:101 workflow = Workflow(102 tenant_id=app_model.tenant_id,103 app_id=app_model.id,104 type=WorkflowType.from_app_mode(app_model.mode).value,105 version="draft",106 graph=json.dumps(graph),107 features=json.dumps(features),108 created_by=account.id,109 environment_variables=environment_variables,110 conversation_variables=conversation_variables,111 )112 db.session.add(workflow)113 # update draft workflow if found114 else:115 workflow.graph = json.dumps(graph)116 workflow.features = json.dumps(features)117 workflow.updated_by = account.id118 workflow.updated_at = datetime.now(timezone.utc).replace(tzinfo=None)119 workflow.environment_variables = environment_variables120 workflow.conversation_variables = conversation_variables121 122 # commit db session changes123 db.session.commit()124 125 # trigger app workflow events126 app_draft_workflow_was_synced.send(app_model, synced_draft_workflow=workflow)127 128 # return draft workflow129 return workflow130 131 def publish_workflow(self, app_model: App, account: Account, draft_workflow: Optional[Workflow] = None) -> Workflow:132 """133 Publish workflow from draft134 135 :param app_model: App instance136 :param account: Account instance137 :param draft_workflow: Workflow instance138 """139 if not draft_workflow:140 # fetch draft workflow by app_model141 draft_workflow = self.get_draft_workflow(app_model=app_model)142 143 if not draft_workflow:144 raise ValueError("No valid workflow found.")145 146 # create new workflow147 workflow = Workflow(148 tenant_id=app_model.tenant_id,149 app_id=app_model.id,150 type=draft_workflow.type,151 version=str(datetime.now(timezone.utc).replace(tzinfo=None)),152 graph=draft_workflow.graph,153 features=draft_workflow.features,154 created_by=account.id,155 environment_variables=draft_workflow.environment_variables,156 conversation_variables=draft_workflow.conversation_variables,157 )158 159 # commit db session changes160 db.session.add(workflow)161 db.session.flush()162 db.session.commit()163 164 app_model.workflow_id = workflow.id165 db.session.commit()166 167 # trigger app workflow events168 app_published_workflow_was_updated.send(app_model, published_workflow=workflow)169 170 # return new workflow171 return workflow172 173 def get_default_block_configs(self) -> list[dict]:174 """175 Get default block configs176 """177 # return default block config178 default_block_configs = []179 for node_type, node_class in node_type_classes_mapping.items():180 default_config = node_class.get_default_config()181 if default_config:182 default_block_configs.append(default_config)183 184 return default_block_configs185 186 def get_default_block_config(self, node_type: str, filters: Optional[dict] = None) -> Optional[dict]:187 """188 Get default config of node.189 :param node_type: node type190 :param filters: filter by node config parameters.191 :return:192 """193 node_type_enum: NodeType = NodeType(node_type)194 195 # return default block config196 node_class = node_type_classes_mapping.get(node_type_enum)197 if not node_class:198 return None199 200 default_config = node_class.get_default_config(filters=filters)201 if not default_config:202 return None203 204 return default_config205 206 def run_draft_workflow_node(207 self, app_model: App, node_id: str, user_inputs: dict, account: Account208 ) -> WorkflowNodeExecution:209 """210 Run draft workflow node211 """212 # fetch draft workflow by app_model213 draft_workflow = self.get_draft_workflow(app_model=app_model)214 if not draft_workflow:215 raise ValueError("Workflow not initialized")216 217 # run draft workflow node218 start_at = time.perf_counter()219 220 try:221 node_instance, generator = WorkflowEntry.single_step_run(222 workflow=draft_workflow,223 node_id=node_id,224 user_inputs=user_inputs,225 user_id=account.id,226 )227 228 node_run_result: NodeRunResult | None = None229 for event in generator:230 if isinstance(event, RunCompletedEvent):231 node_run_result = event.run_result232 233 # sign output files234 node_run_result.outputs = WorkflowEntry.handle_special_values(node_run_result.outputs)235 break236 237 if not node_run_result:238 raise ValueError("Node run failed with no run result")239 240 run_succeeded = True if node_run_result.status == WorkflowNodeExecutionStatus.SUCCEEDED else False241 error = node_run_result.error if not run_succeeded else None242 except WorkflowNodeRunFailedError as e:243 node_instance = e.node_instance244 run_succeeded = False245 node_run_result = None246 error = e.error247 248 workflow_node_execution = WorkflowNodeExecution()249 workflow_node_execution.tenant_id = app_model.tenant_id250 workflow_node_execution.app_id = app_model.id251 workflow_node_execution.workflow_id = draft_workflow.id252 workflow_node_execution.triggered_from = WorkflowNodeExecutionTriggeredFrom.SINGLE_STEP.value253 workflow_node_execution.index = 1254 workflow_node_execution.node_id = node_id255 workflow_node_execution.node_type = node_instance.node_type256 workflow_node_execution.title = node_instance.node_data.title257 workflow_node_execution.elapsed_time = time.perf_counter() - start_at258 workflow_node_execution.created_by_role = CreatedByRole.ACCOUNT.value259 workflow_node_execution.created_by = account.id260 workflow_node_execution.created_at = datetime.now(timezone.utc).replace(tzinfo=None)261 workflow_node_execution.finished_at = datetime.now(timezone.utc).replace(tzinfo=None)262 263 if run_succeeded and node_run_result:264 # create workflow node execution265 workflow_node_execution.inputs = json.dumps(node_run_result.inputs) if node_run_result.inputs else None266 workflow_node_execution.process_data = (267 json.dumps(node_run_result.process_data) if node_run_result.process_data else None268 )269 workflow_node_execution.outputs = (270 json.dumps(jsonable_encoder(node_run_result.outputs)) if node_run_result.outputs else None271 )272 workflow_node_execution.execution_metadata = (273 json.dumps(jsonable_encoder(node_run_result.metadata)) if node_run_result.metadata else None274 )275 workflow_node_execution.status = WorkflowNodeExecutionStatus.SUCCEEDED.value276 else:277 # create workflow node execution278 workflow_node_execution.status = WorkflowNodeExecutionStatus.FAILED.value279 workflow_node_execution.error = error280 281 db.session.add(workflow_node_execution)282 db.session.commit()283 284 return workflow_node_execution285 286 def convert_to_workflow(self, app_model: App, account: Account, args: dict) -> App:287 """288 Basic mode of chatbot app(expert mode) to workflow289 Completion App to Workflow App290 291 :param app_model: App instance292 :param account: Account instance293 :param args: dict294 :return:295 """296 # chatbot convert to workflow mode297 workflow_converter = WorkflowConverter()298 299 if app_model.mode not in {AppMode.CHAT.value, AppMode.COMPLETION.value}:300 raise ValueError(f"Current App mode: {app_model.mode} is not supported convert to workflow.")301 302 # convert to workflow303 new_app = workflow_converter.convert_to_workflow(304 app_model=app_model,305 account=account,306 name=args.get("name"),307 icon_type=args.get("icon_type"),308 icon=args.get("icon"),309 icon_background=args.get("icon_background"),310 )311 312 return new_app313 314 def validate_features_structure(self, app_model: App, features: dict) -> dict:315 if app_model.mode == AppMode.ADVANCED_CHAT.value:316 return AdvancedChatAppConfigManager.config_validate(317 tenant_id=app_model.tenant_id, config=features, only_structure_validate=True318 )319 elif app_model.mode == AppMode.WORKFLOW.value:320 return WorkflowAppConfigManager.config_validate(321 tenant_id=app_model.tenant_id, config=features, only_structure_validate=True322 )323 else:324 raise ValueError(f"Invalid app mode: {app_model.mode}")325 