pipeline-opportunity/scripts/p3_feasibility_e2e.py

315 lines
18 KiB
Python
Raw Permalink 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.

# -*- coding:utf-8 -*-
"""P3 可行性研究 E2E测试机真机
fp_estimate 规则引擎独立复算 → draft_feasibility 正/负链 → score_candidate →
promote 门禁draft/confirmed 拒绝→审批 approved 放行→立项落库核对)→
八节 PPT 生成 → 项目/org 隔离。钉钉审批通道 monkeypatch 为失败→走降级人工(不真发)。
用法(测试机 /d/pipeline/pipeline-app 下):
setsid nohup ./py3/bin/python pkgs/pipeline-opportunity/scripts/p3_feasibility_e2e.py \
> /tmp/p3_e2e.log 2>&1 &
环境变量:
P3_PROJECT_ID 正链项目(默认 L2NJCn7XuOnzVV89lpXF1owner=user-01
P3_CLUSTER 正链类别(默认 AAbym8nFt1fLH1DayfDec 公众号 39条须属于 P3_PROJECT_ID
P3_CLUSTER_PERSONAL 全personal类别默认 d6D5BAoEHgBn7KeiNbuE1normalize done 无 core/ext
P3_CLUSTER_OTHER 跨项目隔离用类别(默认 3kPqI-COvrU66QQhbEDp_批次平台级
P3_KEEP =1 保留 E2E 建的可行性报告与项目(默认报告保留、立项项目删除)
"""
import asyncio
import importlib.util
import json
import math
import os
import sys
import time
W = "/d/pipeline/pipeline-app"
os.chdir(W)
sys.path.insert(0, W)
from appPublic.folderUtils import ProgramPath
from appPublic.jsonConfig import getConfig
from appPublic.event_dispatcher import EventDispatcher
from sqlor.dbpools import DBPools
from ahserver.serverenv import ServerEnv
from ahserver.globalEnv import initEnv
p = ProgramPath()
c = getConfig(W, {"workdir": W, "ProgramPath": p})
DBPools(c.databases)
se = ServerEnv()
se.event_dispatcher = EventDispatcher()
se.get_module_dbname = lambda m: "pipeline"
initEnv()
PID = os.environ.get("P3_PROJECT_ID", "L2NJCn7XuOnzVV89lpXF1")
CL_OK = os.environ.get("P3_CLUSTER", "AAbym8nFt1fLH1DayfDec") # 公众号 39条 init
CL_PERSONAL = os.environ.get("P3_CLUSTER_PERSONAL", "d6D5BAoEHgBn7KeiNbuE1") # done 全personal
CL_PLATFORM = os.environ.get("P3_CLUSTER_OTHER", "3kPqI-COvrU66QQhbEDp_") # 平台级批次类别
KEEP = os.environ.get("P3_KEEP", "") == "1"
FAIL = []
def check(n, cond, d=""):
print("[%s] %s %s" % ("OK " if cond else "FAIL", n, str(d)[:220]), flush=True)
if not cond:
FAIL.append(n)
# fp_calc 规则引擎skills/global 部署位),用于独立复算
FP_CALC = "/d/pipeline/pipeline-app/skills/global/function-point-counting/scripts/fp_calc.py"
_spec = importlib.util.spec_from_file_location("fp_calc_ref", FP_CALC)
fp_ref = importlib.util.module_from_spec(_spec)
_spec.loader.exec_module(fp_ref)
CTX = {"project_id": PID, "user_id": "user-01", "org_id": "0",
"pipeline_id": "opportunity_general", "space": "opportunity_general",
"session_id": "e2e-p3", "workspace_dir": "/tmp"}
async def sql(q, args=None, one=False):
async with DBPools().sqlorContext("pipeline") as sor:
r = await sor.sqlExe(q, args or {})
await sor.sqlExe("COMMIT", {})
return r[0] if (one and r) else r
async def m():
from pipeline_opportunity import opp_ability
ab = opp_ability.register_opp_ability()
async def call(t, prm, ctx=None):
async with DBPools().sqlorContext("pipeline") as sor:
return await ab.handlers[t](sor, prm, ctx or CTX)
print("== P3 E2E start %s project=%s ==" % (time.strftime("%F %T"), PID), flush=True)
# ── A. fp_estimate 规则引擎(独立复算,不经 LLM──
funcs = [
{"name": "用户档案", "type": "ILF", "det": 12, "ret": 2, "ftr": 0, "confidence": "high"},
{"name": "订单录入", "type": "EI", "det": 8, "ret": 0, "ftr": 1, "confidence": "high"},
{"name": "报表导出", "type": "EO", "det": 9, "ret": 0, "ftr": 2, "confidence": "medium"},
{"name": "余额查询", "type": "EQ", "det": 4, "ret": 0, "ftr": 1, "confidence": "medium"},
{"name": "外部支付接口", "type": "EIF", "det": 6, "ret": 1, "ftr": 0, "confidence": "low"},
]
out = await call("opp_fp_estimate", {"functions_json": json.dumps(funcs),
"project_name": "E2E样例"})
d = json.loads(out) if out.startswith("{") else {}
ref = fp_ref.calc({"project": "E2E样例", "functions": [
dict(f, ret=f.get("ret") or (1 if f["type"] in ("ILF", "EIF") else 0)) for f in funcs]})
check("A1 fp_total=独立复算", d.get("总FP") == ref["total_fp"],
"tool=%s ref=%s" % (d.get("总FP"), ref["total_fp"]))
check("A2 分类型计数一致", d.get("分类型") == ref["counts"], "%s vs %s" % (d.get("分类型"), ref["counts"]))
check("A3 四件套markdown非空", bool((d.get("四件套") or "").strip()), len(d.get("四件套") or ""))
bad = await call("opp_fp_estimate", {"functions_json": json.dumps(
[{"name": "x", "type": "XXX", "det": 5}, {"name": "y", "type": "EI", "det": 6, "ftr": 1}])})
db_ = json.loads(bad) if bad.startswith("{") else {}
# EI det=6(档5-15) ftr=1(档0-1) → 低复杂度 → 3 FP独立复算核对
ref_bad = fp_ref.calc({"project": "", "functions": [
{"name": "y", "type": "EI", "det": 6, "ftr": 1}]})
check("A4 非法type进pending不炸", db_.get("总FP") == ref_bad["total_fp"] == 3
and len(db_.get("待确认") or []) == 1,
"fp=%s ref=%s pending=%d" % (db_.get("总FP"), ref_bad["total_fp"],
len(db_.get("待确认") or [])))
e1 = await call("opp_fp_estimate", {"functions_json": "[]"})
check("A5 空数组拒绝", e1.startswith("ERROR:"), e1[:80])
# ── B. draft_feasibility 负链 ──
b1 = await call("opp_draft_feasibility", {"cluster_id": CL_PERSONAL}, CTX)
check("B1 全personal类别拒绝(样本不足)", b1.startswith("ERROR:") and "样本不足" in b1, b1[:160])
ctxB = dict(CTX); ctxB["project_id"] = "FAKE_PROJECT_X"
b2 = await call("opp_draft_feasibility", {"cluster_id": CL_OK}, ctxB)
check("B2 跨项目类别拒绝", b2.startswith("ERROR:") and "其他项目" in b2, b2[:160])
ctxC = dict(CTX); ctxC["org_id"] = "org-evil"
b3 = await call("opp_draft_feasibility", {"cluster_id": CL_PLATFORM}, ctxC)
check("B3 伪org拒绝", b3.startswith("ERROR:") and "机构" in b3, b3[:160])
b4 = await call("opp_draft_feasibility", {"cluster_id": "NO_SUCH_CLUSTER"})
check("B4 类别不存在拒绝", b4.startswith("ERROR:") and "不存在" in b4, b4[:120])
# 未归一化的其它 init 类别(限当前会话可见:平台级批次或本项目批次)
r = await sql(
"SELECT c.id FROM opp_clusters c LEFT JOIN opp_mining_batches b ON b.id=c.batch_id "
"WHERE c.normalize_status='init' AND c.id<>${a}$ AND (b.project_id='' OR b.project_id=${p}$) "
"LIMIT 1", {"a": CL_OK, "p": PID}, one=True)
if r:
b5 = await call("opp_draft_feasibility", {"cluster_id": r.id})
check("B5 未归一化拒绝(引导先normalize)", b5.startswith("ERROR:") and "共性提取" in b5, b5[:150])
# ── C. normalize 正链类别(真实 LLM+VDB 路径E2E 前置)──
st = await sql("SELECT normalize_status FROM opp_clusters WHERE id=${c}$", {"c": CL_OK}, one=True)
if getattr(st, "normalize_status", "") != "done":
print("-- normalize %s (may take minutes) --" % CL_OK, flush=True)
t0 = time.time()
nout = await call("opp_normalize_cluster", {"cluster_id": CL_OK})
check("C1 normalize完成", nout.startswith("OK:"), "%s (%.0fs)" % (nout[:150], time.time() - t0))
else:
check("C1 normalize完成(已done)", True, "skip")
cov = await call("opp_coverage_report", {"cluster_id": CL_OK})
covd = json.loads(cov) if cov.startswith("{") else {}
tiers = covd.get("层级统计") or {}
n_core, n_ext = tiers.get("核心(≥60%)", 0), tiers.get("扩展(≥40%)", 0)
check("C2 有core/ext共性(正链可走)", (n_core + n_ext) > 0,
"core=%s ext=%s" % (n_core, n_ext))
# ── D. draft_feasibility 正链 ──
t0 = time.time()
dout = await call("opp_draft_feasibility", {"cluster_id": CL_OK})
check("D1 draft_feasibility成功", dout.startswith("OK:") and "report_id=" in dout,
"%s (%.0fs)" % (dout[:150], time.time() - t0))
RID = ""
if "report_id=" in dout:
RID = dout.split("report_id=")[1].split("")[0].split("")[0].strip()
rep = await sql("SELECT * FROM opp_reports WHERE id=${r}$", {"r": RID}, one=True) if RID else None
check("D2 落库report_type=feasibility", rep is not None and rep.report_type == "feasibility",
getattr(rep, "report_type", None))
content = getattr(rep, "content", "") or ""
SECS = ("## 市场需求", "## 共性需求", "## 需求规格", "## 架构与可行性",
"## FP与成本", "## 商业价值", "## 风险", "## 结论建议")
miss = [s for s in SECS if s not in content]
check("D3 八节齐全", not miss, "缺:%s" % miss if miss else "8/8")
check("D4 挂项目%s" % PID[:8], getattr(rep, "project_id", "") == PID, rep.project_id if rep else "")
check("D5 status=draft", getattr(rep, "status", "") == "draft", getattr(rep, "status", ""))
# FP/成本与 params 系数一致性
import re as _re
mm = _re.search(r"功能点合计 (\d+(?:\.\d+)?) FP.*?按 ([\d.]+) 人月/FP 估 ([\d.]+) 人月;"
r"按 ([\d.]+) 万元/人月估研发成本 ([\d.]+) 万元", content)
if mm:
fp_v, pm_per, pm_v, wan_per, cost_v = (float(x) for x in mm.groups())
check("D6 成本=FP×系数(params一致)",
abs(pm_v - round(fp_v * pm_per, 1)) < 0.05 and abs(cost_v - round(pm_v * wan_per, 1)) < 0.05
and pm_per == 0.05 and wan_per == 3.0,
"FP=%s pm=%s cost=%s" % (fp_v, pm_v, cost_v))
check("D7 FP>0", fp_v > 0, fp_v)
else:
check("D6 成本=FP×系数", False, "FP与成本节格式不符: %s" % content[content.find("## FP"):content.find("## FP") + 200])
check("D7 FP>0", False, "")
check("D8 架构节非LLM失败占位", "架构分析生成失败" not in content, "")
# ── E. score_candidate ──
sout = await call("opp_score_candidate", {"cluster_id": CL_OK})
sd = json.loads(sout) if sout.startswith("{") else {}
sc = sd.get("score")
check("E1 评分∈[0,100]", isinstance(sc, (int, float)) and 0 <= sc <= 100, sout[:150])
check("E2 三因子齐全", len(sd.get("factors") or {}) == 3, list((sd.get("factors") or {}).keys()))
check("E3 优先级星", bool(sd.get("priority", "").startswith("")), sd.get("priority"))
# ── F. promote 门禁负链 ──
f1 = await call("opp_promote_to_project", {"report_id": RID})
check("F1 draft立项拒绝(未过确认)", f1.startswith("ERROR:") and "确认" in f1, f1[:150])
# research 报告立项拒绝report_type 门禁在 status 门禁之前)
rr = await sql("SELECT id FROM opp_reports WHERE report_type='research' AND id<>${r}$ LIMIT 1",
{"r": RID}, one=True)
if rr:
f2 = await call("opp_promote_to_project", {"report_id": rr.id})
check("F2 research报告立项拒绝", f2.startswith("ERROR:") and "仅可行性" in f2, f2[:120])
# ── G. 提交人工确认 + PPT ──
from pipeline_opportunity.opp_report_capability import submit_for_confirmation, confirm_report, \
initiate_approval, resolve_approval
g1ok, g1msg = await submit_for_confirmation(RID, note="P3 E2E")
check("G1 submit_for_confirmation", g1ok, g1msg[:120])
rep2 = await sql("SELECT confirm_task_id,ppt_path,status FROM opp_reports WHERE id=${r}$",
{"r": RID}, one=True)
ht = getattr(rep2, "confirm_task_id", "") or ""
check("G2 确认任务绑定", bool(ht), ht)
htr = await sql("SELECT task_type,status FROM pipeline_human_tasks WHERE id=${i}$",
{"i": ht}, one=True) if ht else None
check("G3 人工任务pending", htr is not None and htr.status == "pending"
and htr.task_type == "opp_report_confirm",
"%s/%s" % (getattr(htr, "task_type", ""), getattr(htr, "status", "")))
ppt = getattr(rep2, "ppt_path", "") or ""
ok_ppt = ppt and os.path.isfile(ppt) and os.path.getsize(ppt) > 5000
check("G4 PPT已生成", bool(ok_ppt), "%s %sB" % (ppt, os.path.getsize(ppt) if ppt and os.path.isfile(ppt) else -1))
if ok_ppt:
from pptx import Presentation
prs = Presentation(ppt)
all_txt = []
for sl in prs.slides:
for sh in sl.shapes:
if sh.has_text_frame:
all_txt.append(sh.text_frame.text)
joined = "\n".join(all_txt)
sec_hits = [s for s in ("市场需求", "共性需求", "需求规格", "架构与可行性",
"FP与成本", "商业价值", "风险", "结论建议") if s in joined]
check("G5 PPT八节章节页", len(sec_hits) >= 8, "%d/8 %s" % (len(sec_hits), sec_hits))
n_slides = len(prs.slides._sldIdLst)
check("G6 PPT页数>10", n_slides > 10, n_slides)
check("G7 PPT文件名feasibility前缀", os.path.basename(ppt).startswith("feasibility_report_"),
os.path.basename(ppt))
# ── H. 人工确认回流 ──
h1ok, h1msg = await confirm_report(RID, True, operator="user-01")
stt = await sql("SELECT status FROM opp_reports WHERE id=${r}$", {"r": RID}, one=True)
check("H1 confirmed", h1ok and stt.status == "confirmed", stt.status)
ht2 = await sql("SELECT status FROM pipeline_human_tasks WHERE id=${i}$", {"i": ht}, one=True)
check("H2 确认任务done", ht2 and ht2.status == "done", getattr(ht2, "status", ""))
# ── I. confirmed 但未审批 → promote 拒绝 ──
i1 = await call("opp_promote_to_project", {"report_id": RID})
check("I1 confirmed未审批拒绝", i1.startswith("ERROR:") and "审批" in i1, i1[:150])
# ── J. 发起审批(钉钉 monkeypatch 失败 → 降级人工兜底,不真发)──
async def _fake_submit(biz, aid, title, applicant, org):
return {"success": False, "message": "E2E monkeypatch: 跳过钉钉"}
se.submit_approval = _fake_submit
jok, aid = await initiate_approval(RID, note="P3 E2E 审批", created_by="user-01",
applicant_id="user-01")
check("J1 审批单发起", jok, str(aid)[:80])
stj = await sql("SELECT status FROM opp_reports WHERE id=${r}$", {"r": RID}, one=True)
apr = await sql("SELECT status,report_id FROM opp_approvals WHERE id=${a}$", {"a": aid}, one=True)
check("J2 报告approval_initiated+审批单initiated",
stj.status == "approval_initiated" and apr and apr.status == "initiated",
"%s/%s" % (stj.status, getattr(apr, "status", "")))
dht = await sql("SELECT id,status FROM pipeline_human_tasks WHERE project_id=${p}$ AND task_type='opp_dev_approval' AND status='pending' LIMIT 1",
{"p": PID}, one=True)
check("J3 降级人工审批任务兜底", dht is not None, getattr(dht, "id", ""))
# ── K. 审批回流 approved ──
kok, kmsg = await resolve_approval(aid, "approved", operator="user-01")
stk = await sql("SELECT status FROM opp_reports WHERE id=${r}$", {"r": RID}, one=True)
apk = await sql("SELECT status FROM opp_approvals WHERE id=${a}$", {"a": aid}, one=True)
check("K1 报告approved+审批单approved", kok and stk.status == "approved" and apk.status == "approved",
"%s/%s" % (stk.status, apk.status))
# ── L. 立项 ──
lout = await call("opp_promote_to_project", {"report_id": RID})
check("L1 promote成功", lout.startswith("OK:") and "project_id=" in lout, lout[:160])
NEWPID = lout.split("project_id=")[1].split("")[0].strip() if "project_id=" in lout else ""
if NEWPID:
np = await sql("SELECT name,project_type,pipeline_id,org_id,created_by,status FROM sd_projects WHERE id=${i}$",
{"i": NEWPID}, one=True)
check("L2 项目落库(demand_mining/opportunity_general)",
np is not None and np.project_type == "demand_mining"
and np.pipeline_id == "opportunity_general" and np.org_id == "0",
"%s/%s/%s" % (getattr(np, "project_type", ""), getattr(np, "pipeline_id", ""), getattr(np, "org_id", "")))
# 幂等守卫:重复 promote 必须复用既有项目,不再新建
l2 = await call("opp_promote_to_project", {"report_id": RID})
DUPPID = ""
if l2.startswith("OK:") and "project_id=" in l2:
cand = l2.split("project_id=")[1].split("")[0].split(" ")[0].strip()
if cand and cand != NEWPID:
DUPPID = cand
n_prj = await sql("SELECT COUNT(*) n FROM sd_projects WHERE description LIKE ${p}$",
{"p": "由需求挖掘可行性报告立项%%[report_id=%s]%%" % RID}, one=True)
check("L3 重复立项幂等(不重复建)", l2.startswith("OK:") and "幂等" in l2
and not DUPPID and getattr(n_prj, "n", 99) == 1,
"%s | dup=%s | n=%s" % (l2[:120], DUPPID or "-", getattr(n_prj, "n", "?")))
# ── M. 清理(报告+PPT保留供浏览器验收立项项目删除──
PROMOTED = [x for x in (NEWPID, locals().get("DUPPID", "")) if x]
if not KEEP:
for ppid in PROMOTED:
wdir = await sql("SELECT workspace_dir FROM sd_projects WHERE id=${i}$",
{"i": ppid}, one=True)
await sql("DELETE FROM sd_projects WHERE id=${i}$ AND project_type='demand_mining'",
{"i": ppid})
wd = getattr(wdir, "workspace_dir", "") or ""
if wd and os.path.isdir(wd):
import shutil
shutil.rmtree(wd, ignore_errors=True)
print("-- cleaned promoted project %s dir=%s" % (ppid, wd), flush=True)
print("\n== SUMMARY: %d FAIL %s ==" % (len(FAIL), FAIL if FAIL else "(ALL PASS)"), flush=True)
print("RID=%s PPT=%s" % (RID, ppt), flush=True)
asyncio.run(m())