Underground-Digital/Workflow-Engine
0
1from extensions.ext_database import db2from libs.infinite_scroll_pagination import InfiniteScrollPagination3from models.enums import WorkflowRunTriggeredFrom4from models.model import App5from models.workflow import (6 WorkflowNodeExecution,7 WorkflowNodeExecutionTriggeredFrom,8 WorkflowRun,9)10 11 12class WorkflowRunService:13 def get_paginate_advanced_chat_workflow_runs(self, app_model: App, args: dict) -> InfiniteScrollPagination:14 """15 Get advanced chat app workflow run list16 Only return triggered_from == advanced_chat17 18 :param app_model: app model19 :param args: request args20 """21 22 class WorkflowWithMessage:23 message_id: str24 conversation_id: str25 26 def __init__(self, workflow_run: WorkflowRun):27 self._workflow_run = workflow_run28 29 def __getattr__(self, item):30 return getattr(self._workflow_run, item)31 32 pagination = self.get_paginate_workflow_runs(app_model, args)33 34 with_message_workflow_runs = []35 for workflow_run in pagination.data:36 message = workflow_run.message37 with_message_workflow_run = WorkflowWithMessage(workflow_run=workflow_run)38 if message:39 with_message_workflow_run.message_id = message.id40 with_message_workflow_run.conversation_id = message.conversation_id41 42 with_message_workflow_runs.append(with_message_workflow_run)43 44 pagination.data = with_message_workflow_runs45 return pagination46 47 def get_paginate_workflow_runs(self, app_model: App, args: dict) -> InfiniteScrollPagination:48 """49 Get debug workflow run list50 Only return triggered_from == debugging51 52 :param app_model: app model53 :param args: request args54 """55 limit = int(args.get("limit", 20))56 57 base_query = db.session.query(WorkflowRun).filter(58 WorkflowRun.tenant_id == app_model.tenant_id,59 WorkflowRun.app_id == app_model.id,60 WorkflowRun.triggered_from == WorkflowRunTriggeredFrom.DEBUGGING.value,61 )62 63 if args.get("last_id"):64 last_workflow_run = base_query.filter(65 WorkflowRun.id == args.get("last_id"),66 ).first()67 68 if not last_workflow_run:69 raise ValueError("Last workflow run not exists")70 71 workflow_runs = (72 base_query.filter(73 WorkflowRun.created_at < last_workflow_run.created_at, WorkflowRun.id != last_workflow_run.id74 )75 .order_by(WorkflowRun.created_at.desc())76 .limit(limit)77 .all()78 )79 else:80 workflow_runs = base_query.order_by(WorkflowRun.created_at.desc()).limit(limit).all()81 82 has_more = False83 if len(workflow_runs) == limit:84 current_page_first_workflow_run = workflow_runs[-1]85 rest_count = base_query.filter(86 WorkflowRun.created_at < current_page_first_workflow_run.created_at,87 WorkflowRun.id != current_page_first_workflow_run.id,88 ).count()89 90 if rest_count > 0:91 has_more = True92 93 return InfiniteScrollPagination(data=workflow_runs, limit=limit, has_more=has_more)94 95 def get_workflow_run(self, app_model: App, run_id: str) -> WorkflowRun:96 """97 Get workflow run detail98 99 :param app_model: app model100 :param run_id: workflow run id101 """102 workflow_run = (103 db.session.query(WorkflowRun)104 .filter(105 WorkflowRun.tenant_id == app_model.tenant_id,106 WorkflowRun.app_id == app_model.id,107 WorkflowRun.id == run_id,108 )109 .first()110 )111 112 return workflow_run113 114 def get_workflow_run_node_executions(self, app_model: App, run_id: str) -> list[WorkflowNodeExecution]:115 """116 Get workflow run node execution list117 """118 workflow_run = self.get_workflow_run(app_model, run_id)119 120 if not workflow_run:121 return []122 123 node_executions = (124 db.session.query(WorkflowNodeExecution)125 .filter(126 WorkflowNodeExecution.tenant_id == app_model.tenant_id,127 WorkflowNodeExecution.app_id == app_model.id,128 WorkflowNodeExecution.workflow_id == workflow_run.workflow_id,129 WorkflowNodeExecution.triggered_from == WorkflowNodeExecutionTriggeredFrom.WORKFLOW_RUN.value,130 WorkflowNodeExecution.workflow_run_id == run_id,131 )132 .order_by(WorkflowNodeExecution.index.desc())133 .all()134 )135 136 return node_executions137 