lili c6d6d39a2a
Some checks failed
contract-gates / contract-gates (push) Has been cancelled
docs-gate / docs-gate (push) Has been cancelled
feat(cfg): W-CFG-KB K3 知识包治理接线——knowledge 配置块+激活发布 Nacos+worker 订阅热物化(复用阶段二三·默认关)
把切片一知识件外置接进受治理配置集,一条「改知识件→定版→激活→worker 物化生效」链,复用配置控制面阶段二/三治理机器零重造:

game-cloud:knowledge 作 content_json 第四内容块(仿 readiness),与 routeA/routeB/readiness 互斥、sanity(packageId/
  activeVersion 非空);走向相反=真下发(readiness 进程内空载荷,knowledge 跨进程),激活编成一条路 B 载荷
  {knowledge.packageId,knowledge.activeVersion}、经 KNOWLEDGE 档独立 dataId gen-hot-params-knowledge publish 给 worker;
  pathA=null→dispatch/reconcile 只走路 B、激活编排一行未改。加 AigcConfigTierEnum.KNOWLEDGE+KnowledgeConfigKeys+NO_ROUTE 文案。

worker:kb_store 加进程级激活版本源槽+attach/detach(对称 genconfig),resolve_active_version 取值链改为治理源(Nacos)>
  env 钉版>最新 committed;新 kb_nacos.py(KbActiveVersionSource 订阅 dataId、按 packageId 匹配、激活变更即重物化)+
  setup_knowledge_activation() 启动物化+(Nacos 启用时)attach,cheap_service_app 紧接 genconfig 热源挂它。复用 genconfig_nacos 启用旗+代理旁路。

默认关字节不变(TIER2_KB_ROOT 未配→不物化不订阅、read_file 回落仓根;无源→resolve 回落 env/最新);缺文件显式失败沿用
(materialize_active 照抛保留 kb_root、kb_nacos 只在 best-effort 边界接住不咬生成);全 best-effort。admin 本轮做前端契约点
sanity.ts,富编辑视图=紧接下一片(与 readiness 一致)。

主控终审:git 完整(HEAD 未动·无破坏操作)+11 文件精确+读码坐实互斥/routeB载荷/取值链/best-effort/默认关/缺文件失败没削弱
+亲跑 Java codec 34+activation 21(BUILD SUCCESS)+Python test_kb_externalize 21 全绿。真 Nacos 激活端到端排窗口。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-05 22:50:10 -07:00

862 lines
49 KiB
Python

