1. memory 工具(add/list/remove): 多租户写入门禁——org_id/user_id强制注入、 scope白名单user/project/pipeline(global/org种子域禁写)、无身份拒写、 remove仅本人条目; 记忆注入改可见性过滤版(修跨机构泄漏) 2. manage_skill(create/patch/write_file/remove_file/delete): skill_live扩展, 只落本租户orgs/users目录, global原版fork-on-write(校验通过才fork,失败零残留), org+user双副本同步改, 产线层拒改, delete只删本租户副本+审计留痕 3. run_command background=true + process工具(poll/log/wait/kill): bg_jobs.py 状态文件化(workspace/.bg/,跨worker可见), 沙箱档位与前台一致(generic强制 strict+无bwrap拒绝), 超时SIGTERM进程组+stale心跳判活 4. delegate_subtask background + subagent工具(list/steer/stop/result): subagents.py 并行上限3/workspace, 深度限制1(子agent禁再委派), 子会话session_isolation=none(不读父历史不写会话表), steer/stop文件传递 每轮tool-loop边界消费, 结果流式落盘(stop即部分结果)
431 lines
22 KiB
Python
431 lines
22 KiB
Python
# -*- 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}
|
||
|
||
|
||
# ═══════════════════════════════════════════════════════════
|
||
# manage_skill_live(2026-09-10,对齐 Hermes skill_manage:
|
||
# create/patch/write_file/remove_file/delete)
|
||
#
|
||
# 隔离铁律(多机构多用户):
|
||
# - 一切写操作只落在本租户目录:机构 orgs/{org}/(org≠0)或 users/{uid}/
|
||
# (org 0 降级 / 个人技能);global、pipelines(产线/角色/项目层)、
|
||
# 他机构 orgs/ 目录物理不可达(realpath 双检 + 白名单);
|
||
# - patch/write_file/remove_file 的目标若是 global 原版技能 →
|
||
# fork-on-write:整目录继承拷贝到本租户目录再改(同名覆盖,只影响本租户);
|
||
# 目标已是本租户既有副本(org 或 user 层)→ 原地改,不迁移不换层;
|
||
# 目标是产线/角色/项目层技能 → 拒绝(平台维护,agent 不可改);
|
||
# - delete 只删本租户副本(org + user 两处都清);global 原版删不到,
|
||
# 删掉副本后原版重新对本租户可见;
|
||
# - 子文件仅限 references/scripts/templates/assets 四目录(与
|
||
# skill_loader.read_linked_file 白名单一致),路径穿越双检。
|
||
# ═══════════════════════════════════════════════════════════
|
||
|
||
_LINKED_DIRS = ('references', 'scripts', 'templates', 'assets')
|
||
|
||
|
||
def _reload_loader():
|
||
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()
|
||
return True
|
||
except Exception as e:
|
||
logger.warning('skill loader reload failed: %s', e)
|
||
return False
|
||
|
||
|
||
def _check_linked_path(rel_path: str) -> str:
|
||
"""子文件路径白名单校验(references/scripts/templates/assets),防穿越。"""
|
||
rel = (rel_path or '').strip().lstrip('/')
|
||
if not rel or '..' in rel.split('/'):
|
||
raise ValueError('非法子文件路径: %r' % rel_path)
|
||
first = rel.split('/')[0]
|
||
if first not in _LINKED_DIRS:
|
||
raise ValueError('子文件只允许放在 %s 目录下(收到 %s)'
|
||
% ('/'.join(_LINKED_DIRS), rel))
|
||
return rel
|
||
|
||
|
||
def _tenant_dirs(org_id: str, user_id: str, name: str, base: str) -> list:
|
||
"""本租户可写的技能目录候选(org 层在前 = 主目标)。名称非法返回 []。"""
|
||
safe_name = _sanitize_id(name)
|
||
if not _NAME_RE.match(safe_name):
|
||
return []
|
||
dirs = []
|
||
org = _sanitize_id(org_id)
|
||
uid = _sanitize_id(user_id)
|
||
if org and org != '0':
|
||
dirs.append(os.path.join(base, 'orgs', org, safe_name))
|
||
if uid:
|
||
dirs.append(os.path.join(base, 'users', uid, safe_name))
|
||
return dirs
|
||
|
||
|
||
def _find_existing_skill(name, org_id, user_id, pipeline_id='', role='', project_id=''):
|
||
"""在租户可见范围内找已存在的技能(返回 Skill 对象或 None)。
|
||
|
||
用 loader 的 get_merged + resolve_skill:本租户副本(org/user)优先,
|
||
产线/角色/项目层次之,global 原版兜底(供 fork-on-write 继承拷贝)。
|
||
"""
|
||
try:
|
||
from pipeline_core.skill_loader import get_skill_loader, resolve_skill
|
||
from pipeline_core.skill_pack import get_skills_base
|
||
loader = get_skill_loader(get_skills_base())
|
||
merged = loader.get_merged(pipeline_id=pipeline_id or '', role=role or '',
|
||
project_id=project_id or '', org_id=org_id or '',
|
||
user_id=user_id or '')
|
||
return resolve_skill(merged, name)
|
||
except Exception as e:
|
||
logger.warning('find existing skill %s failed: %s', name, e)
|
||
return None
|
||
|
||
|
||
def _existing_in_tenant(existing, tenant_dirs: list) -> list:
|
||
"""本租户已存在的技能副本目录列表(loader 生效优先序:user 层在前)。
|
||
|
||
org 层与 user 层可能同时存在同名副本(user 遮蔽 org)——修改必须
|
||
应用到所有副本,否则改了被遮蔽的那份等于没改。
|
||
"""
|
||
dirs = [d for d in tenant_dirs if os.path.isdir(d)]
|
||
src_file = getattr(existing, 'path', '') or ''
|
||
if src_file:
|
||
src = os.path.realpath(os.path.dirname(src_file))
|
||
# loader 实际生效的副本排最前(结果汇报以此为准)
|
||
dirs.sort(key=lambda d: 0 if os.path.realpath(d) == src else 1)
|
||
return dirs
|
||
|
||
|
||
def _assert_inside_tenant(path: str, tenant_dir: str):
|
||
"""realpath 双检:写路径必须在本租户技能目录内(防 symlink/穿越逃逸)。"""
|
||
real = os.path.realpath(path)
|
||
root = os.path.realpath(tenant_dir)
|
||
if not (real == root or real.startswith(root + os.sep)):
|
||
raise ValueError('路径越界(%s 不在租户目录 %s 内),已拒绝' % (path, tenant_dir))
|
||
|
||
|
||
def _atomic_write_text(path: str, text: str):
|
||
tmp = path + '.tmp'
|
||
with open(tmp, 'w', encoding='utf-8') as f:
|
||
f.write(text)
|
||
os.replace(tmp, path)
|
||
|
||
|
||
async def manage_skill_live(action: str, name: str, org_id: str = '', user_id: str = '',
|
||
who: str = '', description: str = '', content: str = '',
|
||
old_string: str = '', new_string: str = '',
|
||
file_path: str = '', file_content: str = '',
|
||
pipeline_id: str = '', role: str = '', project_id: str = ''):
|
||
"""技能增改删(多租户隔离版)。返回 (ok, message, info)。
|
||
|
||
action:
|
||
create 新建/整篇覆盖本租户技能(等价 propose_skill,走 publish_skill_live)
|
||
patch 定向替换片段(old_string 须在目标文件唯一;SKILL.md 或子文件)
|
||
write_file 写子文件(references/scripts/templates/assets)
|
||
remove_file 删子文件
|
||
delete 删除本租户技能副本(global 原版不受影响,删后重新可见)
|
||
"""
|
||
from pipeline_core.skill_pack import get_skills_base
|
||
base = get_skills_base()
|
||
action = (action or '').strip().lower()
|
||
name = (name or '').strip()
|
||
if action not in ('create', 'patch', 'write_file', 'remove_file', 'delete'):
|
||
return False, '未知 action: %r(支持 create/patch/write_file/remove_file/delete)' % action, {}
|
||
if not name:
|
||
return False, '需要技能名 name', {}
|
||
|
||
tenant_dirs = _tenant_dirs(org_id, user_id, name, base)
|
||
if not tenant_dirs:
|
||
# 名称非法,或无 org/user 身份(无人值守 org 0)——与 publish_skill_live 同门禁
|
||
scope, target = resolve_target(org_id, user_id, name)
|
||
if not scope:
|
||
return False, target, {}
|
||
return False, '技能名非法(只允许字母数字._-,且以字母数字开头,≤64字符)', {}
|
||
|
||
try:
|
||
if action == 'create':
|
||
content = (content or '').strip()
|
||
if not content:
|
||
return False, 'create 需要 content(SKILL.md 正文)', {}
|
||
return await publish_skill_live(
|
||
name, description, content, org_id=org_id, user_id=user_id, who=who)
|
||
|
||
existing = _find_existing_skill(name, org_id, user_id,
|
||
pipeline_id, role, project_id)
|
||
if existing is None:
|
||
return False, ('技能 %r 不存在。patch/write_file/remove_file/delete '
|
||
'只能操作已存在的技能;新建请用 action=create' % name), {}
|
||
|
||
if action == 'delete':
|
||
import shutil
|
||
removed = []
|
||
for d in tenant_dirs:
|
||
if os.path.isdir(d):
|
||
_assert_inside_tenant(d, os.path.dirname(d))
|
||
shutil.rmtree(d)
|
||
removed.append(os.path.relpath(d, base))
|
||
if not removed:
|
||
ex_scope = getattr(existing, 'scope', '?')
|
||
if ex_scope == 'global':
|
||
return False, ('%r 是通用(global)技能,本租户没有副本,平台技能不可删除。'
|
||
'如想改掉它的行为,用 action=create 建同名机构副本覆盖' % name), {}
|
||
return False, ('%r 位于 %s 层(平台维护),本租户目录下无副本,不可删除'
|
||
% (name, ex_scope)), {}
|
||
_reload_loader()
|
||
return True, ('OK: 已删除本租户技能副本 %s。若存在同名通用(global)技能,'
|
||
'它重新对本租户可见' % ', '.join(removed)), \
|
||
{'scope': 'tenant', 'path': removed[0], 'removed': removed}
|
||
|
||
# patch / write_file / remove_file:定位操作目录。
|
||
# 本租户可能同时存在 org 层与 user 层副本(user 遮蔽 org)——
|
||
# target_dirs 收集全部需要一致修改的目录,防「改了被遮蔽的那份」。
|
||
#
|
||
# fork 延迟到校验通过后(2026-09-10 review 修复):先 fork 后校验
|
||
# 会在失败路径留下副本残留(如 old_string 未找到仍拷贝了目录),
|
||
# 残留副本遮蔽 global 原版、还会让 delete 误判「有副本可删」。
|
||
# read_dirs = 校验读取目录(副本 or global 原版);
|
||
# ensure_write_dirs() = 校验全过后才 fork 并返回落盘目录。
|
||
target_dirs = _existing_in_tenant(existing, tenant_dirs)
|
||
forked = False
|
||
fork_note = ''
|
||
need_fork = False
|
||
if target_dirs:
|
||
read_dirs = target_dirs
|
||
elif getattr(existing, 'scope', '') == 'global':
|
||
need_fork = True
|
||
read_dirs = [os.path.realpath(os.path.dirname(getattr(existing, 'path', '')))]
|
||
fork_note = ('(原版属 global 层,已继承副本到本租户目录后修改,'
|
||
'同名覆盖仅本租户生效、其他机构不受影响)')
|
||
else:
|
||
return False, ('%r 位于 %s 层(平台/产线维护),agent 不可修改。'
|
||
'可修改的范围:本机构(org)/本人(user)副本与通用(global)技能'
|
||
'(global 修改走继承副本)' % (name, existing.scope)), {}
|
||
|
||
def ensure_write_dirs():
|
||
"""校验通过后落盘目录就位:需要 fork 时才拷贝 global 原版。"""
|
||
nonlocal target_dirs, forked
|
||
if not need_fork:
|
||
return target_dirs
|
||
import shutil
|
||
skill_dir = tenant_dirs[0]
|
||
_assert_inside_tenant(skill_dir, os.path.dirname(skill_dir))
|
||
os.makedirs(os.path.dirname(skill_dir), exist_ok=True)
|
||
if not os.path.exists(skill_dir):
|
||
shutil.copytree(read_dirs[0], skill_dir)
|
||
target_dirs = [skill_dir]
|
||
forked = True
|
||
return target_dirs
|
||
|
||
for d in read_dirs + tenant_dirs:
|
||
_assert_inside_tenant(d, os.path.dirname(d) if d in tenant_dirs else d)
|
||
|
||
if action == 'patch':
|
||
rel = _check_linked_path(file_path) if file_path else 'SKILL.md'
|
||
if old_string == '':
|
||
return False, 'patch 需要 old_string(要替换的原文片段)', {}
|
||
# 第一阶段:只读校验(任一读取目录找不到/不唯一即整体拒绝,零落盘)
|
||
for d in read_dirs:
|
||
fpath = os.path.join(d, rel)
|
||
if not os.path.isfile(fpath):
|
||
return False, '目标文件不存在: %s(%s,技能内文件用 load_skill 查看)' % (rel, os.path.relpath(d, base)), {}
|
||
with open(fpath, 'r', encoding='utf-8') as f:
|
||
text = f.read()
|
||
n = text.count(old_string)
|
||
if n == 0:
|
||
return False, 'patch 失败:old_string 在 %s 中未找到(先用 load_skill 核对原文,注意空白/缩进完全一致)' % rel, {}
|
||
if n > 1:
|
||
return False, 'patch 失败:old_string 在 %s 中出现 %d 次(须唯一,请带更多上下文)' % (rel, n), {}
|
||
# 第二阶段:校验全过,落盘目录就位(此时才 fork),同构路径逐个写
|
||
write_dirs = ensure_write_dirs()
|
||
texts = []
|
||
for d in write_dirs:
|
||
fpath = os.path.join(d, rel)
|
||
_assert_inside_tenant(fpath, d)
|
||
if not os.path.isfile(fpath):
|
||
return False, '目标文件不存在: %s(%s)' % (rel, os.path.relpath(d, base)), {}
|
||
with open(fpath, 'r', encoding='utf-8') as f:
|
||
text = f.read()
|
||
texts.append((fpath, text.replace(old_string, new_string or '', 1)))
|
||
for fpath, text in texts:
|
||
_atomic_write_text(fpath, text)
|
||
_reload_loader()
|
||
n_dirs = len(texts)
|
||
multi = '(org+user 两层副本已同步修改)' if n_dirs > 1 else ''
|
||
return True, 'OK: 已修改 %r 的 %s%s%s,本租户内立即生效' % (name, rel, multi, fork_note), \
|
||
{'scope': 'tenant', 'path': os.path.relpath(texts[0][0], base),
|
||
'paths': [os.path.relpath(t[0], base) for t in texts], 'forked': forked}
|
||
|
||
if action == 'write_file':
|
||
rel = _check_linked_path(file_path)
|
||
if not file_content:
|
||
return False, 'write_file 需要 file_content', {}
|
||
written = []
|
||
for d in ensure_write_dirs():
|
||
fpath = os.path.join(d, rel)
|
||
_assert_inside_tenant(fpath, d)
|
||
os.makedirs(os.path.dirname(fpath), exist_ok=True)
|
||
_atomic_write_text(fpath, file_content)
|
||
written.append(os.path.relpath(fpath, base))
|
||
_reload_loader()
|
||
multi = '(org+user 两层副本已同步写入)' if len(written) > 1 else ''
|
||
return True, 'OK: 已写入 %r 的子文件 %s%s%s,本租户内立即生效' % (name, rel, multi, fork_note), \
|
||
{'scope': 'tenant', 'path': written[0], 'paths': written, 'forked': forked}
|
||
|
||
if action == 'remove_file':
|
||
rel = _check_linked_path(file_path)
|
||
# 先校验读取目录里至少一处存在该文件(否则不 fork、直接拒绝)
|
||
if not any(os.path.isfile(os.path.join(d, rel)) for d in read_dirs):
|
||
return False, '子文件不存在: %s' % rel, {}
|
||
removed_files = []
|
||
for d in ensure_write_dirs():
|
||
fpath = os.path.join(d, rel)
|
||
_assert_inside_tenant(fpath, d)
|
||
if os.path.isfile(fpath):
|
||
os.remove(fpath)
|
||
removed_files.append(os.path.relpath(fpath, base))
|
||
if not removed_files:
|
||
return False, '子文件不存在: %s' % rel, {}
|
||
_reload_loader()
|
||
return True, 'OK: 已删除 %r 的子文件 %s(%d 处副本)' % (name, rel, len(removed_files)), \
|
||
{'scope': 'tenant', 'path': removed_files[0], 'paths': removed_files}
|
||
|
||
except ValueError as e:
|
||
return False, str(e), {}
|
||
except Exception as e:
|
||
logger.exception('manage_skill_live %s %s failed', action, name)
|
||
return False, 'ERROR: %s' % str(e)[:300], {}
|
||
return False, '未实现的动作分支', {}
|