Team Ai
Apppublic

Underground-Digital/Workflow-Engine

sourceHugging Faceupdated 2y agoView on Hugging Face
0likes
workflow_run_service.py137 linesDownload Raw Back to services
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