"""worker/kb_store.py —— 便宜档知识件外置(W-CFG-KB 切片一)· 知识包存储 + 物化。
【这个模块解决什么】
便宜档生成 agent 靠 read_file 按硬编码相对路径自取「知识件」——设计 skill、黄金脚手架、few-shot 起点范例
(见 cheap-worker/cheap_roles.py 的按品类指路清单)。这些文件现在躺在 worker 的 git 检出里,想改一条设计
范式只能改仓、把新检出重部署到 worker 才生效。本模块把知识件外置成【知识包】:一批文件打成版本化的包,
文件全文存进 MinIO、manifest(文件清单 + 各文件内容指纹 + 版本 id)落 MySQL,再按「当前激活版本」物化到
git 检出之外的知识根(TIER2_KB_ROOT);cheap_run.read_file 先查知识根、命中即读。于是改一条范式 = 存一版
新知识包 + 激活,worker 物化即生效、不重部署。
【复用 ASSET-SRC 存储范式,但另立表 / 桶(不合表)】
存储层刻意复用 tier2 源工程落库(worker/store.py 的 BackendStore)那套已被真库坐实的范式:
· manifest-first 半落库时序:MySQL 落 status=pending 行 → MinIO put 文件全文 → 翻 committed;fetch 只
返 committed。中途崩溃只留一条【可见 pending 行】(不是无 manifest 指向的隐形对象),不留隐形孤儿;
· 版本寻址:改一版知识件 → 派生新 versionId,旧版本仍可按 versionId 取回;
· 内容指纹幂等:整包内容未变 → 不重复落一份新版本;
· 缺源显式失败:manifest 列了某文件但 MinIO 取不到 → fetch 标 missingContent,物化据此【显式失败】、
绝不物化半拉子包(与 ASSET-SRC「缺源不静默」同口径)。
但【不共用同一张表、同一个桶】:源工程是「生成产物」、知识件是「生成输入」,两件正交、生命周期不同,故另立
知识包专用表(cheap_knowledge_package_version)+ 专用桶(minio.kb_bucket,默认 cheap-kb)。只有纯内容
寻址工具(_canonical_content_hash / derive_version_id / _is_safe_rel_path)直接复用 store.py 的单一实现,
不另抄一份——这几件是内容寻址的单一事实源。
【物化落点必须在 git 检出之外(deploy-safe)】
dev 部署做 `git reset --hard`,若把知识包物化进检出树会被抹掉。故物化落点 = 一个可配置的专用目录
TIER2_KB_ROOT(env,默认落检出外路径),与 cheap_run.read_file 的知识根 shadow 读的是同一个 env、同一份
物化产物(写侧在此、读侧在 cheap_run)。
【当前激活版本(切片一简版;K3 前留清晰接缝)】
切片一的「当前激活版本」= env TIER2_KB_ACTIVE_VERSION 显式钉版,或回落「最新 committed 版本」。yudao 配置集
治理(改→版本→激活推送→回滚)+ admin 呈现是后续切片(K3):resolve_active_version 就是那条线的接缝点——
K3 把它换成「读 yudao 当前激活的知识包版本」,并由激活推送触发 materialize_active。本模块只把存储 + 物化 +
「取当前激活版」这条闭环补齐,治理面留缝、不做。
设计源:docs/plans/2026-07-05-002-feat-cfg-ext-配置外置-plan.md(§1.2/§1.3 知识件现状、§4 K1-K2 路乙)。
"""
from __future__ import annotations
import hashlib
import json
import os
import shutil
from abc import ABC, abstractmethod
from pathlib import Path
from typing import Any, Optional
# 纯内容寻址工具复用 store.py 的单一实现(整包指纹 / 版本 id 派生 / 相对路径安全校验)——不另抄一份,
# 保证知识包与源工程用同一套「同内容 → 同指纹」「改源 → 新 versionId」语义。store.py 顶层只 import 标准库,
# 任何装了 Python 的环境都能 import(6c6g / cheap venv 皆可),不引入 pymysql / minio 依赖。
from worker.store import ( # noqa: E402
_canonical_content_hash,
_is_safe_rel_path,
derive_version_id,
)
# ── 落库后端常量(表名 / 桶配置键 / 知识包默认 id;集中此处便于 doc↔code 对账)─────────────────
# 这三个是 doc↔code 契约点:.sql 留档 / 凭据档 / 后续 K3 admin 取回都按它们对齐。刻意与 ASSET-SRC 的
# tier2_source_project_version / tier2-src 区别开——两件正交、不合表。
_KB_MYSQL_TABLE = "cheap_knowledge_package_version"
_KB_BUCKET_DEFAULT = "cheap-kb" # 知识包专用 MinIO 桶(infra.yaml minio.kb_bucket 缺省时的回落)
DEFAULT_PACKAGE_ID = "cheap-kb" # 一套便宜档知识语料 = 一个包 id,其下挂多版本
# 物化落点 env 名(写侧在本模块、读侧在 cheap_run._kb_root;两处必须用同一个 env 名,是一条稳定契约)。
KB_ROOT_ENV = "TIER2_KB_ROOT"
# 当前激活版本 env 名(切片一钉版口子;K3 用 yudao 激活版取代它)。
KB_ACTIVE_VERSION_ENV = "TIER2_KB_ACTIVE_VERSION"
# manifest schema 版本(将来演进 manifest 结构时按它区分)。
_KB_MANIFEST_SCHEMA = "kb-package/1"
# ── MySQL 建表 DDL(代码内首次 save 执行 CREATE TABLE IF NOT EXISTS;另有同源 .sql 见 config/schema/)──
# 幂等键 = (package_id, content_hash):同包同内容(改了知识件 content_hash 才变)只落一行,命中不重写。
# version_id 另加唯一索引(同包版本号唯一,取回按 (package_id, version_id) 寻址)。
# 本表是【新建表】,出生即带 status 列(不像 ASSET-SRC 表有历史无 status 列旧表需 ALTER 补列),故无 ALTER。
# in-code DDL 与 config/schema/cheap_knowledge_package_version.sql 双源,由
# cheap-worker/tests/test_kb_externalize.py 规范化对账钉成机器门(改一处漏改另一处即红)。
_DDL_KB_PACKAGE_VERSION = f"""
CREATE TABLE IF NOT EXISTS {_KB_MYSQL_TABLE} (
id BIGINT NOT NULL AUTO_INCREMENT COMMENT '自增主键',
package_id VARCHAR(128) NOT NULL COMMENT '知识包稳定标识(一套便宜档知识语料一个 id;其下挂多版本)',
version_id VARCHAR(96) NOT NULL COMMENT '版本号 vXXXX-hash12(改一版知识件生成新版本)',
content_hash CHAR(64) NOT NULL COMMENT '整包内容指纹 sha256(= 各知识件 path+content 规范化拼接;幂等键)',
manifest_json LONGTEXT NOT NULL COMMENT '知识包 manifest(文件清单 path+各文件 sha256+bytes;不含 content,content 落 MinIO)',
addressing VARCHAR(512) NULL COMMENT '落库寻址摘要(store/sourceUrl;冗余出 manifest 便于检索)',
created_at DATETIME NOT NULL COMMENT '落库时刻(now_ts 转 UTC datetime)',
status VARCHAR(16) NOT NULL DEFAULT 'committed' COMMENT 'manifest-first 落库状态:pending=已占位未提交(半落库可见残行)/committed=文件已全落可物化;fetch 只返 committed',
PRIMARY KEY (id),
UNIQUE KEY uk_pkg_content (package_id, content_hash) COMMENT '幂等键:同包同内容不重复落',
UNIQUE KEY uk_pkg_version (package_id, version_id) COMMENT '版本寻址键:同包版本号唯一',
KEY idx_pkg_created (package_id, created_at) COMMENT 'fetch 取最新 committed 版本走它'
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='便宜档知识件外置知识包版本表(W-CFG-KB 切片一)';
"""
class KnowledgeMaterializeError(RuntimeError):
"""知识包物化的【显式失败】——激活版取不到 / manifest 列了但缺文件 / 文件指纹不匹配。
这是本模块唯一往外抛的异常:物化是安全语义要求,缺文件绝不静默物化半拉子包。抛出时 kb_root 未被触碰
(物化在 staging 里做完并校验才原子换入),故上一份已物化包原样保留 = 天然回落,agent 读到的仍是完整旧包。
"""
def _file_content_hash(content: str) -> str:
"""单文件内容指纹 sha256(hex)。用于 manifest 的逐文件 hash 与物化后逐文件校验。"""
return hashlib.sha256(content.encode("utf-8")).hexdigest()
def _build_manifest(file_list: list[dict], content_hash: str, version_id: str,
package_id: str, note: Optional[str]) -> dict:
"""据源文件全文清单组 manifest(文件清单 + 各文件 sha256 + bytes;不含 content)。
files 按 path 排序落定(与 _canonical_content_hash 的排序口径一致,保证 manifest 确定性)。
content 不进 manifest(它落 MinIO / LocalFs 的 files/),manifest 只记「有哪些文件、各自多大、指纹是啥」。
"""
files = []
for f in sorted(file_list or [], key=lambda x: (x or {}).get("path", "")):
rel = (f or {}).get("path") or ""
content = (f or {}).get("content")
if not rel or not isinstance(content, str):
continue
files.append({
"path": rel,
"sha256": _file_content_hash(content),
"bytes": len(content.encode("utf-8")),
})
manifest = {
"schemaVersion": _KB_MANIFEST_SCHEMA,
"packageId": package_id,
"versionId": version_id,
"contentHash": content_hash,
"files": files,
}
if note:
manifest["note"] = note
return manifest
class KnowledgePackageStore(ABC):
"""知识包落库取回接口(便宜档知识件外置的存储契约 · 后端无关)。
三个动作钉死知识包的存 / 取 / 定位:
- save(file_list, now_ts) → {packageId, versionId, contentHash}:把一批知识件(文件全文清单)持久化成
一个版本化知识包。改一版知识件再 save 得【新 versionId】(旧版本仍在,按版本取回)。整包内容未变则幂等
命中、不重复落。
- fetch(package_id, version_id) → manifest + 各文件 content:按包 id(+ 可选版本)取回。缺文件不静默——
标 missingContent,由物化据此显式失败。
- latest_committed_version(package_id) → 最新 committed 版本 id:切片一「当前激活版」的默认来源。
【实现约定】与 store.py 的 SourceProjectStore 同纪律:幂等(整包内容未变不重落)、失败可追溯(核心/错误
路径都打日志)。故意不继承 SourceProjectStore——知识包不是源工程,save/fetch 的入参形状不同,强行套一套
接口是拼接设计;两者只共享【范式】(manifest-first + 版本寻址 + 内容指纹),不共享类型。
"""
@abstractmethod
def save(self, file_list: list[dict], *, now_ts: float,
package_id: str = DEFAULT_PACKAGE_ID, note: Optional[str] = None) -> dict:
"""持久化一批知识件为一个版本化知识包,返回落库寻址三元组 {packageId, versionId, contentHash}。
Args:
file_list: 知识件全文清单 [{path, content}](path 为相对路径,与 read_file 硬编码的知识路径同口径)。
now_ts: 调用方传入的 Unix 时间戳(秒),派生 versionId(显式传入便于测试复现 + 统一时钟口径)。
package_id: 知识包稳定标识(同一套语料的多版本共享;缺省 = DEFAULT_PACKAGE_ID)。
note: 人读的变更说明(可选,进 manifest,供 K3 审计 diff)。
"""
raise NotImplementedError
@abstractmethod
def fetch(self, *, package_id: str = DEFAULT_PACKAGE_ID,
version_id: Optional[str] = None) -> Optional[dict]:
"""按包 id(+ 可选版本)取回知识包。
Returns:
{packageId, versionId, contentHash, note?, files:[{path, sha256, bytes, content | missingContent}]}
或 None(无此包 / 无此版本 / 无 committed 版本)。缺文件不静默:该文件项带 missingContent(物化据此显式失败)。
"""
raise NotImplementedError
@abstractmethod
def latest_committed_version(self, *, package_id: str = DEFAULT_PACKAGE_ID) -> Optional[str]:
"""返回该包最新 committed 版本 id(无 → None)。切片一 resolve_active_version 的默认来源。"""
raise NotImplementedError
class LocalFsKnowledgePackageStore(KnowledgePackageStore):
"""本地文件系统实现 —— 零外部依赖的知识包落库(6c6g / cheap venv 可 import + 往返自检,不连基建)。
落点形态(store_root 下,mirror store.LocalFsStore):
<root>/<package_id>/index.json —— 该包版本索引([{versionId, contentHash, savedAtTs, status}])
<root>/<package_id>/<versionId>/manifest.json —— 该版本 manifest(文件清单 + 各文件 sha256 + bytes)
<root>/<package_id>/<versionId>/files/<相对路径> —— 该版本各知识件全文(fetch 据 manifest 回填 content)
与 BackendKnowledgePackageStore 对 fetch / materialize 呈现完全一致的语义(含 missingContent 标注),
故单测用它跑「save→物化→read_file 读到物化内容」「缺文件显式失败」全链路,不依赖真 MinIO。
"""
def __init__(self, store_root: Path) -> None:
"""store_root: 落库根目录(单测传临时目录隔离;真用时可指一个持久目录)。"""
self._root = Path(store_root)
def _pkg_dir(self, package_id: str) -> Path:
return self._root / package_id
def _index_path(self, package_id: str) -> Path:
return self._pkg_dir(package_id) / "index.json"
def _read_index(self, package_id: str) -> list[dict]:
"""读某包版本索引(无 → 空列表;损坏 → 空列表 + 告警,绝不抛)。"""
p = self._index_path(package_id)
if not p.exists():
return []
try:
data = json.loads(p.read_text(encoding="utf-8"))
return data if isinstance(data, list) else []
except Exception as e: # noqa: BLE001 —— 索引损坏当作无历史 + 告警(落库出口不中断)
print(f"[kb-store] ⚠ 版本索引解析失败(当作无历史): package={package_id} "
f"{type(e).__name__}: {e}", flush=True)
return []
def _write_index(self, package_id: str, index: list[dict]) -> None:
"""写某包版本索引(原子:先写临时文件再 replace,避免半写损坏)。"""
p = self._index_path(package_id)
p.parent.mkdir(parents=True, exist_ok=True)
tmp = p.with_suffix(".json.tmp")
tmp.write_text(json.dumps(index, ensure_ascii=False, indent=2), encoding="utf-8")
tmp.replace(p)
def save(self, file_list: list[dict], *, now_ts: float,
package_id: str = DEFAULT_PACKAGE_ID, note: Optional[str] = None) -> dict:
content_hash = _canonical_content_hash(file_list or [])
index = self._read_index(package_id)
# 幂等:同包已存在相同 content_hash 的 committed 版本 → 不重复落,返回已有。
for entry in index:
if entry.get("contentHash") == content_hash and entry.get("status") == "committed":
print(f"[kb-store] save 幂等命中(整包未变,不重复落): package={package_id} "
f"versionId={entry.get('versionId')} contentHash={content_hash[:12]}", flush=True)
return {"packageId": package_id, "versionId": entry.get("versionId"),
"contentHash": content_hash}
version_id = derive_version_id(content_hash, now_ts)
manifest = _build_manifest(file_list, content_hash, version_id, package_id, note)
ver_dir = self._pkg_dir(package_id) / version_id
try:
ver_dir.mkdir(parents=True, exist_ok=True)
# ① 落各知识件全文(支持 fetch 真回填 content;后端版这步对应「文件落 MinIO」)。
files_root = ver_dir / "files"
written = 0
for f in (file_list or []):
rel = (f or {}).get("path") or ""
content = (f or {}).get("content")
if not _is_safe_rel_path(rel) or not isinstance(content, str):
continue
dst = files_root / rel
dst.parent.mkdir(parents=True, exist_ok=True)
dst.write_text(content, encoding="utf-8")
written += 1
# ② 落 manifest。
(ver_dir / "manifest.json").write_text(
json.dumps(manifest, ensure_ascii=False, indent=2), encoding="utf-8")
# ③ 追加版本索引(status=committed;LocalFs 无跨系统半落库风险,直接 committed)。
index.append({"versionId": version_id, "contentHash": content_hash,
"savedAtTs": int(now_ts), "status": "committed"})
self._write_index(package_id, index)
print(f"[kb-store] save 落库成功: package={package_id} versionId={version_id} "
f"contentHash={content_hash[:12]} files={written}", flush=True)
except Exception as e: # noqa: BLE001 —— 落库失败响亮记日志(审计点),抛给调用方(不静默成功)
print(f"[kb-store] ❌ save 落库失败: package={package_id} versionId={version_id} "
f"{type(e).__name__}: {e}", flush=True)
raise
return {"packageId": package_id, "versionId": version_id, "contentHash": content_hash}
def fetch(self, *, package_id: str = DEFAULT_PACKAGE_ID,
version_id: Optional[str] = None) -> Optional[dict]:
index = [e for e in self._read_index(package_id) if e.get("status") == "committed"]
if not index:
print(f"[kb-store] fetch 未命中: 无 committed 版本 package={package_id}", flush=True)
return None
if version_id is None:
chosen = max(index, key=lambda e: e.get("savedAtTs") or 0)
else:
chosen = next((e for e in index if e.get("versionId") == version_id), None)
if chosen is None:
print(f"[kb-store] fetch 未命中: package={package_id} 无此 committed 版本 "
f"versionId={version_id}(已有: {[e.get('versionId') for e in index]})", flush=True)
return None
vid = chosen.get("versionId")
ver_dir = self._pkg_dir(package_id) / vid
manifest_path = ver_dir / "manifest.json"
if not manifest_path.exists():
print(f"[kb-store] ❌ fetch 归档损坏: package={package_id} versionId={vid} 缺 manifest.json",
flush=True)
return None
try:
manifest = json.loads(manifest_path.read_text(encoding="utf-8"))
except Exception as e: # noqa: BLE001 —— manifest 损坏返 None + 告警(取回不抛,交调用方兜底)
print(f"[kb-store] ❌ fetch manifest 解析失败: package={package_id} versionId={vid} "
f"{type(e).__name__}: {e}", flush=True)
return None
return self._assemble_fetched(manifest, vid, _localfs_reader(ver_dir / "files"))
@staticmethod
def _assemble_fetched(manifest: dict, version_id: str, reader) -> dict:
"""据 manifest + 一个 reader(rel → content|None)组 fetch 返回体,逐文件回填 content / 标 missingContent。
reader 抽象了「文件全文从哪读」:LocalFs 从 files/ 读,Backend 从 MinIO 读。缺文件(reader 返 None)
标 missingContent、不静默——物化据此显式失败。
"""
out_files = []
for item in (manifest.get("files") or []):
rel = (item or {}).get("path") or ""
entry = {"path": rel, "sha256": (item or {}).get("sha256"),
"bytes": (item or {}).get("bytes")}
if not _is_safe_rel_path(rel):
entry["missingContent"] = "路径不安全(含 '..' / 绝对路径),跳过取回"
else:
content = reader(rel)
if content is None:
entry["missingContent"] = "知识包缺该文件全文(manifest 列了但取不到)"
else:
entry["content"] = content
out_files.append(entry)
result = {
"packageId": manifest.get("packageId"),
"versionId": manifest.get("versionId") or version_id,
"contentHash": manifest.get("contentHash"),
"files": out_files,
}
if manifest.get("note"):
result["note"] = manifest["note"]
n_missing = sum(1 for f in out_files if "missingContent" in f)
print(f"[kb-store] fetch 命中: package={result['packageId']} versionId={result['versionId']} "
f"files={len(out_files)} 缺文件={n_missing}", flush=True)
return result
def latest_committed_version(self, *, package_id: str = DEFAULT_PACKAGE_ID) -> Optional[str]:
index = [e for e in self._read_index(package_id) if e.get("status") == "committed"]
if not index:
return None
return max(index, key=lambda e: e.get("savedAtTs") or 0).get("versionId")
def _localfs_reader(files_root: Path):
"""给 LocalFs.fetch 用的文件读取器:rel → content(读得到)或 None(缺文件)。"""
def _read(rel: str) -> Optional[str]:
fp = files_root / rel
if not fp.exists():
return None
try:
return fp.read_text(encoding="utf-8")
except Exception as e: # noqa: BLE001 —— 单文件读失败按缺文件处理(物化显式失败,不静默)
print(f"[kb-store] ⚠ 读知识件失败(按缺文件处理): {rel} {type(e).__name__}: {e}", flush=True)
return None
return _read
class BackendKnowledgePackageStore(KnowledgePackageStore):
"""后端落库真实现 —— manifest 落 MySQL、知识件全文落 MinIO(S3 兼容 OSS),manifest-first 半落库。
════════════════════════════════════════════════════════════════════════
忠实复刻 store.BackendStore 的 manifest-first 范式(见本模块 docstring),但另立表 / 桶:
save —— 三步(半落库不留隐形孤儿):
1) 幂等 / 续跑判定:查 (package_id, content_hash)——命中 committed → 不重写返回已有;命中 pending →
复用该行 versionId 续 put 自愈;未命中 → 派生新 versionId、MySQL 落一行 status=pending(先占位)。
2) 各知识件全文落 MinIO 桶(默认 cheap-kb),object key = <package_id>/<version_id>/<相对路径>。
3) MySQL 把该行 status 翻 committed —— 文件已全落,该版本对 fetch/物化可见。
fetch —— 只返 committed(pending 半成品对取回不可见):读 MySQL manifest + 据 key 从 MinIO 拉文件全文,
回填进 files[].content;缺文件标 missingContent(不静默)。
【惰性 import + 代理旁路(本分支硬约束)】
pymysql / minio 在方法体内惰性 import(模块顶层不 import),保证 cheap venv(未装这俩)能 import 本模块、
跑 LocalFs 单测;真连基建的往返排 mini-infra 窗口。连 MinIO / MySQL 前先 client.install_proxy_bypass
把内网 host 并入 NO_PROXY(内网经本机 clash fake-ip 会被转走 → 502;复用 tier2 既有旁路范式)。
连接 / 读写失败不静默吞:响亮记日志(审计点)+ 抛给调用方(物化侧再据缺文件显式失败兜)。
════════════════════════════════════════════════════════════════════════
"""
def __init__(self) -> None:
self._table = _KB_MYSQL_TABLE
self._schema_ensured = False
# ── 配置读取(经 service.infra_config 三级回落:env > infra.yaml > default;best-effort)──
def _mysql_conf(self) -> dict:
from service import infra_config # noqa: PLC0415 —— 惰性 import(与 store.py 同款,需 gen-worker 在 sys.path)
return {
"host": infra_config.get("mysql", "host", "127.0.0.1"),
"port": infra_config.get("mysql", "port", 3306),
"user": infra_config.get("mysql", "user", "root"),
"password": infra_config.get("mysql", "password", ""),
"database": infra_config.get("mysql", "database", "tier2"),
}
def _minio_conf(self) -> dict:
from service import infra_config # noqa: PLC0415
return {
"endpoint": infra_config.get("minio", "endpoint", "127.0.0.1:9000"),
"access_key": infra_config.get("minio", "access_key", ""),
"secret_key": infra_config.get("minio", "secret_key", ""),
# 另立知识包专用桶:infra.yaml minio.kb_bucket(缺省回落 cheap-kb),绝不用 ASSET-SRC 的 tier2-src。
"bucket": infra_config.get("minio", "kb_bucket", _KB_BUCKET_DEFAULT),
"secure": infra_config.get("minio", "secure", False),
}
def _install_proxy_bypass(self, host_url: str) -> None:
"""连基建前把内网 host 并入 NO_PROXY(复用 worker.client.install_proxy_bypass;失败不致命)。"""
try:
from worker import client # noqa: PLC0415
client.install_proxy_bypass(host_url)
except Exception as e: # noqa: BLE001 —— 装旁路失败不阻断(真连时若被代理转走会在连接处响亮报)
print(f"[kb-store] ⚠ 代理旁路装配失败(继续): {host_url} {type(e).__name__}: {e}", flush=True)
def _connect_mysql(self):
"""惰性连 MySQL(pymysql);连前装代理旁路。失败响亮抛(落库审计点,不静默)。"""
conf = self._mysql_conf()
self._install_proxy_bypass(f"http://{conf['host']}:{conf['port']}")
try:
import pymysql # noqa: PLC0415 —— 惰性 import(cheap venv 无 pymysql 也能 import 本模块)
except Exception as e: # noqa: BLE001
print(f"[kb-store] ❌ 缺 pymysql(惰性 import 失败): {type(e).__name__}: {e};"
f"请在运行机 pip install pymysql。", flush=True)
raise
return pymysql.connect(
host=conf["host"], port=int(conf["port"]), user=conf["user"],
password=conf["password"], database=conf["database"],
charset="utf8mb4", autocommit=True,
cursorclass=pymysql.cursors.DictCursor,
)
def _connect_minio(self):
"""惰性连 MinIO,返回 (client, bucket);连前装代理旁路。best-effort 建桶。失败响亮抛。"""
conf = self._minio_conf()
self._install_proxy_bypass(f"http://{conf['endpoint']}")
try:
from minio import Minio # noqa: PLC0415 —— 惰性 import(cheap venv 无 minio 也能 import)
except Exception as e: # noqa: BLE001
print(f"[kb-store] ❌ 缺 minio(惰性 import 失败): {type(e).__name__}: {e};"
f"请在运行机 pip install minio。", flush=True)
raise
cli = Minio(conf["endpoint"], access_key=conf["access_key"],
secret_key=conf["secret_key"], secure=bool(conf["secure"]))
bucket = conf["bucket"]
try:
if not cli.bucket_exists(bucket):
cli.make_bucket(bucket)
print(f"[kb-store] 建 MinIO 知识包桶: {bucket}", flush=True)
except Exception as e: # noqa: BLE001 —— 建桶/探测失败响亮抛(连不上对象存储,落库无意义)
print(f"[kb-store] ❌ MinIO 桶探测/创建失败: bucket={bucket} {type(e).__name__}: {e}", flush=True)
raise
return cli, bucket
def _ensure_schema(self, conn) -> None:
"""建表(每实例一次;idempotent)。本表出生即带 status 列,无历史无 status 列旧表,故无 ALTER。"""
if self._schema_ensured:
return
with conn.cursor() as cur:
cur.execute(_DDL_KB_PACKAGE_VERSION)
self._schema_ensured = True
def _source_url(self, endpoint: str, secure: bool, bucket: str, prefix: str) -> str:
"""拼 manifest.addressing.sourceUrl(MinIO 版本前缀只读镜像 URL;仅记账,fetch 不靠它寻址)。"""
scheme = "https" if secure else "http"
return f"{scheme}://{endpoint}/{bucket}/{prefix}"
def _put_files(self, client, bucket: str, prefix: str, file_list: list[dict]) -> int:
"""把知识件全文逐个 put 进 MinIO(object key = prefix + 相对路径)。返回成功写入数。路径安全复用共用校验。"""
import io # noqa: PLC0415
written = 0
for f in (file_list or []):
rel = (f or {}).get("path") or ""
content = (f or {}).get("content")
if not _is_safe_rel_path(rel) or not isinstance(content, str):
continue
data = content.encode("utf-8")
client.put_object(bucket, f"{prefix}{rel}", io.BytesIO(data), length=len(data),
content_type="text/plain; charset=utf-8")
written += 1
return written
def save(self, file_list: list[dict], *, now_ts: float,
package_id: str = DEFAULT_PACKAGE_ID, note: Optional[str] = None) -> dict:
import datetime # noqa: PLC0415
content_hash = _canonical_content_hash(file_list or [])
conn = self._connect_mysql()
try:
self._ensure_schema(conn)
# ── 幂等 / 续跑判定:查 (package_id, content_hash) 现状 ──
with conn.cursor() as cur:
cur.execute(
f"SELECT version_id, status FROM {self._table} "
f"WHERE package_id=%s AND content_hash=%s LIMIT 1",
(package_id, content_hash),
)
row = cur.fetchone()
if row and row.get("status") == "committed":
existing_vid = row.get("version_id")
print(f"[kb-store] BackendStore.save 幂等命中(整包未变且已提交): package={package_id} "
f"versionId={existing_vid} contentHash={content_hash[:12]}", flush=True)
return {"packageId": package_id, "versionId": existing_vid, "contentHash": content_hash}
reused_pending = bool(row and row.get("status") == "pending")
version_id = row.get("version_id") if reused_pending else derive_version_id(content_hash, now_ts)
if reused_pending:
print(f"[kb-store] BackendStore.save 复用 pending 版本续跑自愈(同前缀真覆盖): "
f"package={package_id} versionId={version_id}", flush=True)
prefix = f"{package_id}/{version_id}/"
client, bucket = self._connect_minio()
manifest = _build_manifest(file_list, content_hash, version_id, package_id, note)
# ── ① 未命中时先落 pending 行占位(此刻 MinIO 尚无对象,半落库崩溃只留【可见 pending 行】)──
if not reused_pending:
mc = self._minio_conf()
addressing = {"store": "mysql+oss",
"sourceUrl": self._source_url(mc["endpoint"], bool(mc["secure"]), bucket, prefix)}
manifest_with_addr = dict(manifest)
manifest_with_addr["addressing"] = addressing
created_at = datetime.datetime.utcfromtimestamp(int(now_ts)).strftime("%Y-%m-%d %H:%M:%S")
with conn.cursor() as cur:
cur.execute(
f"INSERT INTO {self._table} "
f"(package_id, version_id, content_hash, manifest_json, addressing, created_at, status) "
f"VALUES (%s, %s, %s, %s, %s, %s, 'pending')",
(package_id, version_id, content_hash,
json.dumps(manifest_with_addr, ensure_ascii=False),
addressing["sourceUrl"], created_at),
)
# ── ② 落知识件全文到 MinIO(pending 续跑同前缀真覆盖)──
written = self._put_files(client, bucket, prefix, file_list or [])
# ── ③ 翻 committed(文件已全落 → 该版本对 fetch / 物化可见)──
with conn.cursor() as cur:
cur.execute(
f"UPDATE {self._table} SET status='committed' "
f"WHERE package_id=%s AND version_id=%s",
(package_id, version_id),
)
print(f"[kb-store] BackendStore.save 落库成功(committed): package={package_id} "
f"versionId={version_id} contentHash={content_hash[:12]} files={written} "
f"bucket={bucket} prefix={prefix}", flush=True)
return {"packageId": package_id, "versionId": version_id, "contentHash": content_hash}
except Exception as e: # noqa: BLE001 —— 落库失败响亮记日志(审计点),抛给调用方(不静默成功)
print(f"[kb-store] ❌ BackendStore.save 落库失败: package={package_id} "
f"contentHash={content_hash[:12]} {type(e).__name__}: {e}", flush=True)
raise
finally:
try:
conn.close()
except Exception: # noqa: BLE001
pass
def fetch(self, *, package_id: str = DEFAULT_PACKAGE_ID,
version_id: Optional[str] = None) -> Optional[dict]:
conn = self._connect_mysql()
try:
with conn.cursor() as cur:
if version_id is None:
cur.execute(
f"SELECT version_id, manifest_json FROM {self._table} "
f"WHERE package_id=%s AND status='committed' ORDER BY created_at DESC, id DESC LIMIT 1",
(package_id,),
)
else:
cur.execute(
f"SELECT version_id, manifest_json FROM {self._table} "
f"WHERE package_id=%s AND version_id=%s AND status='committed' LIMIT 1",
(package_id, version_id),
)
row = cur.fetchone()
if not row:
print(f"[kb-store] BackendStore.fetch 未命中: package={package_id} "
f"versionId={version_id or '(latest)'}", flush=True)
return None
vid = row.get("version_id")
try:
manifest = json.loads(row.get("manifest_json") or "{}")
except Exception as e: # noqa: BLE001 —— manifest 损坏 → 抛(取回不返半残数据冒充成功)
print(f"[kb-store] ❌ BackendStore.fetch manifest 解析失败: package={package_id} "
f"versionId={vid} {type(e).__name__}: {e}", flush=True)
raise
finally:
try:
conn.close()
except Exception: # noqa: BLE001
pass
client, bucket = self._connect_minio()
prefix = f"{package_id}/{vid}/"
return LocalFsKnowledgePackageStore._assemble_fetched(
manifest, vid, _minio_reader(client, bucket, prefix))
def latest_committed_version(self, *, package_id: str = DEFAULT_PACKAGE_ID) -> Optional[str]:
conn = self._connect_mysql()
try:
self._ensure_schema(conn)
with conn.cursor() as cur:
cur.execute(
f"SELECT version_id FROM {self._table} "
f"WHERE package_id=%s AND status='committed' ORDER BY created_at DESC, id DESC LIMIT 1",
(package_id,),
)
row = cur.fetchone()
return row.get("version_id") if row else None
finally:
try:
conn.close()
except Exception: # noqa: BLE001
pass
def _minio_reader(client, bucket: str, prefix: str):
"""给 Backend.fetch 用的文件读取器:rel → content(MinIO 取到)或 None(缺对象,标 missingContent)。"""
def _read(rel: str) -> Optional[str]:
object_name = f"{prefix}{rel}"
try:
resp = client.get_object(bucket, object_name)
try:
return resp.read().decode("utf-8")
finally:
resp.close()
resp.release_conn() # MinIO HTTPResponse 必须 close+release_conn 防连接泄漏
except Exception as e: # noqa: BLE001 —— 单对象取回失败按缺文件处理(物化显式失败,不静默)
print(f"[kb-store] ⚠ 取知识件对象失败(按缺文件处理): {object_name} "
f"{type(e).__name__}: {e}", flush=True)
return None
return _read
# ── 当前激活版本解析 + 治理激活版本源(K3 接缝点落地)──────────────────────────────────────
# 进程级单例槽:worker 启动期由 kb_nacos.build_and_attach_active_version_source() 注入一个「激活版本源」
# (读 Nacos 下发的当前激活知识包版本);resolve_active_version 取值链最前端先查它。默认 None(未 attach)=
# 切片一行为:env 钉版 > 最新 committed(字节不变)。attach/detach 对称于 genconfig.attach_hot_source,
# 放本模块(纯 dict 槽、不牵 nacos 依赖,cheap venv 可 import);真正建 client + 订阅 + 变更即物化在 kb_nacos。
_active_version_source: Any = None
def attach_active_version_source(source: Any) -> None:
"""worker 启动期注入知识包激活版本源(K3);此后 resolve_active_version 最先查它。
source 只需满足 get(package_id)->version|None 契约(kb_nacos.KbActiveVersionSource 即是),本模块不依赖其
内部实现——故单测可注入任意 fake 源。幂等可重入:重复 attach 覆盖为最新源。
"""
global _active_version_source
_active_version_source = source
def detach_active_version_source() -> None:
"""卸载激活版本源,resolve_active_version 回落 env/最新 committed(K3 回退开关的进程内落点)。"""
global _active_version_source
_active_version_source = None
def resolve_active_version(store: KnowledgePackageStore, *,
package_id: str = DEFAULT_PACKAGE_ID) -> Optional[str]:
"""解析「当前激活的知识包版本」:治理激活版本源(Nacos 下发) > env TIER2_KB_ACTIVE_VERSION 钉版 > 最新 committed。
W-CFG-KB 治理面(K3)已把接缝落地:worker 启动期若接了知识包配置集的 Nacos 激活版本源(kb_nacos),这里最先查它
——改一版知识件经配置控制台定版 + 激活,激活推送到 Nacos,worker 读到新版本号即据它物化。未接源(默认 / TIER2_KB_ROOT
未配 / Nacos 未启用)则回落切片一行为:env 钉版 > 最新 committed(字节不变)。
治理源取值 best-effort:源不可达 / 无此包激活版 / 取值异常都静默回落 env/最新 committed,绝不让治理源拖垮物化
(→拖垮生成)。返回 None = 无激活版(物化 no-op,read_file 回落仓根)。
"""
# ⓪ 治理激活版本源(K3 最高优先级):Nacos 下发的当前激活知识包版本。
src = _active_version_source
if src is not None:
try:
gov = src.get(package_id)
if gov and str(gov).strip():
return str(gov).strip()
except Exception as e: # noqa: BLE001 —— 治理源异常绝不连累物化:静默回落 env/最新 committed
print(f"[kb-store] ⚠ 读知识包激活版本源失败(回落 env/最新 committed): "
f"{type(e).__name__}: {e}", flush=True)
# ① env 显式钉版。
pinned = os.environ.get(KB_ACTIVE_VERSION_ENV)
if pinned and pinned.strip():
return pinned.strip()
# ② 最新 committed 版本。
return store.latest_committed_version(package_id=package_id)
def kb_root_from_env() -> Optional[Path]:
"""物化落点:读 env TIER2_KB_ROOT(与 cheap_run._kb_root 读侧同一个 env)。未配 → None(知识件外置关闭)。"""
raw = os.environ.get(KB_ROOT_ENV)
if not raw or not raw.strip():
return None
try:
return Path(raw).expanduser()
except Exception: # noqa: BLE001
return None
def default_kb_store() -> KnowledgePackageStore:
"""默认知识包 store 工厂(按 env KB_STORE 切换;mirror store.default_store)。
- KB_STORE=backend → BackendKnowledgePackageStore(MySQL+MinIO;部署 / 窗口真跑时切)。
- 其余(含未设)→ LocalFsKnowledgePackageStore(落 <TIER2_KB_ROOT>/_store 或临时;本地可往返自检)。
"""
if (os.environ.get("KB_STORE") or "").strip().lower() == "backend":
print("[kb-store] default_kb_store → BackendKnowledgePackageStore(KB_STORE=backend)", flush=True)
return BackendKnowledgePackageStore()
# LocalFs 落库根:TIER2_KB_ROOT 旁挂 _store(与物化产物同父目录、隔离);未配 KB 根则落当前目录 _kb-store。
kb_root = kb_root_from_env()
store_root = (kb_root.parent / "_kb-store") if kb_root else Path.cwd() / "_kb-store"
return LocalFsKnowledgePackageStore(store_root)
# ── 物化:按当前激活版从知识包取回、校验、原子换入知识根(缺文件显式失败 + 回落上一份)────────────
def materialize_active(store: KnowledgePackageStore, kb_root: Optional[Path] = None, *,
package_id: str = DEFAULT_PACKAGE_ID,
active_version: Optional[str] = None) -> dict:
"""把「当前激活版知识包」物化到知识根 kb_root(检出之外),供 cheap_run.read_file shadow 读。
物化纪律(安全语义,别 best-effort 掉):
· 缺文件【显式失败】:manifest 列了某文件但取不到(missingContent),抛 KnowledgeMaterializeError,
绝不物化半拉子包;
· 逐文件指纹校验:物化到 staging 时每个文件的 sha256 必须等于 manifest,不等即抛(K2 寻址断言);
· 物化失败【回落上一份】:整个物化在 staging 里做完并校验才原子换入 kb_root,任一步失败 kb_root 未被
触碰 = 上一份已物化包原样保留,agent 读到的仍是完整旧包、绝不半拉子。
Args:
store: 知识包 store(LocalFs / Backend)。
kb_root: 物化落点;None → 读 env TIER2_KB_ROOT。未配则 no-op(知识件外置关闭,read_file 回落仓根)。
package_id: 知识包 id。
active_version: 指定物化哪个版本;None → resolve_active_version(env 钉版 / 最新 committed)。
Returns:
{materialized: bool, ...}:materialized=False 表示 no-op(未配 kb_root / 无激活版);
materialized=True 带 {versionId, files, kbRoot}。
Raises:
KnowledgeMaterializeError: 激活版取不到 / 缺文件 / 指纹不匹配(显式失败,kb_root 不动)。
"""
kb_root = kb_root or kb_root_from_env()
if kb_root is None:
print("[kb-materialize] TIER2_KB_ROOT 未配置 → 知识件外置关闭,不物化(read_file 回落仓根)。", flush=True)
return {"materialized": False, "reason": "TIER2_KB_ROOT 未配置"}
version = active_version or resolve_active_version(store, package_id=package_id)
if version is None:
print(f"[kb-materialize] 无激活知识包版本(package={package_id}) → 不物化(read_file 回落仓根)。",
flush=True)
return {"materialized": False, "reason": "无激活知识包版本"}
pkg = store.fetch(package_id=package_id, version_id=version)
if pkg is None:
# 激活版取不到:显式失败(不静默清空 kb_root),kb_root 不动 = 回落上一份已物化包。
raise KnowledgeMaterializeError(
f"当前激活知识包版本取不到: package={package_id} versionId={version} —— "
f"拒绝物化(kb_root 保留上一份,agent 仍读完整旧包)")
files = pkg.get("files") or []
missing = [f.get("path") for f in files if "missingContent" in f]
if missing:
# 缺文件显式失败:manifest 列了但取不到 → 绝不物化半拉子包(类比 ASSET-SRC「缺源显式失败」)。
raise KnowledgeMaterializeError(
f"知识包缺文件,拒绝物化半拉子包: package={package_id} versionId={version} "
f"{len(missing)} 个: {missing[:5]}{'' if len(missing) > 5 else ''}")
kb_root = kb_root.resolve()
staging = kb_root.parent / f"{kb_root.name}.staging-{version}"
# 清残留 staging(上次物化中途崩溃留下的),重新来过。
if staging.exists():
shutil.rmtree(staging, ignore_errors=True)
try:
# ── 物化到 staging + 逐文件指纹校验(校验不过即抛,kb_root 未动)──
written = 0
for f in files:
rel = f.get("path") or ""
content = f.get("content")
if not _is_safe_rel_path(rel) or not isinstance(content, str):
raise KnowledgeMaterializeError(
f"物化异常文件项(路径不安全或 content 缺失): package={package_id} "
f"versionId={version} path={rel!r}")
got = _file_content_hash(content)
want = f.get("sha256")
if want and got != want:
# K2 寻址断言:物化后各文件 hash 必须等于 manifest(内容在途损坏即拒)。
raise KnowledgeMaterializeError(
f"知识件指纹不匹配 manifest,拒绝物化: package={package_id} versionId={version} "
f"path={rel} 期望={str(want)[:12]}… 实得={got[:12]}")
dst = staging / rel
dst.parent.mkdir(parents=True, exist_ok=True)
dst.write_text(content, encoding="utf-8")
written += 1
# 落一份物化元信息(供诊断 / K3 对账「知识根现在是哪个版本」)。
(staging / ".kb-materialized.json").write_text(
json.dumps({"packageId": package_id, "versionId": version, "files": written},
ensure_ascii=False, indent=2), encoding="utf-8")
# ── 原子换入 kb_root(旧的挪到 .prev 作回落备份)──
_atomic_swap(kb_root, staging)
print(f"[kb-materialize] 物化成功: package={package_id} versionId={version} "
f"files={written} kbRoot={kb_root}", flush=True)
return {"materialized": True, "versionId": version, "files": written, "kbRoot": str(kb_root)}
except Exception:
# 任何异常:清 staging(kb_root 未换入 = 保留上一份),把异常抛给调用方(物化失败要可见)。
shutil.rmtree(staging, ignore_errors=True)
raise
def _atomic_swap(kb_root: Path, staging: Path) -> None:
"""把物化好并校验过的 staging 原子换入 kb_root;旧包挪到 <kb_root>.prev 作「上一份」回落备份。
时序:清旧 .prev → kb_root 存在则改名到 .prev → staging 改名到 kb_root。同一文件系统内 rename 是原子操作。
若换入(staging → kb_root)失败,尽力把 .prev 恢复回 kb_root(回落上一份),再抛。
"""
backup = kb_root.parent / (kb_root.name + ".prev")
kb_root.parent.mkdir(parents=True, exist_ok=True)
if backup.exists():
shutil.rmtree(backup, ignore_errors=True)
moved_to_backup = False
if kb_root.exists():
os.replace(str(kb_root), str(backup)) # 旧包挪走留作回落
moved_to_backup = True
try:
os.replace(str(staging), str(kb_root)) # staging 换入
except Exception:
# 换入失败:恢复上一份(把 .prev 挪回 kb_root),让 agent 读到的仍是完整旧包。
if moved_to_backup and backup.exists() and not kb_root.exists():
try:
os.replace(str(backup), str(kb_root))
except Exception: # noqa: BLE001 —— 恢复失败也别吞原始异常
pass
raise
# ── 知识件打包清单(切片一默认集 = cheap_roles.py read_file 指向的知识件)──────────────────────
# 便宜档生成 agent 会 read_file 的知识路径(与 cheap_roles.py 硬编码一致;read_file shadow 按这些相对
# 路径在知识根 shadow 命中)。切片一先外置这批「设计 skill + 品类过门正例」;黄金脚手架 / tier2 Phaser
# 多文件目录树可后续按需扩(store 接受任意文件树,加进本清单即可)。
DEFAULT_KNOWLEDGE_PATHS = [
".agents/skills/littlejs-game-dev.md",
".agents/skills/sim-business-game-design.md",
".agents/skills/narrative-game-design.md",
".agents/skills/trpg-game-design.md",
".agents/skills/heritage-game-design.md",
".agents/skills/puzzle-game-design.md",
"game-runtime/games/_fewshot-feiyi/src/game-logic.js",
"game-runtime/games/_fewshot-puzzle/src/game-logic.js",
"game-runtime/games/_template-story/src/game-logic.js",
]
def collect_knowledge_files(repo_root: Path,
paths: Optional[list[str]] = None) -> list[dict]:
"""从 repo 检出读一批知识件全文,组 [{path, content}] 供 save 打包。缺文件跳过 + 告警(打包 best-effort:
只打进检出里真存在的;缺文件的【显式失败】是物化侧的事,不是打包侧)。"""
repo_root = Path(repo_root)
out = []
for rel in (paths or DEFAULT_KNOWLEDGE_PATHS):
fp = repo_root / rel
if not fp.is_file():
print(f"[kb-collect] ⚠ 知识件缺失,跳过打包: {rel}", flush=True)
continue
try:
out.append({"path": rel, "content": fp.read_text(encoding="utf-8")})
except Exception as e: # noqa: BLE001
print(f"[kb-collect] ⚠ 知识件读失败,跳过: {rel} {type(e).__name__}: {e}", flush=True)
return out