diff --git a/pipeline_service/agent_loop.py b/pipeline_service/agent_loop.py index 9b33fe9..6f10661 100644 --- a/pipeline_service/agent_loop.py +++ b/pipeline_service/agent_loop.py @@ -2171,6 +2171,29 @@ async def _build_fallback_deliverable(space_dir, written_files, role): } +# ── 技能树节流刷新(2026-09-08 技能提议实时生效配套)── +# v1 角色 agent 跑在 poller 长进程里,loader 是全局单例;agent 实时发布的 +# 机构技能(orgs/{org_id}/,NFS 共享)要靠 reload 才能被长进程看到。 +# 全树扫描有成本(数百技能目录),30s 节流:同进程 30s 内多次任务只扫一次。 +_skills_reload_ts = 0.0 +_SKILLS_RELOAD_TTL = 30 + + +def _reload_skills_throttled(loader): + """节流刷新技能树(30s TTL)。失败静默(用旧树,不阻塞任务)。""" + global _skills_reload_ts + import time as _t + now = _t.time() + if now - _skills_reload_ts < _SKILLS_RELOAD_TTL: + return + try: + loader.reload() + _skills_reload_ts = now + except Exception as e: + logger.warning(f"skills reload failed: {e}") + _skills_reload_ts = now # 失败也记时,防每轮重试打爆日志 + + async def _build_role_skills_block(sor, project_id, role, org_id=""): """为角色 agent 构建技能目录块(分层导入第一层)。 @@ -2185,6 +2208,7 @@ async def _build_role_skills_block(sor, project_id, role, org_id=""): from pipeline_core.skill_pack import get_skills_base skills_dir = get_skills_base() loader = get_skill_loader(skills_dir) + _reload_skills_throttled(loader) merged = loader.get_merged(pipeline_id=pid, role=role, project_id=project_id, org_id=org_id or '0') skills = sorted(merged.values(), @@ -2259,6 +2283,7 @@ async def _load_skill_by_name(sor, project_id, role, org_id, name, file_path=Non from pipeline_core.skill_pack import get_skills_base skills_dir = get_skills_base() loader = get_skill_loader(skills_dir) + _reload_skills_throttled(loader) merged = loader.get_merged(pipeline_id=pid, role=role, project_id=project_id, org_id=org_id or '0') from pipeline_core.skill_loader import resolve_skill, normalize_skill_name, format_skill_names @@ -3412,7 +3437,7 @@ __ROLE_SKILLS__ ## 工具 - load_skill(name) — 加载技能全文(复盘流程在 role 技能里) - project_retrospective_data() — 取本项目全部问题素材(冒泡/退回重做/编排缺口/Bug 四类,含解决方法)。素材已附在本轮输入中,此工具供你重新拉取。 -- propose_skill(name, description, content) — 提交技能提议(content 为 SKILL.md 草稿,头部标 ,四段:触发条件/问题现象/根因/处理方法) +- propose_skill(name, description, content) — 沉淀技能并实时发布到项目所属机构技能目录(机构内立即生效、同名覆盖通用技能、其他机构不受影响;content 为 SKILL.md 正文,头部标 ,四段:触发条件/问题现象/根因/处理方法。平台缺省机构 org 0 的提议自动转人工审核) - write_file(path, content) — 写复盘报告 - read_file(path, offset?) / list_files(path) — 读文件/列目录(docx/pdf 自动解析;大文件按截断提示用 offset 续读) diff --git a/pipeline_service/agent_loop_v2.py b/pipeline_service/agent_loop_v2.py index 03b9c5a..d98121e 100644 --- a/pipeline_service/agent_loop_v2.py +++ b/pipeline_service/agent_loop_v2.py @@ -1354,18 +1354,32 @@ class AgentExecutor: return "FAIL: 需要技能名和内容(content 为 SKILL.md 草稿正文)" try: from appPublic.uniqueID import getID + # 实时发布(2026-09-08 用户拍板):写提议者所属机构技能目录 + # skills/orgs/{org_id}/{name}/,org scope 同名覆盖 global, + # 机构内立即生效、其他机构不受影响;org 0 有用户身份时降级 + # users/{user_id}/(只影响本人),org 0 无人值守拒绝(防影响全平台)。 + from .skill_live import publish_skill_live + ok, msg, info = await publish_skill_live( + name, description, content, + org_id=self.org_id or "", user_id=self.user_id or "", + who=self.user_id or "") await sor.C("skill_proposals", { "id": getID(), "name": name, "description": description, "content": content, "source": "agent", - "status": "pending", + # 实时发布成功=published(审计痕迹);被拒=保持 pending 人工审核链 + "status": "published" if ok else "pending", + "feedback": msg if ok else "", "org_id": self.org_id or "", "pipeline_id": self.pipeline_id or "", "created_by": self.user_id or "", }) - return f"OK: 已提交技能提议 '{name}'(待审核,暂不生效)" + if ok: + return msg + # 实时发布被拒(org 0 无人值守等)→ 保持旧行为:提议待审核 + return f"OK: 已提交技能提议 '{name}'({msg},转人工审核,暂不生效)" except Exception as e: return f"ERROR: {str(e)[:300]}" diff --git a/pipeline_service/capability_tools.py b/pipeline_service/capability_tools.py index 68b0ef5..d9abde9 100644 --- a/pipeline_service/capability_tools.py +++ b/pipeline_service/capability_tools.py @@ -225,9 +225,9 @@ TOOL_SCHEMAS = { }, "propose_skill": { "module": "retrospective_capability", - "description": "提交技能提议(复盘产出,进技能管理建议列表待人工审核,不直接生效)", - "params": {"name": "技能名", "description": "一句话描述", - "content": "SKILL.md 草稿(头部 定位,四段:触发条件/问题现象/根因/处理方法)"}, + "description": "沉淀技能(复盘产出):实时发布到项目所属机构技能目录(orgs/{org_id}/,同名覆盖通用技能,机构内立即生效,其他机构不受影响);平台缺省机构(org 0)的无人值守提议转人工审核不实时生效", + "params": {"name": "技能名(字母数字._-,≤64字符)", "description": "一句话描述", + "content": "SKILL.md 正文(头部 定位,四段:触发条件/问题现象/根因/处理方法)"}, "required": ["name", "content"], }, } diff --git a/pipeline_service/retrospective_capability.py b/pipeline_service/retrospective_capability.py index 8c4e37f..1f9fb4d 100644 --- a/pipeline_service/retrospective_capability.py +++ b/pipeline_service/retrospective_capability.py @@ -140,10 +140,15 @@ async def project_retrospective_data(project_id, who=None, agent_id=None, 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 草稿,头部应带目标定位注释 + 2026-09-08 用户拍板改为实时生效:写提议所属机构的技能目录 + skills/orgs/{org_id}/{name}/(org scope 同名覆盖 global,机构内立即生效, + 其他机构不受影响)。org 0(平台缺省机构,orgs/0/ 全机构共享)的无人值守 + 提议不允许实时发布——保持原 pending 人工审核链(pending→testing→approved + →published)。skill_proposals 照落一条留审计(实时发布=published+路径)。 + + content 为 SKILL.md 草稿,头部应带目标定位注释 ()。 """ name = (name or '').strip() @@ -164,6 +169,13 @@ async def propose_skill(project_id, name, description="", content="", if precs: org_id = getattr(precs[0], 'org_id', '') or '0' pipeline_id = getattr(precs[0], 'pipeline_id', '') or '' + + # 实时发布(复盘是无人值守:who 是 agent.pm 不是人,无 user_id—— + # org 0 项目会被 resolve_target 拒绝实时写,自动回落人工审核链) + from .skill_live import publish_skill_live + ok, msg, info = await publish_skill_live( + name, description, content, org_id=org_id, user_id='', who=who or 'agent.pm') + pid = getID() await sor.C('skill_proposals', { 'id': pid, @@ -171,11 +183,14 @@ async def propose_skill(project_id, name, description="", content="", 'description': (description or '')[:500], 'content': content, 'source': 'agent', - 'status': 'pending', + 'status': 'published' if ok else 'pending', + 'feedback': msg if ok else '', '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}'(待审核,暂不生效)" + logger.info("propose_skill: %s project=%s by=%s live=%s", name, project_id, who, ok) + if ok: + return True, msg + return True, f"已提交技能提议 '{name}'({msg},转人工审核,暂不生效)" diff --git a/pipeline_service/skill_live.py b/pipeline_service/skill_live.py new file mode 100644 index 0000000..9ba8d13 --- /dev/null +++ b/pipeline_service/skill_live.py @@ -0,0 +1,136 @@ +# -*- coding: utf-8 -*- +"""skill_live.py - agent 技能提议实时生效(2026-09-08 用户需求变更) + +旧机制:propose_skill 只写 skill_proposals 表(status=pending),人工门禁 +pending→testing→approved→published 后才生效——滞后数小时到数天。 + +新机制(用户拍板):**实时生效,但只写提议者所属机构的技能目录**: +- 客户机构(org_id≠'0')→ skills/orgs/{org_id}/{name}/SKILL.md + org scope(优先级3)在 get_merged 同名覆盖 global(优先级0)——机构内全用户 + 立即生效,其他机构完全不受影响(目录物理隔离)。 +- ⚠️ org 0 陷阱:get_merged 里 orgs/0/(ocai 缺省机构)的技能是**所有机构 + 共享**的(先加载 org0 再被各机构同名覆盖)。org_id='0' 的提议若写 orgs/0/ + 会影响全平台——违反「不影响其他机构」。故 org 0 且有 user_id 时降级写 + skills/users/{user_id}/(user scope 优先级6,只影响该用户本人); + org 0 且无 user_id(无人值守复盘)→ 拒绝实时写,保持原 pending 人工审核链。 + +落库不变:skill_proposals 照写一条(status=published + feedback 记路径), +留审计痕迹;技能管理页可看到「已实时发布」的记录。 + +跨进程实时生效:v2 会话 agent 每条消息 executor init 都 reload(已有); +v1 角色 agent poller 是长进程,loader 是全局单例——本模块写盘后同进程 +reload(),跨进程靠 v1 调用点的 loader.reload() 兜底(每轮任务开始时刷新)。 +多 worker 主机共享 NFS 技能树,各进程 reload 即可见。 + +安全:技能名/机构ID/用户ID 全部做路径穿越清洗(只允许 [A-Za-z0-9._-]), +frontmatter 的 scope 字段剥离(目录决定 scope,防 agent 自称 global)。 +""" +import logging +import os +import re + +logger = logging.getLogger("pipeline.skill_live") + +_NAME_RE = re.compile(r'^[A-Za-z0-9][A-Za-z0-9._-]{0,63}$') + + +def _sanitize_id(raw: str) -> str: + """清洗 org_id/user_id/skill name:只允许字母数字._-,防路径穿越。""" + return re.sub(r'[^A-Za-z0-9._-]', '', (raw or '').strip().replace('/', '_').replace('..', '_'))[:64] + + +def _strip_scope_frontmatter(content: str) -> str: + """剥离 frontmatter 的 scope 字段——scope 由目录决定,防 agent 自称 global/user。""" + if not content.startswith('---'): + return content + end = content.find('\n---', 3) + if end < 0: + return content + fm = content[3:end] + body = content[end + 4:] + fm_lines = [ln for ln in fm.split('\n') if not re.match(r'^\s*scope\s*:', ln)] + return '---' + '\n'.join(fm_lines) + '\n---' + body + + +def _ensure_frontmatter(content: str, name: str, description: str) -> str: + """确保 SKILL.md 带合法 frontmatter(name/description 补齐)。""" + content = _strip_scope_frontmatter(content) + if content.startswith('---'): + return content + desc = (description or '').strip().replace('\n', ' ')[:200] + return ('---\nname: %s\ndescription: %s\nversion: 1.0.0\n---\n\n%s' + % (name, desc or name, content)) + + +def resolve_target(org_id: str, user_id: str, name: str): + """决定实时写入目标。返回 (scope, target_dir) 或 ('', 拒绝原因)。 + + scope: 'org' → orgs/{org_id}/{name}/;'user' → users/{user_id}/{name}/ + """ + safe_name = _sanitize_id(name) + if not _NAME_RE.match(safe_name): + return '', '技能名非法(只允许字母数字._-,且以字母数字开头,≤64字符)' + org = _sanitize_id(org_id) + uid = _sanitize_id(user_id) + from pipeline_core.skill_pack import get_skills_base + base = get_skills_base() + if org and org != '0': + # 客户机构:机构目录实时生效,机构间物理隔离 + return 'org', os.path.join(base, 'orgs', org, safe_name) + # org 0(平台/ocai 缺省机构):orgs/0/ 是全机构共享缺省——绝不实时写 + if uid: + # 有用户身份(会话场景):降级到用户个人技能,只影响本人 + return 'user', os.path.join(base, 'users', uid, safe_name) + # 无人值守(复盘 poller)且 org 0:拒绝实时写,走人工审核 + return '', ('org 0(平台缺省机构)的无人值守提议不允许实时发布' + '(orgs/0/ 技能全机构共享,实时写会影响所有机构)——保持人工审核链') + + +async def publish_skill_live(name: str, description: str, content: str, + org_id: str = '', user_id: str = '', + who: str = ''): + """实时发布技能到机构/用户目录 + loader 刷新。 + + 返回 (ok: bool, message: str, info: dict)。 + info: {'scope': 'org'|'user', 'path': 相对技能根路径}(ok 时有值) + """ + name = (name or '').strip() + content = (content or '').strip() + if not name or not content: + return False, '需要技能名和内容', {} + scope, target = resolve_target(org_id, user_id, name) + if not scope: + return False, target, {} # target 此时是拒绝原因 + safe_name = os.path.basename(target) + try: + os.makedirs(target, exist_ok=True) + skill_path = os.path.join(target, 'SKILL.md') + final = _ensure_frontmatter(content, safe_name, description) + # 原子写:先写临时文件再 rename(防 reload 读到半截文件) + tmp_path = skill_path + '.tmp' + with open(tmp_path, 'w', encoding='utf-8') as f: + f.write(final) + os.replace(tmp_path, skill_path) + except Exception as e: + logger.error('publish_skill_live write failed: %s', e) + return False, '写入技能目录失败: %s' % str(e)[:200], {} + + # 同进程立即生效 + reloaded = False + try: + from pipeline_core.skill_loader import get_skill_loader + from pipeline_core.skill_pack import get_skills_base + get_skill_loader(get_skills_base()).reload() + reloaded = True + except Exception as e: + logger.warning('skill loader reload failed: %s', e) + + from pipeline_core.skill_pack import get_skills_base + rel = os.path.relpath(skill_path, get_skills_base()) + scope_cn = '机构' if scope == 'org' else '个人' + msg = ('OK: 技能 %r 已实时发布到%s目录(%s),%s内立即生效,' + '同名覆盖通用(global)技能,其他机构不受影响。%s' + % (safe_name, scope_cn, rel, + '本机构' if scope == 'org' else '仅你本人', + '' if reloaded else '(loader 刷新失败,下轮会话生效)')) + return True, msg, {'scope': scope, 'path': rel}