# -*- 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}