refactor(skill): 单一技能树——七级 scope 全在全局 skills/ 下,去掉 ensure_org_skills 复制副本

- skill_loader: 加 project_common scope(projects/ 项目通用,不分 project_id) + project 私有改挂 orgs/{org_id}/{project_id}/,SCOPE_PRIORITY 七级
- skill_pack: ensure_org_skills 去复制逻辑改薄封装返回 get_skills_base(),删 _sync_tree/get_org_skills_dir dead code
- skill_pack.dspy/skill_proposals.dspy: base_dir 改全局 skills/,删同步所有机构循环
This commit is contained in:
yumoqing 2026-08-21 19:18:45 +08:00
parent 70b9df480c
commit b9998ed00d
4 changed files with 61 additions and 128 deletions

View File

@ -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()),

View File

@ -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
返回 <workspace_base>/<org_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():

View File

@ -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()},

View File

@ -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: