376 lines
15 KiB
Python
376 lines
15 KiB
Python
#!/usr/bin/env python
|
|
# -*- coding: utf-8 -*-
|
|
"""
|
|
审批节点处理器
|
|
处理审批任务的创建和完成逻辑
|
|
"""
|
|
import logging
|
|
from typing import Dict, Any, Optional, TYPE_CHECKING
|
|
|
|
from sqlalchemy import select
|
|
|
|
from online_dev.workflow.engine.handlers.base_handler import BaseNodeHandler
|
|
|
|
if TYPE_CHECKING:
|
|
from online_dev.workflow.engine.base import ExecutionContext
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
class ApprovalHandler(BaseNodeHandler):
|
|
"""
|
|
审批节点处理器
|
|
|
|
支持:
|
|
- 或签(any):一人通过即可
|
|
- 会签(parallel):所有人都要通过
|
|
- 依次审批(sequential):按顺序审批
|
|
"""
|
|
|
|
async def execute(self, context: 'ExecutionContext', node: Dict) -> None:
|
|
"""
|
|
执行审批节点:创建审批任务
|
|
"""
|
|
from online_dev.workflow.engine.assignee_resolver import assignee_resolver
|
|
|
|
node_id = self.get_node_id(node)
|
|
node_name = self.get_node_name(node)
|
|
node_config = self.get_node_config(node)
|
|
|
|
logger.info(f"创建审批任务 - 节点: {node_id}, 名称: {node_name}")
|
|
|
|
# 解析审批人
|
|
assignee_ids = await assignee_resolver.resolve(
|
|
context.db,
|
|
node_config,
|
|
context.instance,
|
|
context.form_data,
|
|
)
|
|
|
|
logger.info(f"解析到的审批人: {assignee_ids}")
|
|
|
|
if not assignee_ids:
|
|
# 空审批人处理
|
|
empty_action = node_config.get('emptyAssignee', 'error')
|
|
if empty_action == 'skip':
|
|
logger.warning(f"节点 {node_id} 没有审批人,自动跳过")
|
|
await self.advance_to_next(context, node)
|
|
return
|
|
elif empty_action == 'admin':
|
|
logger.warning(f"节点 {node_id} 没有审批人,转交管理员(未实现)")
|
|
await self.advance_to_next(context, node)
|
|
return
|
|
else:
|
|
logger.warning(f"节点 {node_id} 没有审批人,自动跳过")
|
|
await self.advance_to_next(context, node)
|
|
return
|
|
|
|
# 更新实例当前节点
|
|
await self.update_instance_node(context, node)
|
|
|
|
# 创建任务
|
|
multi_approval = node_config.get('multiApproval', 'any')
|
|
|
|
if multi_approval == 'sequential':
|
|
# 依次审批:只创建第一个人的任务
|
|
await self.create_task(context, node, assignee_ids[0], 'approval')
|
|
else:
|
|
# 或签/会签:创建所有人的任务
|
|
for assignee_id in assignee_ids:
|
|
await self.create_task(context, node, assignee_id, 'approval')
|
|
|
|
async def handle_approval(self, context: 'ExecutionContext', task: Any) -> None:
|
|
"""
|
|
处理审批通过后的流程推进
|
|
|
|
Args:
|
|
context: 执行上下文
|
|
task: 已完成的任务
|
|
"""
|
|
from online_dev.workflow.model import WorkflowTask
|
|
from online_dev.workflow.engine.utils import FlowUtils
|
|
|
|
# 检查是否是加签任务
|
|
if task.sign_type:
|
|
await self._handle_sign_task_completion(context, task)
|
|
return
|
|
|
|
node = FlowUtils.find_node_by_id(context.flow_definition, task.node_id)
|
|
if not node:
|
|
logger.error(f"找不到节点: {task.node_id}")
|
|
return
|
|
|
|
node_config = self.get_node_config(node)
|
|
multi_approval = node_config.get('multiApproval', 'any')
|
|
|
|
# 检查多人审批逻辑
|
|
if multi_approval == 'parallel':
|
|
# 会签:检查是否所有人都已审批
|
|
stmt = select(WorkflowTask).where(
|
|
WorkflowTask.instance_id == str(context.instance.id),
|
|
WorkflowTask.node_id == task.node_id,
|
|
WorkflowTask.status == 'pending',
|
|
WorkflowTask.sign_type == '',
|
|
)
|
|
result = await context.db.execute(stmt)
|
|
pending_tasks = result.scalars().all()
|
|
|
|
if len(pending_tasks) > 0:
|
|
logger.info(f"会签模式,还有 {len(pending_tasks)} 人未审批")
|
|
return
|
|
|
|
elif multi_approval == 'sequential':
|
|
# 依次审批:检查是否还有下一个人
|
|
next_assignee = await self._get_next_sequential_assignee(context, task)
|
|
if next_assignee:
|
|
await self.create_task(context, node, next_assignee, 'approval')
|
|
return
|
|
|
|
# 或签模式:取消该节点其他 pending 的普通任务
|
|
if multi_approval == 'any':
|
|
stmt = select(WorkflowTask).where(
|
|
WorkflowTask.instance_id == str(context.instance.id),
|
|
WorkflowTask.node_id == task.node_id,
|
|
WorkflowTask.status == 'pending',
|
|
WorkflowTask.id != task.id,
|
|
)
|
|
result = await context.db.execute(stmt)
|
|
other_pending = result.scalars().all()
|
|
if other_pending:
|
|
logger.info(f"或签模式,取消其他 {len(other_pending)} 个待审批任务")
|
|
canceled_task_ids = []
|
|
for t in other_pending:
|
|
canceled_task_ids.append(str(t.id))
|
|
t.status = 'canceled'
|
|
t.comment = '__or_sign_canceled__'
|
|
context.db.add(t)
|
|
await context.db.flush()
|
|
|
|
for tid in canceled_task_ids:
|
|
try:
|
|
from core.message.service import NotifyService
|
|
await NotifyService.complete_dingtalk_todo(context.db, "workflow_task", tid)
|
|
except Exception as e:
|
|
logger.warning(f"或签取消-清理钉钉待办失败 task={tid}: {e}")
|
|
|
|
# 检查是否还有未完成的加签任务
|
|
# 1. waiting 状态的任务(前加签产生的原任务)
|
|
stmt = select(WorkflowTask).where(
|
|
WorkflowTask.instance_id == str(context.instance.id),
|
|
WorkflowTask.node_id == task.node_id,
|
|
WorkflowTask.status == 'waiting',
|
|
)
|
|
result = await context.db.execute(stmt)
|
|
waiting_tasks = result.scalars().all()
|
|
|
|
# 2. 后加签任务
|
|
stmt = select(WorkflowTask).where(
|
|
WorkflowTask.instance_id == str(context.instance.id),
|
|
WorkflowTask.node_id == task.node_id,
|
|
WorkflowTask.sign_type == 'after',
|
|
WorkflowTask.status == 'pending',
|
|
)
|
|
result = await context.db.execute(stmt)
|
|
after_sign_tasks = result.scalars().all()
|
|
|
|
if len(waiting_tasks) > 0:
|
|
logger.info(f"还有 {len(waiting_tasks)} 个任务在等待前加签完成")
|
|
if multi_approval == 'any':
|
|
canceled_ids = []
|
|
for wt in waiting_tasks:
|
|
canceled_ids.append(str(wt.id))
|
|
wt.status = 'canceled'
|
|
wt.comment = '__or_sign_canceled__'
|
|
context.db.add(wt)
|
|
await context.db.flush()
|
|
for tid in canceled_ids:
|
|
try:
|
|
from core.message.service import NotifyService
|
|
await NotifyService.complete_dingtalk_todo(context.db, "workflow_task", tid)
|
|
except Exception as e:
|
|
logger.warning(f"或签取消waiting任务-清理钉钉待办失败 task={tid}: {e}")
|
|
|
|
if len(after_sign_tasks) > 0:
|
|
logger.info(f"还有 {len(after_sign_tasks)} 个后加签任务未完成")
|
|
if multi_approval == 'any':
|
|
canceled_ids = []
|
|
for ast in after_sign_tasks:
|
|
canceled_ids.append(str(ast.id))
|
|
ast.status = 'canceled'
|
|
ast.comment = '__or_sign_canceled__'
|
|
context.db.add(ast)
|
|
await context.db.flush()
|
|
for tid in canceled_ids:
|
|
try:
|
|
from core.message.service import NotifyService
|
|
await NotifyService.complete_dingtalk_todo(context.db, "workflow_task", tid)
|
|
except Exception as e:
|
|
logger.warning(f"或签取消后加签任务-清理钉钉待办失败 task={tid}: {e}")
|
|
|
|
# 推进到下一节点
|
|
await self.advance_to_next(context, node)
|
|
|
|
async def _handle_sign_task_completion(self, context: 'ExecutionContext', task: Any) -> None:
|
|
"""
|
|
处理加签任务完成后的逻辑
|
|
"""
|
|
from online_dev.workflow.model import WorkflowTask
|
|
from online_dev.workflow.engine.utils import FlowUtils
|
|
|
|
sign_type = task.sign_type
|
|
parent_task_id = task.parent_task_id
|
|
|
|
logger.info(f"加签任务完成 - 类型: {sign_type}, 父任务ID: {parent_task_id}")
|
|
|
|
if sign_type == 'before':
|
|
# 前加签完成:检查是否所有前加签任务都完成了
|
|
stmt = select(WorkflowTask).where(
|
|
WorkflowTask.instance_id == str(context.instance.id),
|
|
WorkflowTask.node_id == task.node_id,
|
|
WorkflowTask.sign_type == 'before',
|
|
WorkflowTask.parent_task_id == parent_task_id,
|
|
WorkflowTask.status == 'pending',
|
|
)
|
|
result = await context.db.execute(stmt)
|
|
pending_before_signs = result.scalars().all()
|
|
|
|
if len(pending_before_signs) > 0:
|
|
logger.info(f"还有 {len(pending_before_signs)} 个前加签任务未完成")
|
|
return
|
|
|
|
# 所有前加签任务完成,恢复原任务
|
|
if parent_task_id:
|
|
stmt = select(WorkflowTask).where(WorkflowTask.id == parent_task_id)
|
|
result = await context.db.execute(stmt)
|
|
parent_task = result.scalar_one_or_none()
|
|
if parent_task and parent_task.status == 'waiting':
|
|
parent_task.status = 'pending'
|
|
context.db.add(parent_task)
|
|
await context.db.flush()
|
|
logger.info(f"前加签完成,恢复原任务: {parent_task_id}")
|
|
|
|
elif sign_type == 'after':
|
|
# 后加签完成:检查是否所有后加签任务都完成了
|
|
stmt = select(WorkflowTask).where(
|
|
WorkflowTask.instance_id == str(context.instance.id),
|
|
WorkflowTask.node_id == task.node_id,
|
|
WorkflowTask.sign_type == 'after',
|
|
WorkflowTask.parent_task_id == parent_task_id,
|
|
WorkflowTask.status == 'pending',
|
|
)
|
|
result = await context.db.execute(stmt)
|
|
pending_after_signs = result.scalars().all()
|
|
|
|
if len(pending_after_signs) > 0:
|
|
logger.info(f"还有 {len(pending_after_signs)} 个后加签任务未完成")
|
|
return
|
|
|
|
# 所有后加签任务完成,推进流程
|
|
node = FlowUtils.find_node_by_id(context.flow_definition, task.node_id)
|
|
if node:
|
|
await self.advance_to_next(context, node)
|
|
|
|
elif sign_type == 'parallel':
|
|
# 并行加签完成
|
|
node = FlowUtils.find_node_by_id(context.flow_definition, task.node_id)
|
|
if not node:
|
|
return
|
|
|
|
node_config = self.get_node_config(node)
|
|
multi_approval = node_config.get('multiApproval', 'any')
|
|
|
|
stmt = select(WorkflowTask).where(
|
|
WorkflowTask.instance_id == str(context.instance.id),
|
|
WorkflowTask.node_id == task.node_id,
|
|
WorkflowTask.status == 'pending',
|
|
)
|
|
result = await context.db.execute(stmt)
|
|
pending_tasks = result.scalars().all()
|
|
|
|
if multi_approval == 'any':
|
|
await self.advance_to_next(context, node)
|
|
elif multi_approval == 'parallel' and len(pending_tasks) == 0:
|
|
await self.advance_to_next(context, node)
|
|
|
|
elif sign_type == 'delegate':
|
|
# 委托任务完成
|
|
if parent_task_id:
|
|
stmt = select(WorkflowTask).where(WorkflowTask.id == parent_task_id)
|
|
result = await context.db.execute(stmt)
|
|
parent_task = result.scalar_one_or_none()
|
|
if parent_task and parent_task.status == 'delegated':
|
|
parent_task.status = 'pending'
|
|
parent_task.comment = f'委托人已审批通过,请确认'
|
|
context.db.add(parent_task)
|
|
await context.db.flush()
|
|
logger.info(f"委托任务完成,恢复原任务: {parent_task_id}")
|
|
|
|
elif sign_type == 'transfer':
|
|
# 转交任务完成
|
|
node = FlowUtils.find_node_by_id(context.flow_definition, task.node_id)
|
|
if node:
|
|
node_config = self.get_node_config(node)
|
|
multi_approval = node_config.get('multiApproval', 'any')
|
|
|
|
if multi_approval == 'any':
|
|
await self.advance_to_next(context, node)
|
|
elif multi_approval == 'parallel':
|
|
stmt = select(WorkflowTask).where(
|
|
WorkflowTask.instance_id == str(context.instance.id),
|
|
WorkflowTask.node_id == task.node_id,
|
|
WorkflowTask.status == 'pending',
|
|
)
|
|
result = await context.db.execute(stmt)
|
|
pending_tasks = result.scalars().all()
|
|
if len(pending_tasks) == 0:
|
|
await self.advance_to_next(context, node)
|
|
else:
|
|
await self.advance_to_next(context, node)
|
|
|
|
async def _get_next_sequential_assignee(
|
|
self,
|
|
context: 'ExecutionContext',
|
|
current_task: Any,
|
|
) -> Optional[str]:
|
|
"""
|
|
获取依次审批的下一个审批人
|
|
"""
|
|
from online_dev.workflow.model import WorkflowTask
|
|
from online_dev.workflow.engine.assignee_resolver import assignee_resolver
|
|
from online_dev.workflow.engine.utils import FlowUtils
|
|
|
|
node = FlowUtils.find_node_by_id(context.flow_definition, current_task.node_id)
|
|
if not node:
|
|
return None
|
|
|
|
node_config = self.get_node_config(node)
|
|
|
|
# 获取所有审批人列表
|
|
all_assignees = await assignee_resolver.resolve(
|
|
context.db,
|
|
node_config,
|
|
context.instance,
|
|
context.form_data,
|
|
)
|
|
|
|
if not all_assignees:
|
|
return None
|
|
|
|
# 获取已处理的任务
|
|
stmt = select(WorkflowTask.assignee_id).where(
|
|
WorkflowTask.instance_id == str(context.instance.id),
|
|
WorkflowTask.node_id == current_task.node_id,
|
|
WorkflowTask.status.in_(['approved', 'transferred']),
|
|
).order_by(WorkflowTask.sys_create_datetime)
|
|
result = await context.db.execute(stmt)
|
|
handled_ids = [str(uid) for uid in result.scalars().all()]
|
|
|
|
handled_set = set(handled_ids)
|
|
|
|
# 找到下一个未处理的审批人
|
|
for assignee_id in all_assignees:
|
|
if assignee_id not in handled_set:
|
|
return assignee_id
|
|
|
|
return None
|