import datetime import json from typing import List from typing import TypedDict, Union from langgraph.func import entrypoint, task from langgraph.types import interrupt from app.model.ApproveModel import ApproveService from app.model.UserModel import UserService, UserServiceModel from app.utils.db_utils import AsyncSessionDep from app.utils.postgres_checkpointer import AsyncPostgresSaverDep # 创建审批工作流 def create_approve_workflow( checkpointer: AsyncPostgresSaverDep, session: AsyncSessionDep, ): @entrypoint(checkpointer=checkpointer) async def approve_workflow(input_dict: ApproveWorkflowInputs) -> bool: # 申请人id reimburse_user_id = input_dict.get('reimburse_user_id') # 审批单id approve_id = input_dict.get('approve_id') # 立即拟造一条审批结果,审批人为申请人自己,审批结果为通过 approve_result: ApproveResult = {"flag": True, "user_id": reimburse_user_id, } # 循环判断审批结果,如果为通过,则继续触发下一个上级主管审批 while approve_result.get('flag'): # 审批通过 # 继续触发下一个上级主管审批 has_supervisor_approve = await task_supervisor_approve( user_id=approve_result.get('user_id'), approve_id=approve_id ) if has_supervisor_approve: # 触发中断,等待上级审批恢复中断,拿到审批结果 approve_result = interrupt({}) # 中断恢复之后进入下一个循环 else: # 已经没有上级主管了,退出循环 break # 退出循环之后,如果审批结果是驳回,则结束审批流程 if not approve_result.get('flag'): # 审批驳回 await update_approve( approve_id=approve_id, session=session, log_content=approve_result.get('reason'), # 不再需要审批人 # 审批状态为驳回 approve_dict={"user_id": "", "status": "rejected"}, ) return False else: # 审批通过 await update_approve( approve_id=approve_id, session=session, log_content=approve_result.get('reason'), # 不再需要审批人 # 审批状态为通过 approve_dict={"user_id": "", "status": "approved"}, ) return True # 执行上级主管审批,返回结果为布尔值,意思是是否触发了主管审批 @task async def task_supervisor_approve(user_id: str, approve_id: str) -> bool: # 直接找审批人的上级来审批 supervisor_user = await aget_supervisor_user(user_id=user_id, session=session) print('\n\n', ":::::::::::task_supervisor_approve start::::::::::::", '\n\n') print(supervisor_user.full_name, user_id, approve_id) print('\n\n', ":::::::::::task_supervisor_approve end::::::::::::", '\n\n') if not supervisor_user: return False await update_approve( approve_id=approve_id, session=session, log_content=f"「{supervisor_user.full_name}」处理审批", approve_dict={ "user_id": supervisor_user.id, # 审批人为上级主管 "status": "approving", # 审批状态为处理中 }, ) return True return approve_workflow async def update_approve( approve_id: str, session: AsyncSessionDep, log_content: str, approve_dict: dict, ): # 查询审批单信息,准备更新审批单中的日志信息 approve_cls = await ApproveService.query_item(session=session, row_dict={"id": approve_id}) # 更新审批日志 approve_logs: List[ApproveLog] = json.loads(approve_cls.logs or '[]') approve_logs.append({"content": log_content, "datetime": datetime.datetime.now().strftime("%Y-%m-%d %H:%M:%S")}) # 更新审批单信息,包括审批人以及审批日志,此时审批状态为驳回 approve_cls = await ApproveService.item_update(session=session, row_dict={ "id": approve_id, "logs": json.dumps(approve_logs, ensure_ascii=False), **approve_dict, }) return approve_cls # 找到上级主管用户信息 async def aget_supervisor_user(user_id: str, session: AsyncSessionDep) -> Union[UserServiceModel, None]: # 目标用户信息 user_cls: UserServiceModel = await UserService.query_item(session=session, row_dict={"id": user_id}) # 目标用户的上级职位编码 parent_code = user_cls.pos.parent_code # 根据上级职位编码查询上级用户信息 return await UserService.query_item(session=session, row_dict={"pos_code": parent_code}) # approve workflow审批流程的输入参数类型 class ApproveWorkflowInputs(TypedDict): reimburse_user_id: str approve_id: str # 审批回执数据类型 class ApproveResult(TypedDict): # 审批标识,是审批通过还是审批驳回 flag: bool # 审批驳回原因 reason: str # 审批人id user_id: str # 审批日志数据类型 class ApproveLog(TypedDict): content: str datetime: str