diff --git a/pipeline_core/skill_loader.py b/pipeline_core/skill_loader.py index d8124d6..697b48d 100644 --- a/pipeline_core/skill_loader.py +++ b/pipeline_core/skill_loader.py @@ -34,19 +34,21 @@ logger = logging.getLogger("pipeline.skill_loader") # scope 优先级(数值越大优先级越高) SCOPE_PRIORITY = { "global": 0, - "org": 1, + "project_common": 1, "pipeline": 2, - "project": 3, - "role": 4, - "user": 5, + "org": 3, + "project": 4, + "role": 5, + "user": 6, } SCOPE_TAG = { "global": "通用", - "org": "组织", + "project_common": "项目通用", "pipeline": "产线", - "role": "角色", + "org": "组织", "project": "项目", + "role": "角色", "user": "个人", } @@ -190,15 +192,16 @@ def _parse_frontmatter(content: str) -> Optional[dict]: # ── 六级技能加载器 ── class SkillLoader: - """六级技能加载器:global → org → pipeline → project → role → user""" + """七级技能加载器:global → project_common → pipeline → org → project → role → user""" def __init__(self, base_dir: str = ""): self.base_dir = base_dir self._global: Dict[str, Skill] = {} + self._projects_common: Dict[str, Skill] = {} # 项目通用(projects/,所有项目共享) self._orgs: Dict[str, Dict[str, Skill]] = {} # {org_id: {name: Skill}} self._pipelines: Dict[str, Dict[str, Skill]] = {} # {pipeline_id: {name: Skill}} self._roles: Dict[str, Dict[str, Skill]] = {} # {pipeline_id:role: {name: Skill}} - self._projects: Dict[str, Dict[str, Skill]] = {} # {project_id: {name: Skill}} + self._projects: Dict[str, Dict[str, Skill]] = {} # {org_id:project_id: {name: Skill}} 项目私有(挂在机构下) self._users: Dict[str, Dict[str, Skill]] = {} # {user_id: {name: Skill}} if base_dir: self.reload() @@ -212,6 +215,7 @@ class SkillLoader: def reload(self): self._global = {} + self._projects_common = {} self._orgs = {} self._pipelines = {} self._roles = {} @@ -226,13 +230,10 @@ class SkillLoader: if os.path.isdir(global_dir): self._global = self._scan_scope_dir(global_dir, "global") - # 组织: skills/orgs/{org_id}/ - orgs_dir = os.path.join(self.base_dir, "orgs") - if os.path.isdir(orgs_dir): - for org_id in os.listdir(orgs_dir): - p = os.path.join(orgs_dir, org_id) - if os.path.isdir(p): - self._orgs[org_id] = self._scan_scope_dir(p, "org", org_id) + # 项目通用: skills/projects/(所有项目共享,不分 project_id) + projects_dir = os.path.join(self.base_dir, "projects") + if os.path.isdir(projects_dir): + self._projects_common = self._scan_scope_dir(projects_dir, "project_common") # 产线: skills/pipelines/{pipeline_id}/common/ + roles/{role}/ pipelines_dir = os.path.join(self.base_dir, "pipelines") @@ -254,13 +255,27 @@ class SkillLoader: key = f"{pid}:{role}" self._roles[key] = self._scan_scope_dir(rdir, "role", key) - # 项目/会话: skills/projects/{project_id}/ - projects_dir = os.path.join(self.base_dir, "projects") - if os.path.isdir(projects_dir): - for project_id in os.listdir(projects_dir): - p = os.path.join(projects_dir, project_id) - if os.path.isdir(p): - self._projects[project_id] = self._scan_scope_dir(p, "project", project_id) + # 组织 + 项目私有: skills/orgs/{org_id}/ + # 第一层含 SKILL.md 的是 org 技能;否则是 {project_id} 项目私有子目录(挂在机构下) + orgs_dir = os.path.join(self.base_dir, "orgs") + if os.path.isdir(orgs_dir): + for org_id in os.listdir(orgs_dir): + org_path = os.path.join(orgs_dir, org_id) + if not os.path.isdir(org_path): + continue + self._orgs[org_id] = {} + for entry in os.listdir(org_path): + entry_path = os.path.join(org_path, entry) + if not os.path.isdir(entry_path): + continue + if os.path.isfile(os.path.join(entry_path, "SKILL.md")): + # org 技能(第一层直接是 skill 目录) + skill = Skill.from_file(os.path.join(entry_path, "SKILL.md"), "org", org_id) + self._orgs[org_id][skill.name] = skill + else: + # {project_id} 项目私有目录 + key = f"{org_id}:{entry}" + self._projects[key] = self._scan_scope_dir(entry_path, "project", key) # 用户: skills/users/{user_id}/ users_dir = os.path.join(self.base_dir, "users") @@ -271,6 +286,7 @@ class SkillLoader: self._users[user_id] = self._scan_scope_dir(p, "user", user_id) logger.debug(f"SkillLoader reloaded: global={len(self._global)}, " + f"projects_common={len(self._projects_common)}, " f"orgs={sum(len(s) for s in self._orgs.values())}, " f"pipelines={sum(len(s) for s in self._pipelines.values())}, " f"roles={sum(len(s) for s in self._roles.values())}, " @@ -302,20 +318,26 @@ class SkillLoader: def get_merged(self, pipeline_id: str = "", role: str = "", project_id: str = "", org_id: str = "", user_id: str = "") -> Dict[str, Skill]: - """按优先级合并六级技能。同名技能高优先级覆盖低优先级。 + """按优先级合并七级技能。同名技能高优先级覆盖低优先级。 - 优先级(低→高):global → org → pipeline → project → role → user + 优先级(低→高):global → project_common → pipeline → org → project → role → user """ merged = dict(self._global) - if org_id and org_id in self._orgs: - merged.update(self._orgs[org_id]) + # 项目通用(projects/,所有项目共享) + merged.update(self._projects_common) if pipeline_id and pipeline_id in self._pipelines: merged.update(self._pipelines[pipeline_id]) - if project_id and project_id in self._projects: - merged.update(self._projects[project_id]) + if org_id and org_id in self._orgs: + merged.update(self._orgs[org_id]) + + # 项目私有(orgs/{org_id}/{project_id}/) + if org_id and project_id: + proj_key = f"{org_id}:{project_id}" + if proj_key in self._projects: + merged.update(self._projects[proj_key]) if pipeline_id and role: role_key = f"{pipeline_id}:{role}" @@ -391,6 +413,7 @@ class SkillLoader: def list_by_scope(self) -> Dict[str, int]: return { "global": len(self._global), + "project_common": len(self._projects_common), "orgs": sum(len(s) for s in self._orgs.values()), "pipelines": sum(len(s) for s in self._pipelines.values()), "roles": sum(len(s) for s in self._roles.values()), diff --git a/pipeline_core/skill_pack.py b/pipeline_core/skill_pack.py index ed8931d..0de6bcf 100644 --- a/pipeline_core/skill_pack.py +++ b/pipeline_core/skill_pack.py @@ -35,83 +35,13 @@ ALL_DIR = os.path.join(SKILLS_BASE, "global") # 公共技能 = 全局模板( PIPELINES_DIR = os.path.join(SKILLS_BASE, "pipelines") # 产线技能(common/roles),源头在 build.sh 6b 复制 -def _sync_tree(src_dir, dst_dir): - """递归增量同步目录树:src 里 dst 缺失的文件/子目录复制过去(不覆盖已存在的)。 - - 用于产线技能 pipelines/(结构是 pipelines/{pid}/common/{skill} 与 - pipelines/{pid}/roles/{role}/{skill},比 global/packs 深一层)。 - """ - if not os.path.isdir(src_dir): - return - os.makedirs(dst_dir, exist_ok=True) - for name in os.listdir(src_dir): - src = os.path.join(src_dir, name) - dst = os.path.join(dst_dir, name) - if os.path.isdir(src): - if not os.path.exists(dst): - shutil.copytree(src, dst) - else: - _sync_tree(src, dst) - else: - if not os.path.exists(dst): - shutil.copy2(src, dst) - - -def get_org_skills_dir(workspace_base, org_id): - """机构工作目录下的技能根目录(多租户隔离)。 - - workspace_base: 机构工作目录根(如 /d/pipeline/workspaces),从 params 表 workspace_base 读。 - org_id: 机构 ID。 - 返回 //skills - """ - return os.path.normpath(os.path.join(workspace_base, str(org_id or '0'), "skills")) - - def ensure_org_skills(workspace_base, org_id): - """确保机构工作目录的 skills/ 已初始化(从全局模板复制公共技能 + 技能集清单)。 + """【单一技能树,2026-08-21 重构】技能改为所有机构共享读全局 skills/,不再复制到机构工作目录。 - 多租户隔离:运行时技能必须落在机构工作目录,不能读全局应用目录的 skills/。 - 增量同步:全局模板新增的技能/清单,会同步到机构目录(已有机构也能拿到新技能)。 - 返回机构技能根目录。 + 保留此函数仅为兼容旧调用(返回全局技能根目录),复制逻辑已删除。 + 参数 workspace_base / org_id 已不再使用。 """ - org_skills = get_org_skills_dir(workspace_base, org_id) - os.makedirs(org_skills, exist_ok=True) - # 公共技能(global):首次全量复制,之后增量同步(全局模板新增的技能复制过来) - org_global = os.path.join(org_skills, "global") - if not os.path.isdir(org_global): - if os.path.isdir(ALL_DIR): - shutil.copytree(ALL_DIR, org_global) - elif os.path.isdir(ALL_DIR): - for name in os.listdir(ALL_DIR): - src = os.path.join(ALL_DIR, name) - dst = os.path.join(org_global, name) - if not os.path.exists(dst): - if os.path.isdir(src): - shutil.copytree(src, dst) - else: - shutil.copy2(src, dst) - # 技能集清单(packs):首次全量复制,之后增量同步 - org_packs = os.path.join(org_skills, "packs") - if not os.path.isdir(org_packs): - if os.path.isdir(PACKS_DIR): - shutil.copytree(PACKS_DIR, org_packs) - elif os.path.isdir(PACKS_DIR): - for name in os.listdir(PACKS_DIR): - src = os.path.join(PACKS_DIR, name) - dst = os.path.join(org_packs, name) - if not os.path.exists(dst): - if os.path.isdir(src): - shutil.copytree(src, dst) - else: - shutil.copy2(src, dst) - # 产线技能(pipelines/:common + roles):首次全量复制,之后递归增量同步 - org_pipelines = os.path.join(org_skills, "pipelines") - if not os.path.isdir(org_pipelines): - if os.path.isdir(PIPELINES_DIR): - shutil.copytree(PIPELINES_DIR, org_pipelines) - else: - _sync_tree(PIPELINES_DIR, org_pipelines) - return org_skills + return get_skills_base() def list_packs(): diff --git a/wwwroot/api/skill_pack.dspy b/wwwroot/api/skill_pack.dspy index 607b2f0..106c4f9 100644 --- a/wwwroot/api/skill_pack.dspy +++ b/wwwroot/api/skill_pack.dspy @@ -33,16 +33,10 @@ except Exception: if not org_id: return json.dumps({"success": False, "error": "用户未关联机构"}, ensure_ascii=False) -from pipeline_core.skill_pack import list_packs, install_pack, uninstall_pack, installed_packs, ensure_org_skills -from pipeline_service.workspace import get_workspace_base +from pipeline_core.skill_pack import list_packs, install_pack, uninstall_pack, installed_packs, get_skills_base -# 机构技能根目录(多租户隔离:机构工作目录/skills/,不在全局应用目录) -try: - async with DBPools().sqlorContext(dbname) as sor: - ws_base = await get_workspace_base(sor) - base_dir = ensure_org_skills(ws_base, org_id) -except Exception: - base_dir = os.environ.get("PIPELINE_SKILLS_BASE", "skills") +# 机构技能根目录 = 全局 skills/(单一技能树,所有机构共享读,2026-08-21 重构) +base_dir = get_skills_base() if action == 'list': return json.dumps({"success": True, "is_admin": is_admin, "packs": list_packs()}, diff --git a/wwwroot/api/skill_proposals.dspy b/wwwroot/api/skill_proposals.dspy index 50573cb..4c65849 100644 --- a/wwwroot/api/skill_proposals.dspy +++ b/wwwroot/api/skill_proposals.dspy @@ -91,29 +91,15 @@ async with DBPools().sqlorContext(dbname) as sor: return json.dumps({"success": False, "error": "提议缺少 name 或 content"}, ensure_ascii=False) - # 写 SKILL.md 到全局模板 skills/global/{name}/ - from pipeline_core.skill_pack import get_skills_base, ensure_org_skills + # 写 SKILL.md 到全局模板 skills/global/{name}/(单一技能树,所有机构共享读,无需同步副本) + from pipeline_core.skill_pack import get_skills_base skill_dir = os.path.join(get_skills_base(), "global", name) os.makedirs(skill_dir, exist_ok=True) with open(os.path.join(skill_dir, "SKILL.md"), "w", encoding="utf-8") as f: f.write(content) - # 增量同步到所有机构工作目录 - from pipeline_service.workspace import get_workspace_base - ws_base = os.path.normpath(await get_workspace_base(sor)) - synced = 0 - if os.path.isdir(ws_base): - for org_dir in os.listdir(ws_base): - org_path = os.path.join(ws_base, org_dir) - if os.path.isdir(org_path) and not org_dir.startswith('.'): - try: - ensure_org_skills(ws_base, org_dir) - synced += 1 - except Exception: - pass - await sor.U('skill_proposals', {'id': qid, 'status': 'published'}) - return json.dumps({"success": True, "name": name, "synced_orgs": synced}, + return json.dumps({"success": True, "name": name}, ensure_ascii=False) else: