diff --git a/pipeline_service/agent_loop_v2.py b/pipeline_service/agent_loop_v2.py index d4a112f..4673214 100644 --- a/pipeline_service/agent_loop_v2.py +++ b/pipeline_service/agent_loop_v2.py @@ -139,10 +139,12 @@ class AgentExecutor: session_id: str = "", # 会话内唯一标识(web 多 tab 独立会话;历史隔离键) default_pipeline_id: str = "", # 无当前项目时的默认产线(入口指定,如 bidding_general) delegate_depth: int = 0, # 委派深度(0=顶层会话;子agent=1,禁止再委派) + pii_session_key: str = "", # PII ephemeral 缓存键(gateway _session_key;空=无 PII 层) ): self.config = config self.project_id = project_id self.user_id = user_id + self.pii_session_key = pii_session_key self.org_id = "" # loaded from project context self.pipeline_id = default_pipeline_id or "" # 默认产线(入口指定);有项目时被项目 pipeline_id 覆盖 self.role = role # 产线内角色(技能/工具/prompt 按角色加载) @@ -1609,6 +1611,20 @@ class AgentExecutor: except Exception as e: logger.warning(f"secret prepare_command failed (run without secrets): {e}") + # PII 执行边界(2026-09-17):@@pii:占位符 → $PIPELINE_PII_*,真值经 env 注入, + # 与凭据同构;缓存缺失时引用保留(命令取空,不静默改语义——与凭据不同:PII + # 缺失是过期常态,拒绝执行会误伤,取空+输出可见更安全)。 + pii_env = None + if self.pii_session_key and "@@pii:" in cmd: + try: + from . import pii_guard + cmd, pii_env = await pii_guard.prepare_command( + cmd, self.pii_session_key) + except Exception as e: + logger.warning(f"pii prepare_command failed (run without pii env): {e}") + if pii_env: + secret_env = {**(secret_env or {}), **pii_env} + # 后台执行(2026-09-10,对齐 Hermes terminal background): # 状态文件化到 workspace/.bg/,跨 worker 进程可 poll; # 沙箱档位与前台完全一致(generic 强制 strict,无 bwrap 上面已拒)。 @@ -2457,6 +2473,21 @@ class AgentExecutor: if self.config.session_isolation == "none": return + # PII 落库掩码(2026-09-17):会话表是长期留存+session_search 可捞,占位符 + # 虽非明文但缓存 24h 内可还原 → 落库一律换掩码形态(138****5678), + # 缓存过期后掩码仍是唯一可见形态(合规:留存数据不含可还原 PII)。 + if self.pii_session_key and "@@pii:" in (user_input + reply): + try: + from . import pii_guard + if "@@pii:" in user_input: + user_input = await pii_guard.restore_outbound( + user_input, self.pii_session_key, masked=True) + if "@@pii:" in reply: + reply = await pii_guard.restore_outbound( + reply, self.pii_session_key, masked=True) + except Exception as e: + logger.warning(f"pii mask on save failed: {e}") + try: from sqlor.dbpools import DBPools from appPublic.uniqueID import getID diff --git a/pipeline_service/gateway.py b/pipeline_service/gateway.py index 3796a6c..0a2fbf9 100644 --- a/pipeline_service/gateway.py +++ b/pipeline_service/gateway.py @@ -172,6 +172,22 @@ class Gateway: for note in secret_notes: yield json.dumps({"type": "progress", "message": note + "\n"}, ensure_ascii=False) + "\n" + # 0.6 PII 入站门禁(2026-09-17 用户定案:合规,敏感信息不传第三方): + # 手机号/身份证/银行卡/邮箱等个人敏感信息 → AES 密文入 Redis(24h ephemeral, + # **不进机构密码表**)→ 原文替换 @@pii:占位符。与 0.5 凭据门禁同构不同库。 + # 放在模型/落库之前:净化后的 content 才是上行与落库版本。失败降级原文。 + pii_key = self._session_key(channel, user_id, session_id) + try: + from . import pii_guard + _pr = await pii_guard.mask_inbound(content, pii_key) + content = _pr.get("text") or content + pii_notes = list(_pr.get("notes") or []) + except Exception as e: + logger.warning(f"pii inbound gate failed (continue with raw content): {e}") + pii_notes = [] + for note in pii_notes: + yield json.dumps({"type": "progress", "message": note + "\n"}, ensure_ascii=False) + "\n" + # 1. 解析项目上下文(纯通用模式跳过,不挂任何产线插件) if generic: ctx = {"pid": "", "pipeline_id": "", "name": "", "project_model_name": ""} @@ -222,7 +238,7 @@ class Gateway: executor = AgentExecutor( config=config, project_id=ctx["pid"], user_id=user_id, role=role, base_url=base_url, generic=generic, session=sess, session_id=session_id, - default_pipeline_id=pipeline_id) + default_pipeline_id=pipeline_id, pii_session_key=pii_key) # 3.5 原生视觉(2026-09-10):上传图片构造多模态 parts。 # 上限/超限策略在 build_image_parts(8MB/张、4张/条,超限显式告知不静默丢)。 @@ -239,8 +255,21 @@ class Gateway: content = content + "\n\n[系统说明·本次上传图片] " + ";".join(image_notes) # 4. 转发标准化事件流 + # PII 用户侧还原(2026-09-17):模型/落库只见占位符,回用户的 reply/progress + # 还原原文(用户自己贴的信息要能看见);tool_call/tool_result **不还原**—— + # 它们会进模型上下文与 llm_call_trace 上行。还原失败降级为占位符原样。 + from . import pii_guard as _pii async for chunk in executor.run(content, history=history, image_parts=image_parts or None): + if pii_notes and chunk.strip(): + try: + _ev = json.loads(chunk.strip()) + if _ev.get("type") in ("reply", "progress") and "@@pii:" in str(_ev.get("message", "")): + _ev["message"] = await _pii.restore_outbound( + str(_ev["message"]), pii_key) + chunk = json.dumps(_ev, ensure_ascii=False) + "\n" + except Exception: + pass yield chunk async def _secret_inbound_gate(self, user_id: str, content: str): diff --git a/pipeline_service/pii_guard.py b/pipeline_service/pii_guard.py new file mode 100644 index 0000000..2f7f7e2 --- /dev/null +++ b/pipeline_service/pii_guard.py @@ -0,0 +1,258 @@ +"""pii_guard.py - 个人敏感信息(PII)打码守卫(2026-09-17 用户定案)。 + +与凭据(secret_vault)**同构不同库**: + 同构 —— 入站侦测 → 原文替换为占位符 `@@pii:KIND_FP@@` → 模型只见占位符 → + 执行边界以子进程 env 注入真值(`$PIPELINE_PII_*`)→ 用户侧回复还原原文。 + 不同库 —— PII **不进机构密码表**(pipeline_user_secrets):AES 密文存 Redis + (db4,TTL 24h,会话级 ephemeral),过期即焚,不做跨会话资产。 + +语义差异(写死,不靠调用方自觉): + 凭据 = "只用不看":还原只发生在执行 env,用户回复里也只见占位符/掩码。 + PII = "模型不看、用户看":落库/上行/模型上下文全是占位符;回用户侧的 + reply/progress 事件还原原文(用户自己贴的信息要能看见)。 + +检测层只做**确定性格式**(手机号/身份证+MOD 11-2 校验/银行卡+Luhn/邮箱), +不上模型:门禁在消息进模型之前,用 LLM 判断 = 先把原文发给 LLM,自相矛盾。 +人名/地址等无格式实体留二期(本地 NER,只标不杀,事后审计)。 + +合规边界(一期):覆盖 gateway v2 会话(入站门禁 + 用户侧还原 + run_command +env 注入 + 会话落库占位符)。v1 角色 agent 交付件出库打码 = 二期,README 注明。 +""" +import hashlib +import json +import re +import time +from typing import Dict, List, Optional, Tuple + +from appPublic.log import info, warning + +#: 占位符语法(与 @@sec: 同族不同前缀,互不干扰) +PII_PLACEHOLDER_PREFIX = "@@pii:" +PII_PLACEHOLDER_SUFFIX = "@@" +PII_ENV_PREFIX = "PIPELINE_PII_" +_PLACEHOLDER_RE = re.compile(r"@@pii:([A-Z]+_[0-9A-F]{8})@@") + +#: Redis 存储 TTL(24h,会话级 ephemeral;PII 不做长期资产) +PII_TTL_SECONDS = 24 * 60 * 60 +_REDIS_DB = 4 # session 用 db3,PII 独立 db4,互不污染 + +# ══════════════════════ 确定性检测器 ══════════════════════ + +_ID_WEIGHTS = [7, 9, 10, 5, 8, 4, 2, 1, 6, 3, 7, 9, 10, 5, 8, 4, 2] +_ID_CHECKCODES = "10X98765432" + +_RE_IDCARD = re.compile(r"(? bool: + """GB 11643 MOD 11-2 校验。""" + v = v.upper() + if len(v) != 18: + return False + try: + s = sum(int(v[i]) * _ID_WEIGHTS[i] for i in range(17)) + except ValueError: + return False + return _ID_CHECKCODES[s % 11] == v[17] + + +def luhn_ok(v: str) -> bool: + digits = [int(c) for c in v] + odd = digits[-1::-2] + even = digits[-2::-2] + tot = sum(odd) + for d in even: + d *= 2 + tot += d - 9 if d > 9 else d + return tot % 10 == 0 + + +def _mask(kind: str, v: str) -> str: + if kind == "PHONE": + return v[:3] + "****" + v[7:] + if kind == "IDCARD": + return v[:3] + "*" * (len(v) - 7) + v[-4:] + if kind == "BANKCARD": + return v[:4] + "*" * (len(v) - 8) + v[-4:] + if kind == "EMAIL": + local, _, domain = v.partition("@") + keep = local[:1] if local else "" + return keep + "***@" + domain + return "*" * len(v) + + +def detect_pii(text: str) -> List[Dict]: + """扫描文本返回 PII 候选(span 去重,优先级 身份证>银行卡>手机>邮箱)。 + + 返回项:{value, kind, masked, span}。只认确定性格式,零模型依赖。 + """ + if not text or not isinstance(text, str): + return [] + hits: List[Dict] = [] + + def _add(kind: str, value: str, span: Tuple[int, int]): + hits.append({"kind": kind, "value": value, + "masked": _mask(kind, value), "span": span}) + + for m in _RE_IDCARD.finditer(text): + v = m.group(1) + if idcard_checksum_ok(v): + _add("IDCARD", v.upper(), m.span(1)) + for m in _RE_BANKCARD.finditer(text): + v = m.group(1) + if luhn_ok(v): + _add("BANKCARD", v, m.span(1)) + for m in _RE_PHONE.finditer(text): + _add("PHONE", m.group(1), m.span(1)) + for m in _RE_EMAIL.finditer(text): + v = m.group(1) + domain = v.partition("@")[2].lower() + if domain in _EMAIL_JUNK_DOMAINS or v.count("@") != 1: + continue + _add("EMAIL", v, m.span(1)) + + # span 去重:按优先级排序后剔除与已选区间重叠者(身份证 18 位优先于银行卡/手机子串) + prio = {"IDCARD": 0, "BANKCARD": 1, "PHONE": 2, "EMAIL": 3} + hits.sort(key=lambda h: (prio[h["kind"]], h["span"][0])) + chosen: List[Dict] = [] + for h in hits: + a, b = h["span"] + if any(not (b <= ca or a >= cb) for ca, cb in + (c["span"] for c in chosen)): + continue + chosen.append(h) + chosen.sort(key=lambda h: h["span"][0]) + return chosen + + +def _ph_id(kind: str, value: str) -> str: + return kind + "_" + hashlib.sha1(value.encode("utf-8")).hexdigest()[:8].upper() + + +def placeholder_of(kind: str, value: str) -> str: + return PII_PLACEHOLDER_PREFIX + _ph_id(kind, value) + PII_PLACEHOLDER_SUFFIX + + +# ══════════════════════ Redis ephemeral 存储 ══════════════════════ + +_redis_cli = None + + +def _get_redis(): + global _redis_cli + if _redis_cli is not None: + return _redis_cli + import redis.asyncio as aioredis + from appPublic.jsonConfig import getConfig + url = "" + try: + url = (getConfig().website.session_redis.url or "") + except Exception: + url = "" + if not url: + url = "redis://127.0.0.1:6379/3" + # 换 db 号:.../3 → .../4(PII 独立库) + base, _, db = url.rpartition("/") + url = (base + "/" + str(_REDIS_DB)) if base else url + _redis_cli = aioredis.Redis.from_url(url) + return _redis_cli + + +def _rkey(session_key: str, ph_id: str) -> str: + return f"pii:{session_key}:{ph_id}" + + +async def store_pii(session_key: str, kind: str, value: str, masked: str) -> str: + """AES 密文写 Redis(TTL 24h)。返回占位符。绝不写 DB 表。""" + from .secret_vault import encrypt_secret + r = _get_redis() + ph = placeholder_of(kind, value) + payload = json.dumps({"enc": encrypt_secret(value), "masked": masked, + "kind": kind, "created": int(time.time())}, + ensure_ascii=False) + await r.set(_rkey(session_key, _ph_id(kind, value)), payload, ex=PII_TTL_SECONDS) + return ph + + +async def _load(session_key: str, ph_id: str) -> Optional[Dict]: + r = _get_redis() + raw = await r.get(_rkey(session_key, ph_id)) + if not raw: + return None + try: + if isinstance(raw, bytes): + raw = raw.decode("utf-8") + return json.loads(raw) + except Exception: + return None + + +async def mask_inbound(text: str, session_key: str) -> Dict: + """入站门禁:侦测 → 密文入 Redis → 原文替换占位符。 + + 返回 {text, notes, count}。任何异常由调用方兜底降级为原文。 + """ + cands = detect_pii(text) + if not cands: + return {"text": text, "notes": [], "count": 0} + out = text + kinds = [] + for c in reversed(cands): # 从后往前替换,span 不漂移 + ph = await store_pii(session_key, c["kind"], c["value"], c["masked"]) + out = out[:c["span"][0]] + ph + out[c["span"][1]:] + kinds.append(c["kind"]) + note = ("检测到个人敏感信息(" + "、".join(sorted(set(kinds))) + + ")→ 已替换为占位符,原文加密缓存 24 小时(不进模型上下文、不明文落库)," + "回复中会自动还原给你。") + return {"text": out, "notes": [note], "count": len(cands)} + + +async def restore_outbound(text: str, session_key: str, masked: bool = False) -> str: + """用户侧还原:占位符 → 原文(masked=False)或掩码形态(masked=True,日志/导出用)。 + + 缓存过期/缺失 → 占位符原样保留(不猜不编,降级可见)。 + """ + if not text or PII_PLACEHOLDER_PREFIX not in text: + return text + + async def _sub(m): + rec = await _load(session_key, m.group(1)) + if not rec: + return m.group(0) + if masked: + return rec.get("masked", m.group(0)) + from .secret_vault import decrypt_secret + return decrypt_secret(rec.get("enc", "")) + + # 逐个替换(async re 不存在,手动遍历) + out = text + for m in reversed(list(_PLACEHOLDER_RE.finditer(text))): + repl = await _sub(m) + out = out[:m.start()] + repl + out[m.end():] + return out + + +async def prepare_command(command: str, session_key: str) -> Tuple[str, Dict[str, str]]: + """执行边界:占位符 → `$PIPELINE_PII_*` env 引用 + env 真值字典。 + + 与凭据 prepare_command 同构:命令字符串里**只有 env 引用**(落库/上行安全), + 真值只经子进程 env 进入。缺失缓存 → 引用保留、env 不注入(命令自然取空)。 + """ + if not command or PII_PLACEHOLDER_PREFIX not in command: + return command, {} + env: Dict[str, str] = {} + out = command + for m in reversed(list(_PLACEHOLDER_RE.finditer(command))): + rec = await _load(session_key, m.group(1)) + var = PII_ENV_PREFIX + m.group(1) + if rec: + from .secret_vault import decrypt_secret + env[var] = decrypt_secret(rec.get("enc", "")) + out = out[:m.start()] + "$" + var + out[m.end():] + return out, env