diff --git a/scripts/p3_feasibility_e2e.py b/scripts/p3_feasibility_e2e.py new file mode 100644 index 0000000..bf74723 --- /dev/null +++ b/scripts/p3_feasibility_e2e.py @@ -0,0 +1,304 @@ +# -*- 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 正链项目(默认 L2NJCn7XuOnzVV89lpXF1,owner=user-01) + P3_CLUSTER 正链类别(默认 AAbym8nFt1fLH1DayfDec 公众号 39条,须属于 P3_PROJECT_ID) + P3_CLUSTER_PERSONAL 全personal类别(默认 d6D5BAoEHgBn7KeiNbuE1,normalize 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 {} + check("A4 非法type进pending不炸", db_.get("总FP") == 7 and len(db_.get("待确认") or []) == 1, + bad[:150]) + 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", ""))) + # 重复立项行为记录(当前实现无幂等守卫,重复调用会再建项目——如实记录,清理兜底) + l2 = await call("opp_promote_to_project", {"report_id": RID}) + DUP = l2.startswith("OK:") and "project_id=" in l2 + if DUP: + DUPPID = l2.split("project_id=")[1].split("(")[0].strip() + print("-- FINDING: promote 重复调用再建项目 %s(无幂等守卫)--" % DUPPID, flush=True) + check("L3 重复立项行为已记录", l2.startswith("OK:") or l2.startswith("ERROR:"), l2[:120]) + + # ── 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())