ymq 1c04efe1be feat(agent): 会话agent四项Hermes能力对齐(记忆写入/技能管理/后台进程/并行委派)
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即部分结果)
2026-09-10 17:09:45 +08:00

431 lines
22 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

# -*- 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 带合法 frontmattername/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_live2026-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 需要 contentSKILL.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, '未实现的动作分支', {}