pipeline-opportunity/pipeline_opportunity/opp_report_capability.py

319 lines
13 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""商机产线报告/审批能力:研发报告编写 → 人工确认 → 研发审批流程。
状态机(REPORT_TRANSITIONS,opp_common):
draft → confirmed → approval_initiated → approved / rejected
└ rejected → draft(可重来)
门禁纪律:
- confirm 只能 draft → confirmed,且是人工决策(由人工任务回流驱动)
- initiate_approval 只能 confirmed → approval_initiated(未经确认禁止发起审批)
- 审批结论由人工决策回流(approve/reject),agent 不得自批自审
"""
import logging
import os
import re
from ahserver.serverenv import ServerEnv
from .opp_common import (
get_db, new_id, rec_to_dict, rows_to_dicts, can_transition,
RP_DRAFT, RP_CONFIRMED, RP_APPROVAL_INITIATED, RP_APPROVED, RP_REJECTED,
AP_INITIATED, AP_APPROVED, AP_REJECTED,
HT_REPORT_CONFIRM, HT_DEV_APPROVAL,
find_human_task, create_human_task, get_report,
)
logger = logging.getLogger("pipeline.opportunity.report")
async def create_report(project_id, software, title="", analysis="",
created_by="agent.opportunity"):
"""创建研发报告草稿。返回 (True, report_id) 或 (False, 错误)。"""
if not software:
return False, "缺少 software(软件/主题名)"
if not title:
title = "%s 研发机会分析报告" % software
db, dbname = get_db()
rid = new_id()
async with db.sqlorContext(dbname) as sor:
await sor.C("opp_reports", {
"id": rid,
"project_id": project_id or "",
"software": software,
"title": title,
"content": analysis or "",
"status": RP_DRAFT,
"created_by": created_by,
})
await sor.sqlExe("COMMIT", {})
return True, rid
async def update_report(report_id, content):
"""更新草稿报告内容(仅 draft 可改)。"""
db, dbname = get_db()
async with db.sqlorContext(dbname) as sor:
rep = await get_report(sor, report_id)
if not rep:
return False, "报告不存在: %s" % report_id
if rep.get("status") != RP_DRAFT:
return False, "状态 %s 不可修改(仅草稿可编辑)" % rep.get("status")
await sor.U("opp_reports", {"id": report_id, "content": content})
await sor.sqlExe("COMMIT", {})
return True, "已更新"
async def list_reports(project_id="", status=""):
db, dbname = get_db()
async with db.sqlorContext(dbname) as sor:
sql = ("SELECT id, project_id, software, title, status, created_by, "
"created_at, updated_at FROM opp_reports WHERE 1=1")
args = {}
if project_id:
sql += " AND project_id=${pid}$"
args["pid"] = project_id
if status:
sql += " AND status=${st}$"
args["st"] = status
sql += " ORDER BY created_at DESC LIMIT 50"
recs = await sor.sqlExe(sql, args)
await sor.sqlExe("COMMIT", {})
return rows_to_dicts(recs)
async def get_report_full(report_id):
db, dbname = get_db()
async with db.sqlorContext(dbname) as sor:
rep = await get_report(sor, report_id)
return rep
async def _build_report_ppt(sor, rep):
"""为报告生成 PPT,存到项目工作空间 deliverables/ 目录。返回文件路径(失败 '')。"""
from .opp_ppt import build_report_ppt
# 解析项目目录({workspace}/{org}/{space}/projects/{项目}/)
project_id = rep.get("project_id", "") or ""
project_dir = ""
if project_id:
try:
recs = await sor.sqlExe(
"SELECT workspace_dir FROM sd_projects WHERE id=${p}$",
{"p": project_id})
if recs:
project_dir = getattr(recs[0], "workspace_dir", "") or ""
except Exception:
project_dir = ""
if not project_dir:
# 无项目/无目录:兜底放平台 files 目录
project_dir = os.path.join(os.getcwd(), "files", "opp_reports")
rid = rep.get("id", "")
software = re.sub(r"[^\w\u4e00-\u9fff-]", "_", rep.get("software", "") or "report")[:40]
fname = "research_report_%s_%s.pptx" % (software, rid[:8])
out_path = os.path.join(project_dir, "deliverables", fname)
ok, res = build_report_ppt(
rep.get("title", ""), rep.get("content", ""),
software=rep.get("software", ""), out_path=out_path)
if not ok:
logger.warning("build_report_ppt 失败: %s", str(res)[:200])
return ""
return out_path
async def submit_for_confirmation(report_id, note=""):
"""agent 完成报告后提交待人工确认:发人工确认任务(门禁)。
报告保持 draft,创建/复用 pending 人工任务。人工确认后才进 confirmed。
同时生成报告 PPT(存项目工作空间 deliverables/),供确认人下载审阅。
"""
db, dbname = get_db()
async with db.sqlorContext(dbname) as sor:
rep = await get_report(sor, report_id)
if not rep:
return False, "报告不存在: %s" % report_id
if rep.get("status") != RP_DRAFT:
return False, "状态 %s 不可提交确认" % rep.get("status")
# 生成/更新报告 PPT(幂等:重复提交覆盖旧文件)
ppt_path = ""
try:
ppt_path = await _build_report_ppt(sor, rep)
except Exception as e:
logger.warning("报告PPT生成失败(不阻塞提交): %s", str(e)[:200])
if ppt_path:
await sor.U("opp_reports", {"id": report_id, "ppt_path": ppt_path})
# 幂等:本报告已有 pending 确认任务则复用(按 confirm_task_id 绑定,
# 不复用其他报告的任务——确认是逐份报告的独立门禁,
# 且待办页按 confirm_task_id 反查报告正文,绑定丢失则报告无法被确认)
exist_id = rep.get("confirm_task_id") or ""
if exist_id:
t_recs = await sor.sqlExe(
"SELECT id, status FROM pipeline_human_tasks WHERE id=${i}$",
{"i": exist_id})
if t_recs and getattr(t_recs[0], "status", "") == "pending":
return True, "已有待确认任务: %s" % exist_id
hid = await create_human_task(
sor, rep.get("project_id", ""), HT_REPORT_CONFIRM,
"确认研发报告:%s" % rep.get("title", ""),
note or "请审阅研发报告《%s》,确认后发起研发审批。" % rep.get("title", ""))
await sor.U("opp_reports", {
"id": report_id, "confirm_task_id": hid})
await sor.sqlExe("COMMIT", {})
return True, "已提交人工确认任务: %s" % hid
async def confirm_report(report_id, ok, operator=""):
"""人工确认回流(门禁)。ok=True → confirmed;ok=False → 退回草稿待改。"""
db, dbname = get_db()
async with db.sqlorContext(dbname) as sor:
rep = await get_report(sor, report_id)
if not rep:
return False, "报告不存在: %s" % report_id
if rep.get("status") != RP_DRAFT:
return False, "状态 %s 无法确认(仅草稿可确认)" % rep.get("status")
new_status = RP_CONFIRMED if ok else RP_DRAFT
await sor.U("opp_reports", {
"id": report_id, "status": new_status,
"confirmed_by": operator or "human"})
# 完结对应的人工任务
hid = rep.get("confirm_task_id")
if hid:
await sor.U("pipeline_human_tasks", {
"id": hid, "status": "done" if ok else "rejected"})
await sor.sqlExe("COMMIT", {})
return True, "已确认" if ok else "已退回草稿"
async def initiate_approval(report_id, note="", created_by="agent.opportunity",
applicant_id=""):
"""发起研发审批:要求报告已 confirmed(人工确认门禁)。
审批通道优先级:
1. 钉钉审批——dda_approval_configs 配了 biz_type=opp_dev_approval 且启用时,
调 dingdingflow.submit_approval 发起,审批人节点链由配置的
approvers_config 决定;结果经回调写回(见 opp_common 注册的钩子)。
2. 降级人工——无钉钉配置/发起失败时,落平台内人工审批任务兜底,
保证产线闭环不因缺少钉钉环境而卡死。
"""
db, dbname = get_db()
async with db.sqlorContext(dbname) as sor:
rep = await get_report(sor, report_id)
if not rep:
return False, "报告不存在: %s" % report_id
if not await can_transition(rep.get("status"), RP_APPROVAL_INITIATED):
return False, "状态 %s 不可发起审批(须先人工确认)" % rep.get("status")
aid = new_id()
await sor.C("opp_approvals", {
"id": aid,
"report_id": report_id,
"project_id": rep.get("project_id", ""),
"status": AP_INITIATED,
"note": note,
"created_by": created_by,
})
await sor.U("opp_reports", {
"id": report_id, "status": RP_APPROVAL_INITIATED})
await sor.sqlExe("COMMIT", {})
# 尝试钉钉审批通道
dd_ok = False
env = ServerEnv()
submit_approval = getattr(env, "submit_approval", None)
if submit_approval is not None:
try:
org_id = rep.get("org_id", "") or "0"
r = await submit_approval(
"opp_dev_approval", aid,
"研发立项审批:%s" % rep.get("title", ""),
applicant_id or created_by, org_id)
if isinstance(r, dict) and r.get("success") and r.get("instance_id"):
dd_ok = True
logger.info("研发审批已推钉钉: report=%s instance=%s",
report_id, r.get("instance_id"))
else:
logger.warning("钉钉审批发起未成功(降级人工): report=%s msg=%s",
report_id, (r or {}).get("message", ""))
except Exception as e:
logger.warning("钉钉审批发起异常(降级人工): report=%s err=%s",
report_id, str(e)[:200])
if not dd_ok:
# 降级:平台内人工审批任务
db, dbname = get_db()
async with db.sqlorContext(dbname) as sor:
await create_human_task(
sor, rep.get("project_id", ""), HT_DEV_APPROVAL,
"研发审批:%s" % rep.get("title", ""),
note or "请审批《%s》的研发立项。" % rep.get("title", ""))
await sor.sqlExe("COMMIT", {})
return True, aid
return True, aid
async def resolve_approval(approval_id, result, operator=""):
"""审批结论回流(人工决策)。result=approved/rejected。"""
if result not in (AP_APPROVED, AP_REJECTED):
return False, "result 必须是 approved 或 rejected"
db, dbname = get_db()
async with db.sqlorContext(dbname) as sor:
recs = await sor.sqlExe(
"SELECT id, report_id, project_id, status FROM opp_approvals WHERE id=${i}$",
{"i": approval_id})
await sor.sqlExe("COMMIT", {})
if not recs:
return False, "审批单不存在: %s" % approval_id
rec = rec_to_dict(recs[0])
if rec.get("status") != AP_INITIATED:
return False, "审批单状态 %s 无法结论" % rec.get("status")
report_status = RP_APPROVED if result == AP_APPROVED else RP_REJECTED
await sor.U("opp_approvals", {
"id": approval_id, "status": result,
"resolved_by": operator or "human"})
await sor.U("opp_reports", {
"id": rec.get("report_id"), "status": report_status})
# 完结人工审批任务
ht = await find_human_task(sor, rec.get("project_id", ""),
HT_DEV_APPROVAL, status="pending")
if ht:
await sor.U("pipeline_human_tasks", {
"id": ht.get("id"),
"status": "done" if result == AP_APPROVED else "rejected"})
await sor.sqlExe("COMMIT", {})
return True, result
async def on_dingtalk_approval_done(biz_id, status, approval_id, comment):
"""钉钉审批结果回调钩子(dingdingflow register_biz_handler 签名)。
biz_id = opp_approvals.id(发起审批时传入)。
status ∈ approved / rejected / cancelled,cancelled 按驳回处理
(产线不能停在 approval_initiated 永久卡死)。
"""
if status == "approved":
result = AP_APPROVED
else:
result = AP_REJECTED
ok, msg = await resolve_approval(
biz_id, result, operator="dingtalk:%s" % (approval_id or ""))
if ok:
logger.info("钉钉审批结论已回流: opp_approval=%s -> %s", biz_id, result)
else:
logger.warning("钉钉审批结论回流失败: opp_approval=%s msg=%s", biz_id, msg)
async def list_approvals(project_id="", status=""):
db, dbname = get_db()
async with db.sqlorContext(dbname) as sor:
sql = ("SELECT id, report_id, project_id, status, note, created_by, "
"created_at, resolved_at FROM opp_approvals WHERE 1=1")
args = {}
if project_id:
sql += " AND project_id=${pid}$"
args["pid"] = project_id
if status:
sql += " AND status=${st}$"
args["st"] = status
sql += " ORDER BY created_at DESC LIMIT 50"
recs = await sor.sqlExe(sql, args)
await sor.sqlExe("COMMIT", {})
return rows_to_dicts(recs)