diff --git a/app/controller/add_approve_route.py b/app/controller/add_approve_route.py index 468d7c0..637f07e 100644 --- a/app/controller/add_approve_route.py +++ b/app/controller/add_approve_route.py @@ -1,13 +1,12 @@ -from typing import TypedDict - from fastapi import FastAPI -from langgraph.constants import START, END -from langgraph.graph import StateGraph +from langchain_core.runnables import RunnableConfig +from langgraph.func import entrypoint -from app.model.ApproveModel import ApproveModel, ApproveService +from app.model.ApproveModel import ApproveService from app.model.ReimburseModel import ReimburseService, ReimburseModel from app.utils.db_utils import AsyncSessionDep from app.utils.postgres_checkpointer import AsyncPostgresSaverDep +from app.workflow.approve import create_approve_workflow, ApproveWorkflowInputs def add_approve_route(app: FastAPI): @@ -53,36 +52,24 @@ def add_approve_route(app: FastAPI): # 将审批单与报销单管理,设置报销单的approve_id为审批单的id await ReimburseService.item_update(session=session, row_dict={"id": reimburse_id, "approve_id": insert_approve_cls.id}, ) - graph = create_approve_graph(session, checkpointer) + reimburse_workflow = create_reimburse_workflow(session=session, checkpointer=checkpointer) - await graph.ainvoke({ - "input_approve_dict": insert_approve_cls.model_dump(), - "input_reimburse_user_id": reimburse_cls.user_id - }) + await reimburse_workflow.ainvoke({"reimburse_user_id": reimburse_cls.user, "approve_id": insert_approve_cls.id}) return {"message": "报销单提交成功!"} -# 创建一个图来处理审批流 -def create_approve_graph( - session: AsyncSessionDep, +def create_reimburse_workflow( checkpointer: AsyncPostgresSaverDep, + session: AsyncSessionDep, ): - class StateSchema(TypedDict): - input_reimburse_user_id: str - input_approve_dict: dict + approve_workflow = create_approve_workflow(checkpointer, session) - status: str + @entrypoint(checkpointer=checkpointer) + async def reimburse_workflow(input_dict: ApproveWorkflowInputs, config: RunnableConfig): + approve_flag = await approve_workflow.ainvoke(input_dict, config=config) + # 这里审批结束之后可以做一些事情,比如发送邮件、短信、微信消息等通知用户,因为没有实现对应的模块,这里就不做任何处理 + print("报销单审批流执行结束::::::::::>>>>>>>>>>>", approve_flag) + return None - builder = StateGraph(StateSchema) - - async def node_create_approve(state: StateSchema): - pass - - builder.add_node(node_create_approve) - builder.add_edge(START, 'node_create_approve') - builder.add_edge('node_create_approve', END) - - graph = builder.compile(checkpointer=checkpointer) - - return graph + return reimburse_workflow diff --git a/app/workflow/approve.py b/app/workflow/approve.py new file mode 100644 index 0000000..5ed5b99 --- /dev/null +++ b/app/workflow/approve.py @@ -0,0 +1,141 @@ +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({}) + + # 退出循环之后,如果审批结果是驳回,则结束审批流程 + 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) + + if not supervisor_user: + return False + + await update_approve( + approve_id=approve_id, + session=session, + log_content=f"「{supervisor_user}」处理审批", + 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.get('logs', '[]')) + approve_logs.append({"content": log_content, "datetime": datetime.datetime.now().strftime("%Y-%m-%d %H:%M:%S")}) + + # 更新审批单信息,包括审批人以及审批日志,此时审批状态为驳回 + approve_cls = 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 = UserService.query_item(session=session, row_dict={"id": user_id}) + # 目标用户的上级职位编码 + parent_code = user_cls.pos.parent_code + # 根据上级职位编码查询上级用户信息 + return 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