From 0547c373023e9bd6338b20be3a7396e375ca5f9d Mon Sep 17 00:00:00 2001 From: ymq Date: Mon, 31 Aug 2026 16:02:03 +0800 Subject: [PATCH] =?UTF-8?q?QC=E5=A5=91=E5=90=88=E5=BA=A6=E9=97=A8=E7=A6=81?= =?UTF-8?q?+=E5=A4=8D=E7=9B=98=E6=9C=BA=E5=88=B6=EF=BC=9A=E5=88=A0?= =?UTF-8?q?=E9=9C=80=E6=B1=82/=E8=AE=BE=E8=AE=A1=E4=BA=BA=E5=B7=A5?= =?UTF-8?q?=E7=A1=AE=E8=AE=A4=E3=80=81PM=E4=B8=89=E9=98=B6=E6=AE=B5?= =?UTF-8?q?=E7=A6=81=E5=86=85=E5=AE=B9=E5=90=A6=E5=86=B3(=E4=BB=A3?= =?UTF-8?q?=E7=A0=81=E5=85=9C=E5=BA=95=E4=BF=AE=E6=AD=A3)=E3=80=81?= =?UTF-8?q?=E5=A4=8D=E7=9B=98=E4=BB=BB=E5=8A=A1+=E6=89=A7=E8=A1=8C?= =?UTF-8?q?=E5=99=A8+=E6=8A=80=E8=83=BD=E6=8F=90=E8=AE=AE=E5=B7=A5?= =?UTF-8?q?=E5=85=B7=E3=80=81QC=E5=B7=A5=E5=85=B7=E9=93=BE=E8=A1=A5ctx?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pipeline_service/agent_loop.py | 316 +++++++++++++++---- pipeline_service/capability_tools.py | 14 + pipeline_service/retrospective_capability.py | 170 ++++++++++ 3 files changed, 431 insertions(+), 69 deletions(-) create mode 100644 pipeline_service/retrospective_capability.py diff --git a/pipeline_service/agent_loop.py b/pipeline_service/agent_loop.py index d6cb05e..90a6af3 100644 --- a/pipeline_service/agent_loop.py +++ b/pipeline_service/agent_loop.py @@ -605,6 +605,12 @@ PM_SYSTEM_PROMPT = """你是项目经理(PM)。你的职责是项目计划 2. 任务分配:用 create_tasks 工具把子任务派发给对应角色 agent,并记录父子关系与依赖关系(能并行的并行、不能并行的串行)。 3. 任务验收:检查交付件 + 实际产出文件,决定 review_approve / review_reject / review_complete。 +## QC 门禁(质量权威,2026-08-31 起) +- 需求/设计/开发(__ROLE__ 为 requirement/design/develop 时)的交付件到达你之前,已经过 QC 契合度门禁:逐项核对、契合度>9.5 才放行,退回意见精确到核对项。 +- **这三个阶段你不再复审内容质量**,禁止输出 review_reject 做内容否决(质量判定权在 QC);你的职责收敛为编排:确认交付件存在 → 需要拆分的先 create_tasks 派发(design 按模块清单派 develop 任务)→ review_approve 推进。 +- 例外(逃逸阀):发现系统性缺陷(如设计根本性错误影响下游全部成果)可用 review_rollback 回退,rollback 会在后续阶段暴露问题时使用,不是常规手段。 +- deploy_test/test/deploy_prod 阶段仍由你做完整验收(下述审核标准照旧)。 + ## 任务分解与编排(先评估,再拆解) - 派发前先评估:当前任务是否「过于复杂」(涉及多个模块/应用、多个独立交付单元、工作量超单 agent 一次产出)。简单任务直接派发单个任务,不必强行拆分。 - **按设计师划分好的模块派发**:develop 任务以 design 阶段设计师在 modules/<模块名>.md 里划分好的模块为单元——一个模块 = 一个 develop 任务,PM 不要自己重新拆模块。派发时按设计师定义的模块间依赖关系编排:无依赖的模块并行(不填 depends_on),有依赖的模块串行(depends_on 指向其依赖的模块任务)。基础模块(apppublic、sqlor、ahserver、accounting、appbase、rbac 等)已存在、可直接引用,**不派发「重新开发基础模块」的任务**。 @@ -664,9 +670,9 @@ __ROLE_SKILLS__ 回退:{"action":"review_rollback","rollback_role":"develop","comment":"回退原因"} 审核标准: -- requirement:需求是否清晰完整可量化(含部署环境需求是否明确) -- design:方案合理、覆盖需求、技术可行 -- develop:必须用 list_files/read_file 检查代码文件是否实际产出(git 提交由审核通过后系统统一执行) +- requirement:内容质量由 QC 门禁把关(契合度>9.5 才到达你),你只确认交付件存在 + 需要拆分的先 create_tasks 派发 → review_approve +- design:内容质量由 QC 门禁把关,你按 modules/*.md 的模块清单用 create_tasks 派 develop 任务 → review_approve +- develop:QC 把关内容;你只确认代码文件实际产出(list_files 存在性检查)→ review_approve - deploy_test:测试环境部署配置完整、服务可访问 - test:测试覆盖充分、发现问题记录完整 - deploy_prod:生产部署配置完整、可一键部署、有回滚方案 @@ -680,6 +686,13 @@ __ROLE_SKILLS__ QC_SYSTEM_PROMPT = """你是质量控制工程师(QC)。对交付件做合规检查和质量检查,不合规直接退回重做。 +## 评分协议(强制) +1. 先按对应审查技能(load_skill 加载)构建**核对项清单**——每项必须机械可判定(文件在不在/格式对不对/需求点有没有被响应),不做主观打分。 +2. 逐项判定:每项只输出 过/不过 + 证据(文件路径/位置/缺失说明)。 +3. 算分:**契合度 = 10 × 通过项数 / 总项数**,保留两位小数。 +4. 决策:契合度 > 9.5 → review_approve;否则 review_reject(退回意见=未过项清单,编号+位置+问题+改法)。 +5. 审查技能标注的硬门禁项(需求覆盖、应用脚手架禁项等)任一不过 → 无论总分多少一律 review_reject。 + ## 检查维度 1. 项目规范检查:产出是否按 SDLC 仓库标准路径/命名/格式产出——先用 load_skill 加载 project-directory-spec 拿到权威目录结构与路径,再据此检查(交付文件在机构工作空间的 projects/{项目}/docs/、apps/、modules/ 下,不要凭记忆找旧 repos/ 目录) 2. 项目过程规范:是否遵循各阶段流程规范(develop 是否实际产出代码文件、test 是否覆盖充分、deploy 配置是否完整、requirement 是否明确部署环境需求) @@ -693,9 +706,11 @@ QC_SYSTEM_PROMPT = """你是质量控制工程师(QC)。对交付件做合 ## 技能 __ROLE_SKILLS__ 遇到需要具体规范、目录结构、路径、格式、流程的任务,先用 load_skill 加载对应技能全文,不要凭记忆瞎写。 +先 load_skill 加载被查角色对应的审查技能(agent.requirement→review-requirement、agent.design→review-design、agent.develop→review-develop、agent.test→review-test、agent.deploy_*→review-deploy),按其核对项清单执行,不要凭记忆检查。 ## 工具 - load_skill(name, file_path?) — 按需加载技能全文或 references/scripts/templates 子文件(需要具体规范/路径/格式时用) +- list_features() — 查功能清单(审查需求/设计时用于功能落库与需求覆盖核对) - read_file(path) — 读工作空间文件 - list_files(path) — 列目录 - git_status() — 查看git状态 @@ -703,14 +718,14 @@ __ROLE_SKILLS__ ## 输出格式(每次一个JSON) 查看文件:{"action":"tool_call","tool":"read_file","params":{"path":"相对路径"}} -检查通过:{"action":"review_approve","comment":"合规,质量达标"} -不合规退回:{"action":"review_reject","comment":"退回原因","questions":"具体问题清单,逐条列出"} +检查通过:{"action":"review_approve","comment":"契合度 9.75(通过 39/总 40):……摘要……","score":9.75,"passed":39,"total":40} +不合规退回:{"action":"review_reject","comment":"契合度 8.50(通过 17/总 20)未达 9.5","score":8.5,"passed":17,"total":20,"questions":"[#1] 位置:…… 问题:…… 改法:……(逐条列出未过项)"} ## 检查流程 -1. 读交付件内容,检查是否按 SDLC 规范产出(路径/命名/格式) -2. 检查过程合规(文件实际产出、目录结构、必填项) -3. 判断产出质量是否达标(完整、可量化、可验收) -4. 合规则 review_approve;不合规则 review_reject 并逐条列出问题清单""" +1. load_skill 加载对应审查技能,构建核对项清单 +2. 逐项机械判定(用 list_features/read_file/list_files/run_shell 核实真实性,不轻信交付说明) +3. 算契合度 = 10 × 通过项/总项 +4. > 9.5 且无硬门禁未过项 → review_approve;否则 review_reject 并逐条列未过项""" # ── Agent 核心逻辑 ── @@ -2120,7 +2135,11 @@ async def _load_skill_by_name(sor, project_id, role, org_id, name, file_path=Non async def role_agent_run(project_id, role, agent_id=None, model_name=None): role, _role_specific, _ = await _resolve_role(project_id, role) if role == 'pm': - return {"status": "idle", "message": "PM请使用 pm_review_run"} + # PM 常规任务走 pm_review_run(pm_poller,state='review')。 + # 例外:项目复盘任务(task_kind=retrospective)是 submitted/role_task 队列, + # 由 agent_poller 派发到这里。复盘执行器自认领(CAS 防双执行); + # 无复盘任务时返回 idle,与旧行为一致(PM 提示语仅保留给常规路径)。 + return await retrospective_run(project_id, agent_id=agent_id, model_name=model_name) db = _get_db() async with db.sqlorContext("pipeline") as sor: @@ -2844,6 +2863,17 @@ async def pm_review_run(project_id, agent_id=None, model_name=None): status = 'approved' decision['status'] = 'approved' + # 防御(2026-08-31 QC 门禁改造):需求/设计/开发三阶段的质量判定权在 QC(契合度>9.5 放行), + # PM 不再复审内容。LLM 偶发无视 prompt 输出 review_reject 时,代码兜底修正为 approved 推进—— + # 否则 PM 否决与 QC 放行互相打架,交付件在 review↔submitted 间震荡,任务链停摆。 + # 系统性缺陷仍允许走 review_rollback(逃逸阀),这里只拦内容否决型 rejected。 + if status == 'rejected' and task_role in ('agent.requirement', 'agent.design', 'agent.develop'): + logger.warning(f"pm_review_run {task_role} 交付件已过 QC 门禁,PM 内容否决无效,修正为 approved: task={task_id}") + comment = f"(QC 门禁已通过,PM 内容否决被系统修正){comment}" + status = 'approved' + decision['status'] = 'approved' + decision['comment'] = comment + if status == 'approved': from appPublic.uniqueID import getID pm_did = getID() @@ -2899,64 +2929,10 @@ async def pm_review_run(project_id, agent_id=None, model_name=None): logger.info(f"模块级 develop approved,跳过 deploy_test(部署以应用为单位): {task_id}") return {"status": "approved", "task_id": task_id, "comment": comment, "next_role": "", "skip_next": "module_not_deployed_independently"} - # 人工确认节点:requirement/design 审核通过后,不直接派发下一角色任务, - # 而是发「需求/设计确认」人工任务给项目 owner。owner 确认通过 → confirm_stage_gate - # 继续派发下一角色;不确定 → 退回原角色重做(带修改意见)。 - if task_role in ('agent.requirement', 'agent.design'): - confirm_type = 'requirement_confirmation' if task_role == 'agent.requirement' else 'design_confirmation' - confirm_label = '需求' if task_role == 'agent.requirement' else '设计' - # 确认去重(2026-08-27 修复):同迭代已确认过(status=done)且本任务非回退重做 - # → 不再重复发确认任务,直接走确认后派发。元景故障:错位的重复 design 任务各被 - # approved 一次,owner 连续收到两条「设计确认」待办。回退重做(有 rollback_from) - # 产出的是新交付件,必须重新确认,不走去重。 - iter_id = '' - if task_iter: - _it2 = await sor.sqlExe( - "SELECT id FROM sd_iterations WHERE project_id=${pid}$ AND iteration_name=${n}$", - {"pid": project_id, "n": task_iter}) - iter_id = getattr(_it2[0], 'id', '') if _it2 else '' - is_redo = bool(_tp.get('rollback_from')) - if iter_id and not is_redo: - _done_conf = await sor.sqlExe( - "SELECT id FROM pipeline_human_tasks WHERE iteration_id=${iid}$ " - "AND task_type=${ct}$ AND status='done' LIMIT 1", - {"iid": iter_id, "ct": confirm_type}) - await sor.sqlExe("COMMIT", {}) - if _done_conf: - logger.info( - f"确认去重:{confirm_type} 本迭代已确认过(非回退重做)," - f"跳过重复确认,直接派发 {next_role}: task={task_id}") - next_tid, next_title = await _create_next_task( - sor, project_id, task, next_role, comment, - next_title=decision.get('next_title', '') or '', - next_desc=decision.get('next_desc', '') or '') - return {"status": "approved", "task_id": task_id, - "next_task_id": next_tid, "next_role": next_role, - "comment": comment, "confirm_skipped": "already_confirmed"} - _own = await sor.sqlExe( - "SELECT created_by FROM sd_projects WHERE id=${pid}$", {"pid": project_id}) - await sor.sqlExe("COMMIT", {}) - owner_id = getattr(_own[0], 'created_by', '') if _own else '' - if not owner_id: - owner_id = 'user-01' # 历史项目 NULL 回填 admin - from .human_task_capability import create_human_task - ok, hid = await create_human_task( - project_id, - f"{confirm_label}确认:{title}", - f"请项目 owner 确认「{confirm_label}」是否满足要求。\n\n" - f"· 确认通过 → 继续后续任务({next_role})。\n" - f"· 不确定 → 请填写修改意见,将退回 {task_role} 重做。\n\n" - f"原任务:{title}\nPM 审核意见:{comment or '通过'}", - task_type=confirm_type, - assignee_id=owner_id, - iteration_id=iter_id, - created_by='agent.pm', - task_id=task_id, - ) - logger.info(f"人工确认节点:{task_role} approved → 创建 {confirm_type} " - f"给 owner={owner_id}, human_task={'OK' if ok else 'FAIL:' + str(hid)}") - return {"status": "approved", "task_id": task_id, "next_role": next_role, - "human_confirm": hid if ok else "", "comment": comment} + # 需求/设计人工确认节点已移除(2026-08-31 QC 门禁改造): + # 需求分析/设计评审不再设人工审核,QC 门禁(qc_review 阶段,契合度>9.5 放行) + # 是这两个阶段的质量权威;QC 通过后 PM 只做编排派发,直接 _create_next_task + # 创建下一角色任务(原 requirement_confirmation/design_confirmation 人工确认删除)。 # 把 PM 给出的 next_task_title / next_task_description 传下去(原实现只传 comment, # PM 的这两个字段被静默丢弃 → 派生任务只能继承上游 title/description,正是 # 「design 任务名跟需求任务名一样」和「应用级 develop 拿到需求描述」的根源)。 @@ -2976,6 +2952,8 @@ async def pm_review_run(project_id, agent_id=None, model_name=None): logger.warning(f"_check_orchestration_gaps failed: {_e}") return {"status": "approved", "task_id": task_id, "next_task_id": next_tid, "next_role": next_role, "comment": comment} else: + # 项目完结(无下一阶段)→ 后置创建项目复盘任务(非阻塞,失败不影响完结) + await _spawn_retrospective(sor, project_id, source_task_id=task_id) return {"status": "completed", "task_id": task_id, "comment": comment or "项目完成"} elif status == 'rejected': @@ -3023,6 +3001,8 @@ async def pm_review_run(project_id, agent_id=None, model_name=None): {"cm": comment, "tid": task_id}) from .task_capability import complete_task await complete_task(task_id, project_id, who="agent.pm", agent_id=agent_id) + # 项目完结(review_complete 路径)→ 后置创建项目复盘任务(幂等,非阻塞) + await _spawn_retrospective(sor, project_id, source_task_id=task_id) return {"status": "completed", "task_id": task_id, "comment": comment or "项目完成"} @@ -3108,6 +3088,12 @@ async def qc_review_run(project_id, agent_id=None, model_name=None): msgs = [{"role": "system", "content": qc_system}] msgs.append({"role": "user", "content": f"请检查以下交付件(类型:{deliverable_type}):\n\n{content_preview}"}) + # 能力工具上下文(list_features 需要 project_id;QC 审查需求/设计时核对功能落库) + capability_ctx = { + "project_id": project_id, "iteration_id": "", "who": "agent.qc", + "agent_id": agent_id, "task_id": task_id, "org_id": org_id or '0', + } + from .llm_bridge import llm_call_msgs decision = None @@ -3141,7 +3127,7 @@ async def qc_review_run(project_id, agent_id=None, model_name=None): if tool == 'load_skill': result = await _load_skill_by_name(sor, project_id, 'qc', org_id, params.get('name', ''), params.get('file_path') or None) else: - result = await _exec_agent_tool(tool, params, space_dir) + result = await _exec_agent_tool(tool, params, space_dir, capability_ctx) msgs.append({"role": "assistant", "content": raw}) msgs.append({"role": "user", "content": f"工具 {tool} 结果:\n{result}"}) else: @@ -3186,6 +3172,198 @@ async def qc_review_run(project_id, agent_id=None, model_name=None): return {"status": "rejected", "task_id": task_id, "comment": comment, "question": rejection_q} +# ── 项目复盘(2026-08-31 新增:项目完结后 PM 回顾 → 技能提议)── + +RETRO_SYSTEM_PROMPT = """你是项目经理(PM),正在执行「项目复盘」任务。项目已执行完毕,你负责回顾执行过程,把遇到的问题和处理方法沉淀为技能提议,让同类问题在下个项目被预防。 + +## 工作环境 +工作空间:__WORKSPACE__ +项目目录:__PROJECT_DIR__ + +## 技能 +__ROLE_SKILLS__ +先 load_skill 加载 role 技能(含「项目复盘」章节的流程与纪律),再执行。 + +## 工具 +- load_skill(name) — 加载技能全文(复盘流程在 role 技能里) +- project_retrospective_data() — 取本项目全部问题素材(冒泡/退回重做/编排缺口/Bug 四类,含解决方法)。素材已附在本轮输入中,此工具供你重新拉取。 +- propose_skill(name, description, content) — 提交技能提议(content 为 SKILL.md 草稿,头部标 ,四段:触发条件/问题现象/根因/处理方法) +- write_file(path, content) — 写复盘报告 +- read_file(path) / list_files(path) — 读文件/列目录 + +## 流程 +1. 通读下方问题素材,逐条判定可复用性(跨项目会再发生=可复用;本项目特有的业务偏差=一次性)。 +2. 可复用问题:每条(同目标同要点合并)调 propose_skill 提交;目标技能按问题发生角色定位。 +3. 一次性问题:不提提议,只在报告中写明理由。 +4. write_file 写复盘报告到 projects/__PROJECT_NAME__/docs/02-retrospective/retrospective.md(五段结构:统计/问题清单/提议清单/未提议问题及理由/结论,对照 sdlc-repo-standard 技能)。 +5. 输出 deliver 提交。 + +## 输出格式(每次一个JSON) +调用工具:{"action":"tool_call","tool":"propose_skill","params":{"name":"...","description":"...","content":"..."}} +提交:{"action":"deliver","summary":"复盘完成:N 条提议","result":"复盘摘要"}""" + + +async def _spawn_retrospective(sor, project_id, source_task_id=''): + """项目完结后创建项目复盘任务(幂等:每项目一个未取消的复盘任务)。 + + 复盘是后置非阻塞任务:失败不影响项目 completed 状态,只走 failed 冒泡。 + """ + try: + dup = await sor.sqlExe( + "SELECT id FROM pipeline_tasks WHERE tenant_id=${pid}$ AND role='agent.pm' " + "AND params LIKE ${pat}$ AND state<>'cancelled' LIMIT 1", + {"pid": project_id, "pat": '%"task_kind": "retrospective"%'}) + await sor.sqlExe("COMMIT", {}) + if dup: + return '' + from appPublic.uniqueID import getID + tid = getID() + params = { + "task_kind": "retrospective", + "source_task_id": source_task_id or '', + "description": ("项目执行完毕,对执行过程做复盘:收集本项目遇到的问题和处理方法," + "判定可复用性,可复用问题提交技能提议(进技能管理建议列表)," + "产出复盘报告落盘(格式见 sdlc-repo-standard 复盘章节)。"), + } + await sor.C('pipeline_tasks', { + 'id': tid, 'tenant_id': project_id, 'pipeline_id': 'role_task', + 'owner_id': 'pm', 'title': '项目复盘', + 'params': json.dumps(params, ensure_ascii=False), + 'role': 'agent.pm', 'state': 'submitted', 'claimed_by': None, + }) + await sor.sqlExe("COMMIT", {}) + logger.info(f"retrospective spawned: task={tid} project={project_id}") + return tid + except Exception as e: + logger.warning(f"_spawn_retrospective failed: project={project_id} err={e}") + return '' + + +async def retrospective_run(project_id, agent_id=None, model_name=None): + """PM 项目复盘执行器:认领复盘任务 → 收集问题素材 → LLM 判定可复用性并提交技能提议 + → 写复盘报告 → 任务 completed。失败走 mark_failed(failed poller 冒泡),不影响项目状态。 + """ + from appPublic.uniqueID import getID + db = _get_db() + async with db.sqlorContext("pipeline") as sor: + model_name, org_id = await _resolve_llm_context(sor, project_id, 'pm', model_name) + # 认领(CAS:claimed_by IS NULL 防双执行) + recs = await sor.sqlExe( + "SELECT id, title, params FROM pipeline_tasks " + "WHERE tenant_id=${pid}$ AND state='submitted' AND role='agent.pm' " + "AND params LIKE ${pat}$ AND (claimed_by IS NULL OR claimed_by='') " + "ORDER BY created_at ASC LIMIT 1", + {"pid": project_id, "pat": '%"task_kind": "retrospective"%'}) + await sor.sqlExe("COMMIT", {}) + if not recs: + return {"status": "idle", "message": "没有待复盘任务"} + task_id = getattr(recs[0], 'id', '') + claim_token = agent_id or getID() + await sor.sqlExe( + "UPDATE pipeline_tasks SET state='running', claimed_by=${cb}$, updated_at=NOW() " + "WHERE id=${tid}$ AND state='submitted' AND (claimed_by IS NULL OR claimed_by='')", + {"cb": claim_token, "tid": task_id}) + chk = await sor.sqlExe("SELECT state FROM pipeline_tasks WHERE id=${tid}$", {"tid": task_id}) + await sor.sqlExe("COMMIT", {}) + if not chk or getattr(chk[0], 'state', '') != 'running': + return {"status": "idle", "message": "复盘任务认领竞争失败"} + + space_dir = await _get_space_dir(sor, project_id) + project_name = await _get_project_name(sor, project_id) + project_dir = os.path.join(space_dir, 'projects', project_name) if project_name else space_dir + await sor.sqlExe("COMMIT", {}) + + # 问题素材(代码统一收集,PM 不重复实现查询) + from .retrospective_capability import project_retrospective_data + materials = await project_retrospective_data(project_id) + + role_skills = await _build_role_skills_block(sor, project_id, _normalize_role('pm'), org_id) + retro_system = RETRO_SYSTEM_PROMPT\ + .replace('__WORKSPACE__', space_dir)\ + .replace('__PROJECT_DIR__', project_dir)\ + .replace('__PROJECT_NAME__', project_name or '')\ + .replace('__ROLE_SKILLS__', role_skills) + + msgs = [{"role": "system", "content": retro_system}] + msgs.append({"role": "user", "content": f"本项目的问题素材(JSON):\n\n{materials}"}) + + capability_ctx = { + "project_id": project_id, "iteration_id": "", "who": "agent.pm", + "agent_id": agent_id, "task_id": task_id, "org_id": org_id or '0', + } + + from .llm_bridge import llm_call_msgs + deliverable = None + proposals_made = 0 + for turn in range(10): + await sor.sqlExe( + "UPDATE pipeline_tasks SET updated_at=NOW() WHERE id=${tid}$ AND state='running'", + {"tid": task_id}) + await sor.sqlExe("COMMIT", {}) + if turn >= 7: + msgs.append({"role": "user", "content": + "已执行足够轮次。现在必须立即收尾:write_file 写复盘报告(若未写),然后输出 deliver 提交,禁止再调其它工具。"}) + try: + raw = await llm_call_msgs(msgs, model=model_name, temperature=0.3, org_id=org_id) + except Exception as e: + err_msg = f"{type(e).__name__}: {str(e)[:400]}" + from .task_capability import mark_failed + await mark_failed(task_id, project_id, who="agent.pm", agent_id=agent_id, error=err_msg) + return {"status": "failed", "task_id": task_id, "error": str(e)[:200]} + + act = _parse_agent_action(raw) + if act.get('action') == 'deliver': + deliverable = act + break + elif act.get('action') == 'tool_call': + tool = act.get('tool', '') + params = act.get('params', {}) + if tool == 'load_skill': + result = await _load_skill_by_name(sor, project_id, 'pm', org_id, params.get('name', ''), params.get('file_path') or None) + else: + result = await _exec_agent_tool(tool, params, space_dir, capability_ctx) + if tool == 'propose_skill' and str(result).startswith('OK'): + proposals_made += 1 + msgs.append({"role": "assistant", "content": raw}) + msgs.append({"role": "user", "content": f"工具 {tool} 结果:\n{result}"}) + else: + deliverable = {"result": raw} + break + + if not deliverable: + from .task_capability import mark_failed + await mark_failed(task_id, project_id, who="agent.pm", agent_id=agent_id, + error="复盘任务轮次耗尽未产出交付件") + return {"status": "failed", "task_id": task_id, "error": "轮次耗尽未交付"} + + # 复盘报告落盘:规范路径 + deliverables/pm/ + result_text = deliverable.get("result") or deliverable.get("summary") or "" + report_rel = os.path.join('projects', project_name or '', 'docs', '02-retrospective', 'retrospective.md') + try: + os.makedirs(os.path.dirname(os.path.join(space_dir, report_rel)), exist_ok=True) + with open(os.path.join(space_dir, report_rel), 'w', encoding='utf-8') as f: + f.write(result_text or '') + except Exception as e: + logger.warning(f"retrospective report write failed: {e}") + from appPublic.uniqueID import getID as _gid + did = _gid() + await sor.C("pipeline_deliverables", { + "id": did, "project_id": project_id, "task_id": task_id, + "deliverable_type": "retrospective", "title": "项目复盘报告", + "content": result_text[:60000], + "file_path": os.path.join(project_dir, 'deliverables', 'pm', f"{task_id}_retrospective.md"), + "quality_score": 100, "review_status": "approved", "created_by": agent_id or "pm", + }) + await sor.sqlExe("COMMIT", {}) + # 复盘任务直接 completed(后置任务不走 QC/PM 审核环;失败另有 mark_failed 冒泡) + from .task_capability import set_task_state + await set_task_state(task_id, project_id, 'running', 'completed', + who="agent.pm", agent_id=agent_id, + detail=f"复盘完成,提交技能提议 {proposals_made} 条") + logger.info(f"retrospective_run completed: task={task_id} proposals={proposals_made}") + return {"status": "completed", "task_id": task_id, "proposals": proposals_made} + + # ── 失败任务处理(第二层:PM/cockpit 判定重跑,最多3次,超限报故障给用户)── _RETRY_HINTS = ( diff --git a/pipeline_service/capability_tools.py b/pipeline_service/capability_tools.py index c3aa3a3..b7d72f1 100644 --- a/pipeline_service/capability_tools.py +++ b/pipeline_service/capability_tools.py @@ -216,6 +216,20 @@ TOOL_SCHEMAS = { "params": {"env_type": "类型(可选 test/prod)", "status": "状态(可选)"}, "required": [], }, + # ── retrospective_capability(项目复盘:问题收集 + 技能提议)── + "project_retrospective_data": { + "module": "retrospective_capability", + "description": "取本项目全部问题素材(冒泡/退回重做/编排缺口/Bug 四类,含解决方法),复盘分析用", + "params": {}, + "required": [], + }, + "propose_skill": { + "module": "retrospective_capability", + "description": "提交技能提议(复盘产出,进技能管理建议列表待人工审核,不直接生效)", + "params": {"name": "技能名", "description": "一句话描述", + "content": "SKILL.md 草稿(头部 定位,四段:触发条件/问题现象/根因/处理方法)"}, + "required": ["name", "content"], + }, } diff --git a/pipeline_service/retrospective_capability.py b/pipeline_service/retrospective_capability.py new file mode 100644 index 0000000..ccc056b --- /dev/null +++ b/pipeline_service/retrospective_capability.py @@ -0,0 +1,170 @@ +"""复盘能力 — 项目执行完后的问题收集 + 技能提议(通用、产线无关)。 + +职责(方案 P2/P3 的代码层): +- project_retrospective_data:四类问题源统一收集(冒泡问题/退回重做/编排缺口/Bug), + 返回结构化素材给 PM 复盘任务,PM 不重复实现查询(「重复逻辑写 py」原则)。 +- propose_skill:PM 复盘产出的技能提议落 skill_proposals(source=agent, status=pending), + 进技能管理建议列表待人工审核,不直接改技能库(门禁守源头)。 +""" + +import json +import logging + +from sqlor.dbpools import DBPools +from appPublic.uniqueID import getID + +DBNAME = "pipeline" +logger = logging.getLogger("pipeline.retrospective_capability") + + +def _get_db(): + db = DBPools() + if not db.databases: + from appPublic.jsonConfig import getConfig + config = getConfig() + if config and config.databases: + db.databases = config.databases + return db, DBNAME + + +def _rows_to_dicts(recs, limit=200): + out = [] + for r in (recs or [])[:limit]: + d = {} + for k, v in vars(r).items(): + if k.startswith('_') or callable(v): + continue + d[k] = v + out.append(d) + return out + + +async def project_retrospective_data(project_id, who=None, agent_id=None, + task_id=None, iteration_id=None): + """收集本项目执行中的全部问题素材(四类源合并 + 统计)。返回结构化文本(JSON)。 + + 四类问题源(全部已落库,本函数只查不写): + ① 冒泡问题 pipeline_agent_questions:question + 解决记录(answer/answered_by) + ② 退回重做 audit_log:qc_reject/reject 记录(detail=退回意见,who=被退角色) + ③ 编排缺口 pipeline_pm_notices:gap_text + ④ Bug sd_bugs:title/description/fix_description(按项目迭代关联) + """ + if not project_id: + return "FAIL: 缺少 project_id" + db, dbname = _get_db() + async with db.sqlorContext(dbname) as sor: + # ① 冒泡问题(含解决记录) + bubbles = _rows_to_dicts(await sor.sqlExe( + "SELECT problem_type, question, from_role, current_handler_role, status, " + "answer, answered_by, created_at FROM pipeline_agent_questions " + "WHERE tenant_id=${pid}$ ORDER BY created_at ASC LIMIT 100", + {"pid": project_id})) + # ② 退回重做(审计:谁被退、退回意见) + reworks = _rows_to_dicts(await sor.sqlExe( + "SELECT action, who, detail, created_at FROM audit_log " + "WHERE tenant_id=${pid}$ AND entity='pipeline_tasks' " + "AND action IN ('qc_reject','reject') ORDER BY created_at ASC LIMIT 100", + {"pid": project_id})) + # ③ 编排缺口 + gaps = _rows_to_dicts(await sor.sqlExe( + "SELECT gap_text, status FROM pipeline_pm_notices " + "WHERE tenant_id=${pid}$ ORDER BY created_at ASC LIMIT 50", + {"pid": project_id})) + # ④ Bug(项目 → 迭代 → Bug) + iters = _rows_to_dicts(await sor.sqlExe( + "SELECT id, iteration_name FROM sd_iterations WHERE project_id=${pid}$", + {"pid": project_id})) + bugs = [] + if iters: + ids = ",".join("'" + str(i.get('id', '')).replace("'", "") + "'" for i in iters) + bugs = _rows_to_dicts(await sor.sqlExe( + "SELECT title, description, severity, status, fix_description, created_at " + "FROM sd_bugs WHERE iteration_id IN (" + ids + ") " + "ORDER BY created_at ASC LIMIT 100", {})) + # 统计:工期 + precs = await sor.sqlExe( + "SELECT name, created_at, status FROM sd_projects WHERE id=${pid}$", + {"pid": project_id}) + await sor.sqlExe("COMMIT", {}) + proj = _rows_to_dicts(precs) + proj = proj[0] if proj else {} + + problems = [] + for b in bubbles: + problems.append({ + "type": "bubble", "problem": b.get('question', ''), + "solution": b.get('answer', '') or '(未记录解决方法)', + "problem_type": b.get('problem_type', ''), + "roles": "%s→%s" % (b.get('from_role', ''), b.get('current_handler_role', '')), + "time": str(b.get('created_at', ''))}) + for r in reworks: + problems.append({ + "type": "rework", "problem": r.get('detail', '') or '(无退回意见)', + "solution": '(重做后通过,对比新旧交付件可见改法)', + "role": r.get('who', ''), "via": r.get('action', ''), + "time": str(r.get('created_at', ''))}) + for g in gaps: + problems.append({ + "type": "gap", "problem": g.get('gap_text', ''), + "solution": '(PM 处置记录见任务链)', "time": ''}) + for bug in bugs: + problems.append({ + "type": "bug", "problem": "%s:%s" % (bug.get('title', ''), (bug.get('description', '') or '')[:300]), + "solution": bug.get('fix_description', '') or '(未记录修复说明)', + "severity": bug.get('severity', ''), "status": bug.get('status', ''), + "time": str(bug.get('created_at', ''))}) + + data = { + "project": proj.get('name', ''), "project_status": proj.get('status', ''), + "project_created_at": str(proj.get('created_at', '')), + "problems": problems, + "stats": { + "bubble_count": len(bubbles), "rework_count": len(reworks), + "gap_count": len(gaps), "bug_count": len(bugs), + "total": len(problems), + }, + } + return json.dumps(data, ensure_ascii=False, default=str) + + +async def propose_skill(project_id, name, description="", content="", + who=None, agent_id=None, task_id=None, iteration_id=None): + """提交技能提议(复盘产出)。落 skill_proposals,source=agent,status=pending。 + + 提议只是草稿:进技能管理建议列表待人工审核(pending→testing→approved→published), + 不直接改技能库。content 为 SKILL.md 草稿,头部应带目标定位注释 + ( 或 )。 + """ + name = (name or '').strip() + content = (content or '').strip() + if not name or not content: + return False, "需要技能名和内容(content 为 SKILL.md 草稿正文)" + if len(content) < 80: + return False, "提议内容过短(<80字):SKILL.md 草稿须含触发条件/问题现象/根因/处理方法" + db, dbname = _get_db() + async with db.sqlorContext(dbname) as sor: + # org/pipeline 从项目取(提议按机构隔离 + 产线归属) + precs = await sor.sqlExe( + "SELECT org_id, pipeline_id FROM sd_projects WHERE id=${pid}$", + {"pid": project_id}) + await sor.sqlExe("COMMIT", {}) + org_id = '0' + pipeline_id = '' + if precs: + org_id = getattr(precs[0], 'org_id', '') or '0' + pipeline_id = getattr(precs[0], 'pipeline_id', '') or '' + pid = getID() + await sor.C('skill_proposals', { + 'id': pid, + 'name': name[:100], + 'description': (description or '')[:500], + 'content': content, + 'source': 'agent', + 'status': 'pending', + 'org_id': org_id, + 'pipeline_id': pipeline_id, + 'created_by': who or 'agent.pm', + }) + await sor.sqlExe("COMMIT", {}) + logger.info("propose_skill: %s project=%s by=%s", name, project_id, who) + return True, f"已提交技能提议 '{name}'(待审核,暂不生效)"