aps-agent/server/knowledge/assets.py

274 lines
17 KiB
Python
Raw Permalink Normal View History

2026-07-21 11:05:57 +08:00
# ============================================================
# 知识资产库(moduleId: knowledge-assets, 可重生 ✅, 黄金测试 tests/golden/test_m3_knowledge.py)
# plan.md §8.1:四类资产(算法知识/企业 SOP/历史案例/工艺知识),
# 版本化 + 入库审批标记 + 检索强制带出处。
# M3 形态:单文件 JSON + 种子资产;接口稳定,后续可换向量库不改调用方。
# ============================================================
from __future__ import annotations # 前向类型引用
import json # 序列化
import os # 路径
import tempfile # 原子写
import threading # 互斥
import uuid # 资产 ID
from datetime import datetime # 时间戳
from typing import Any # 类型标注
from server.timeutil import fmt_dt # 时间格式化
# 存储路径(与世界状态同目录)
_PATH_ENV = "APS_KNOWLEDGE_PATH"
def _seed_assets() -> list[dict[str, Any]]:
"""种子知识资产(演示语料 + 机加工工艺模式;真实项目可由导入管线追加)。"""
2026-07-21 11:05:57 +08:00
now = fmt_dt(datetime.now()) # 统一时间戳
def asset(kind: str, title: str, content: str, tags: list[str]) -> dict[str, Any]:
return { # 资产结构(DatasetAsset 的知识侧近亲)
"assetId": uuid.uuid4().hex[:10], # 资产 ID(@知识 引用用)
"kind": kind, # algorithm / sop / case / process
"title": title, # 标题(引用与出处展示)
"content": content, # 正文(M3 单块;后续分 chunk)
"tags": tags, # 检索加权标签
"version": "v1", # 版本号(重复入库递增)
"approved": True, # 审批标记(§8.1 入库需审批)
"createdAt": now, # 入库时间
}
base = [
2026-07-21 11:05:57 +08:00
asset("sop", "换线标准SOP",
"产线换型必须遵守:① 同产品族连续生产优先,减少换线;② 换线作业标准 30 分钟,"
"跨产品族换线 60 分钟;③ 每日 16:00 后不安排换线(夜班技师不足);"
"④ 换线前必须完成上一批次尾检,未尾检不得拆线。",
["换线", "SOP", "产品族", "夜班"]),
asset("sop", "插单审批流程",
"VIP 客户插单:计划员评估影响(延迟/成本)→ 主管审批 → 发布新版本。"
"非 VIP 插单原则上排入下一个计划周期;确需插入的走例外审批,须附影响评估报告。",
["插单", "审批", "VIP", "例外"]),
asset("sop", "主数据录入顺序",
"对齐聚制云操作手册:主数据必须按依赖顺序录入,不可跳步。"
"① 生产模型:工厂→车间→工段→产线→工位(左侧树状图维护上下级);"
"② 设备模型:设备挂到产线/工位;"
"③ 工艺模型:先物料(成品/半成品/原材料),再产品BOM(树形增删改),再工序库,再工艺路线(新建后「去配置」拖拽工序并填工时);"
"④ 产线-产品绑定:产品与产线多对多配置完成后,才允许订单排产。"
"缺 BOM、缺路线、未绑定产线时,排产应提示失败环节与缺失项,而不是给出无依据的结果。"
"现场数据导入后可用对话「查主数据」「查订单」「查工艺路线」核对。",
["主数据", "录入顺序", "BOM", "工艺路线", "产线绑定", "聚制云", "查询"]),
asset("sop", "工艺路线配置规范",
"工艺路线 = 产品的有序工序流。配置步骤:"
"① 在工序管理维护标准工序;② 新建工艺路线(名称/编码/产品);"
"③ 进入「去配置」:将工序拖入流程图并用箭头串顺序;"
"④ 为每步填写准备时间、单件工时、是否外协;外协步骤在订单分解时生成委外建议;"
"⑤ 审核通过后设为默认版本,仅 ACTIVE+default 参与排产与追溯。"
"修改用量或工时后,必须重新分解/试排,追溯页应显示新版本号。"
"机加工行业通用生成模式见知识库「机加工工艺路线生成总则」及车铣钻磨装配等分册。",
["工艺路线", "工序", "去配置", "外协", "默认版本", "机加工"]),
asset("sop", "订单排产与追溯",
"计划链(对齐聚制云 §4.9):销售订单 → 生产订单/采购订单/外协订单 → 生产工单。"
"排产前检查:产线产品配置、默认BOM、默认工艺路线、库存是否覆盖关键料。"
"排产成功后:在生产订单看排产日历与工单列表;工单可回指销售订单与产品BOM。"
"本系统用「计划追溯」聚合同一视图:主数据依据(BOM用量+库存缺口、工艺步骤)、"
"分解建议(采购/委外及状态)、排产产物(PO/WO)、各产线负荷占用分钟、冲突。"
"口令示例:「追溯 102285668」「查订单」「查 BOM」。",
["订单", "排产", "追溯", "库存", "负荷", "MRP", "工单"]),
asset("sop", "库存与齐套口径",
"齐套净需求 = BOM毛需求 − 现有库存 − 在途。缺口>0 生成采购建议,"
"建议下单日 = 交期 − 采购前置期 − 1天缓冲。关键料缺口在追溯页标出。"
"库存字段在物料主数据维护(stock/inTransit/safetyStock);"
"改库存后需重新分解与试排,负荷与缺料冲突才会更新。"
"排产引擎对缺料产生 MATERIAL_SHORTAGE 冲突,不阻断出草案,但必须在追溯/冲突列表可见。",
["库存", "齐套", "缺料", "采购前置期", "安全库存"]),
asset("sop", "对话查询主数据与订单",
"自然语言可直接检索当前世界数据(只读):"
"「有哪些订单」「查订单 102285668」「查物料」「查工序库」"
"「查 BOM」「查工艺路线」「主数据概览」「柔性资源」。"
"追溯全链用「追溯 订单号」。问工艺怎么生成用「机加工工艺路线」「装配工艺模式」走知识库。"
"写入仍走确认卡(新建订单/物料等)。",
["对话", "查询", "订单", "主数据", "BOM", "工艺路线"]),
2026-07-21 11:05:57 +08:00
asset("algorithm", "交期优先策略说明",
"DELIVERY_FIRST 按交期升序 + 客户等级加权排序订单,优先保障最紧交期。"
"适用:交付压力大、订单交期集中。代价:可能牺牲产能均衡与换线次数。"
"参数:客户等级权重(VIP=3/重要=2/普通=1),同交期按下单时间先到先排。",
["交期优先", "DELIVERY_FIRST", "策略", "权重"]),
asset("algorithm", "产能均衡策略说明",
"CAPACITY_BALANCE 以产线负荷方差最小为目标分配订单,避免忙闲不均。"
"适用:多产线可互替、瓶颈不固定。代价:个别紧急订单可能被排后。"
"建议与交期承诺看板联用,均衡后逐单核对交付风险。追溯页的「负荷占用分钟」可核对是否均衡。",
2026-07-21 11:05:57 +08:00
["产能均衡", "CAPACITY_BALANCE", "负荷", "方差"]),
asset("algorithm", "柔性能力池排产说明",
"PoolEngine(柔性轨)按设备能力池动态组虚拟产线,与固定产线 RuleEngine 单据链并行。"
"输入:flexOrders + flexBom + flexRoutings + flexEquipment;"
"输出:flexVirtualLines + flexWorkOrders,可在柔性甘特与聊天 flex-schedule 块下钻。"
"瓶颈产能法按能力池日产能找限制性瓶颈。固定轨商务订单与柔性轨订单不要混为一谈。",
["柔性", "能力池", "虚拟产线", "瓶颈", "PoolEngine"]),
2026-07-21 11:05:57 +08:00
asset("process", "产线A设备能力",
"产线A(L001)适合小批量多品种:SMT 贴片 + 波峰焊;单班产能 400 标准件/日;"
"维保窗口每周三 08:00-10:00 固定停机;夜班仅一班技师,复杂换型不排夜班。"
"(演示电子装配口径;机加工请检索「机加工工艺路线生成总则」。)",
2026-07-21 11:05:57 +08:00
["产线A", "L001", "产能", "维保"]),
asset("process", "高压线束与PDU共线要点",
"柔性场景案例:高压线束/PDU/充电枪/连接器共线。"
"瓶颈工序:压接端子、激光焊接(全厂最紧)、ESD装配、插拔力测试。"
"压接机可跨区移动;排产须计换型与移动时间。物料齐套看 flexBom 关键料。"
"机械装配类现场单请用「城轨机构装配路线模板」与「查工艺路线」。",
["线束", "PDU", "瓶颈", "压接", "焊接"]),
2026-07-21 11:05:57 +08:00
asset("case", "6月底比亚迪插单复盘",
"6月28日比亚迪 VIP 插单 2000 件:采用交期优先重排,牺牲产线B均衡度 12%,"
"全部订单按期交付。经验:VIP 插单后应立即跑一次沙盒对比,确认对存量订单影响 <5% 再发布。"
"发布前用计划追溯核对库存缺口与负荷占用。",
2026-07-21 11:05:57 +08:00
["比亚迪", "插单", "复盘", "沙盒"]),
]
from server.knowledge.machining_patterns import machining_process_assets
return base + machining_process_assets(asset)
2026-07-21 11:05:57 +08:00
class KnowledgeStore:
"""知识资产仓:加载/列出/取回/新增(单文件 JSON + 原子写;P0 读 / P1 新增)。"""
def __init__(self, path: str | None = None) -> None:
"""初始化:加载资产,空库自动播种。权力等级 P0(构造只读+首次播种)。"""
self.path = path or os.environ.get(_PATH_ENV, "server/data/knowledge.json") # 存储路径
self._lock = threading.Lock() # 并发保护
self.assets: list[dict[str, Any]] = self._load() # 资产列表
if not self.assets: # 空库播种
self.assets = _seed_assets() # 种子资产
self._write() # 落盘
def _load(self) -> list[dict[str, Any]]:
"""从磁盘加载;缺失/损坏返回空(触发播种)。"""
try:
with open(self.path, "r", encoding="utf-8") as f: # 读文件
return json.load(f).get("assets", []) # 取资产数组
except (FileNotFoundError, json.JSONDecodeError): # 缺失/损坏
return []
def _write(self) -> None:
"""原子写盘(tmp + os.replace,与 WorldStore 同法)。"""
os.makedirs(os.path.dirname(self.path) or ".", exist_ok=True) # 确保目录
fd, tmp = tempfile.mkstemp(dir=os.path.dirname(self.path) or ".", suffix=".tmp") # 临时文件
try:
with os.fdopen(fd, "w", encoding="utf-8") as f: # 写临时文件
json.dump({"assets": self.assets}, f, ensure_ascii=False, indent=1) # 序列化
os.replace(tmp, self.path) # 原子替换
except BaseException: # 失败清理
if os.path.exists(tmp):
os.unlink(tmp)
raise
def list_meta(self) -> list[dict[str, Any]]:
"""资产元信息清单(不含正文;知识面板/命令面板数据源,P0)。"""
return [{k: a[k] for k in ("assetId", "kind", "title", "version", "createdAt")}
for a in self.assets]
def get(self, asset_id: str) -> dict[str, Any] | None:
"""按 ID 取完整资产(@知识 引用解析用,P0)。"""
return next((a for a in self.assets if a["assetId"] == asset_id), None)
def find_by_title(self, title: str) -> dict[str, Any] | None:
"""按标题模糊取资产(@知识:换线标准SOP 的解析路径,P0)。"""
t = title.strip() # 清理输入
return next((a for a in self.assets if t and t in a["title"]), None)
def add(self, kind: str, title: str, content: str, tags: list[str] | None = None) -> dict[str, Any]:
"""新增资产(P1:只新增知识,不碰世界状态;报告入库走此口 §9.10 规则3)。"""
return self.add_with_chunks(kind=kind, title=title, content=content, chunks=None, tags=tags)
def add_with_chunks(
self, *, kind: str, title: str, content: str,
chunks: list[dict[str, Any]] | None = None,
tags: list[str] | None = None,
source: str | None = None,
approved: bool | None = None,
) -> dict[str, Any]:
"""新增资产并可附带切块(RAG 导入);同名标题递增版本。"""
with self._lock:
existing = self.find_by_title(title)
version = "v1"
if existing and existing["title"] == title:
version = f"v{int(existing['version'].lstrip('v')) + 1}"
asset = {
2026-07-21 11:05:57 +08:00
"assetId": uuid.uuid4().hex[:10], "kind": kind, "title": title,
"content": content, "tags": tags or [], "version": version,
"approved": (kind != "report") if approved is None else bool(approved),
2026-07-21 11:05:57 +08:00
"createdAt": fmt_dt(datetime.now()),
"source": source or "",
"chunks": list(chunks or []),
2026-07-21 11:05:57 +08:00
}
# 给 chunk 补 asset 引用
for c in asset["chunks"]:
c.setdefault("assetId", asset["assetId"])
c.setdefault("title", title)
c.setdefault("version", version)
c.setdefault("kind", kind)
self.assets.append(asset)
self._write()
return asset
def iter_search_units(self) -> list[dict[str, Any]]:
"""检索单元:有 chunks 时按 chunk 展开,否则整篇资产。"""
units: list[dict[str, Any]] = []
for a in self.assets:
if not a.get("approved", True):
continue
chs = a.get("chunks") or []
if chs:
for c in chs:
units.append({
"assetId": a["assetId"], "title": a["title"], "kind": a["kind"],
"version": a["version"], "tags": a.get("tags") or [],
"content": c.get("text") or "",
"chunkId": c.get("chunkId"),
"heading": c.get("heading"), "page": c.get("page"),
})
else:
units.append({
"assetId": a["assetId"], "title": a["title"], "kind": a["kind"],
"version": a["version"], "tags": a.get("tags") or [],
"content": a.get("content") or "",
"chunkId": None, "heading": None, "page": None,
})
return units
2026-07-21 11:05:57 +08:00
def ensure_seed_assets(self) -> int:
"""幂等补种:按标题补齐缺失的种子资产(旧 knowledge.json 不丢用户数据)。返回新增条数。"""
have = {a["title"] for a in self.assets}
added = 0
with self._lock:
for seed in _seed_assets():
if seed["title"] in have:
continue
# 重新发号,避免与现有 assetId 冲突
seed = dict(seed)
seed["assetId"] = uuid.uuid4().hex[:10]
self.assets.append(seed)
have.add(seed["title"])
added += 1
if added:
self._write()
return added
# ---------------- 按租户隔离的实例 ----------------
_stores: dict[str, KnowledgeStore] = {}
_stores_lock = threading.Lock()
2026-07-21 11:05:57 +08:00
def get_knowledge() -> KnowledgeStore:
"""Return the current tenant's knowledge store."""
from server.auth.context import get_identity
from server.state.store import world_path_for
tenant_uuid = get_identity().tenant_uuid
with _stores_lock:
store = _stores.get(tenant_uuid)
if store is None:
if tenant_uuid == "platform":
path = os.environ.get(_PATH_ENV, "server/data/knowledge.json")
else:
world_path = world_path_for("knowledge", tenant_uuid)
path = os.path.join(os.path.dirname(os.path.dirname(world_path)), "knowledge.json")
store = KnowledgeStore(path)
store.ensure_seed_assets()
_stores[tenant_uuid] = store
return store