feat: 审批流基本结构
This commit is contained in:
@@ -1,13 +1,12 @@
|
|||||||
from typing import TypedDict
|
|
||||||
|
|
||||||
from fastapi import FastAPI
|
from fastapi import FastAPI
|
||||||
from langgraph.constants import START, END
|
from langchain_core.runnables import RunnableConfig
|
||||||
from langgraph.graph import StateGraph
|
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.model.ReimburseModel import ReimburseService, ReimburseModel
|
||||||
from app.utils.db_utils import AsyncSessionDep
|
from app.utils.db_utils import AsyncSessionDep
|
||||||
from app.utils.postgres_checkpointer import AsyncPostgresSaverDep
|
from app.utils.postgres_checkpointer import AsyncPostgresSaverDep
|
||||||
|
from app.workflow.approve import create_approve_workflow, ApproveWorkflowInputs
|
||||||
|
|
||||||
|
|
||||||
def add_approve_route(app: FastAPI):
|
def add_approve_route(app: FastAPI):
|
||||||
@@ -53,36 +52,24 @@ def add_approve_route(app: FastAPI):
|
|||||||
# 将审批单与报销单管理,设置报销单的approve_id为审批单的id
|
# 将审批单与报销单管理,设置报销单的approve_id为审批单的id
|
||||||
await ReimburseService.item_update(session=session, row_dict={"id": reimburse_id, "approve_id": insert_approve_cls.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({
|
await reimburse_workflow.ainvoke({"reimburse_user_id": reimburse_cls.user, "approve_id": insert_approve_cls.id})
|
||||||
"input_approve_dict": insert_approve_cls.model_dump(),
|
|
||||||
"input_reimburse_user_id": reimburse_cls.user_id
|
|
||||||
})
|
|
||||||
|
|
||||||
return {"message": "报销单提交成功!"}
|
return {"message": "报销单提交成功!"}
|
||||||
|
|
||||||
|
|
||||||
# 创建一个图来处理审批流
|
def create_reimburse_workflow(
|
||||||
def create_approve_graph(
|
|
||||||
session: AsyncSessionDep,
|
|
||||||
checkpointer: AsyncPostgresSaverDep,
|
checkpointer: AsyncPostgresSaverDep,
|
||||||
|
session: AsyncSessionDep,
|
||||||
):
|
):
|
||||||
class StateSchema(TypedDict):
|
approve_workflow = create_approve_workflow(checkpointer, session)
|
||||||
input_reimburse_user_id: str
|
|
||||||
input_approve_dict: dict
|
|
||||||
|
|
||||||
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)
|
return reimburse_workflow
|
||||||
|
|
||||||
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
|
|
||||||
|
|||||||
@@ -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
|
||||||
Reference in New Issue
Block a user