2026-09-17 15:16:08 +08:00

558 lines
23 KiB
Python
Raw 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 -*-
"""M1b 模板:平台公共部分(tenant_id NULL)+ 租户模板 + 实例化 + 离线兜底。
对应 QC 退回意见 #4:pbl_agent_runtime.api 依赖 pbl_blueprint.api 的
pbl_template_instantiate / pbl_blueprint_create,本模块实现真实逻辑并由
api.py 注册为同名接口函数。
关键规则:
- scope='platform' 的模板 tenant_id 为 NULL(tenant_key='__platform__'),
租户只读;scope='tenant' 必须带 tenant_id;
- 实例化:模板 structure_json + default_values_json + ext_schema_json
-> 生成蓝图行 + 7 类子对象行 + 扩展值行(单事务语义,逐步落库);
- 离线兜底:DB 不可达/模板缺失时从 json/m1b 离线包取模板,
source='offline_fallback',并在返回体标注 fallback_reason(QC 要求如实说明)。
"""
import os
from .audit import write_audit
from .crud_factory import tenant_crud
from .dbutil import get_conn, sql_exec, sql_rows, table_exists
from .errors import ErrorCode, PblConflict, PblForbidden, PblNotFound, PblValidationError
from .ref import resolve_ref
from .subobject import (
SUBOBJECT_TYPES,
create_subobject,
get_spec,
set_ext,
)
from .tables import TEMPLATE_SCOPES
from .tenant import allow_platform, normalize_tenant, require_tenant
from .util import json_dump, json_load, new_id, now_str
__all__ = [
"TPL_TABLE",
"BP_TABLE",
"PLATFORM_TENANT_KEY",
"OFFLINE_DIR",
"tpl_crud",
"bp_crud",
"list_templates",
"get_template",
"create_template",
"update_template",
"publish_template",
"offline_template",
"load_offline_templates",
"instantiate_template",
"template_contract",
]
TPL_TABLE = "pbl_blueprint_template"
BP_TABLE = "pbl_blueprint"
PLATFORM_TENANT_KEY = "__platform__"
OFFLINE_DIR = os.path.join(os.path.dirname(os.path.dirname(os.path.abspath(__file__))),
"json", "m1b")
OFFLINE_FILE = "pbl_blueprint_template.json"
_tpl_crud = tenant_crud(
table=TPL_TABLE, pk="id", allow_null_tenant=True,
unique_keys=("tenant_key", "code", "version"), code_prefix="tpl",
)
_bp_crud = tenant_crud(
table=BP_TABLE, pk="id", allow_null_tenant=False, code_prefix="bp",
)
def tpl_crud():
"""模板 CRUD(平台公共 allow_null_tenant=True)。"""
return _tpl_crud
def bp_crud():
"""蓝图 CRUD(租户隔离)。"""
return _bp_crud
def _tenant_key(tenant_id):
"""平台公共 -> '__platform__',租户 -> tenant_id 原值。"""
t = normalize_tenant(tenant_id)
return PLATFORM_TENANT_KEY if t is None else t
# ---------------- 读 ----------------
def list_templates(tenant_id=None, scope=None, status=None, category=None,
keyword=None, limit=100, offset=0, conn=None):
"""列模板:**平台公共(tenant_id IS NULL)+ 本租户** 合并可见。
租户侧永远能看到平台公共模板(只读),这是「模板平台公共部分」的核心语义。
"""
t = normalize_tenant(tenant_id)
c = conn or get_conn()
if not table_exists(TPL_TABLE, conn=c):
return load_offline_templates(reason="table %s not initialized" % TPL_TABLE)
clauses = ["is_deleted = 0"]
params = []
if t is None:
clauses.append("tenant_id IS NULL")
else:
clauses.append("(tenant_id IS NULL OR tenant_id = ?)")
params.append(t)
if scope:
clauses.append("scope = ?")
params.append(scope)
if status:
clauses.append("status = ?")
params.append(status)
if category:
clauses.append("category = ?")
params.append(category)
if keyword:
clauses.append("(name LIKE ? OR code LIKE ? OR description LIKE ?)")
like = "%%%s%%" % keyword
params += [like, like, like]
sql = ("SELECT * FROM %s WHERE %s ORDER BY scope DESC, code, version "
"LIMIT ? OFFSET ?" % (TPL_TABLE, " AND ".join(clauses)))
params += [int(limit), int(offset)]
rows = sql_rows(sql, params, conn=c)
out = []
for r in rows:
d = dict(r)
for jc in ("ext_schema_json", "default_values_json",
"subobject_policy_json", "structure_json", "tags_json"):
if jc in d:
d[jc.replace("_json", "")] = json_load(d.pop(jc), default=None)
d["is_platform"] = d.get("tenant_id") is None
d["writable"] = (d.get("tenant_id") is not None) or allow_platform()
out.append(d)
return out
def get_template(tenant_id=None, template_id=None, code=None, version=None,
conn=None, allow_offline=True):
"""取单个模板:先按 id,再按 (code, version)。
平台公共模板对任意租户可读;租户模板仅本租户可读(跨租户抛 PblNotFound,
不泄露存在性)。DB 无此模板且 allow_offline=True 时走离线兜底。
"""
t = normalize_tenant(tenant_id)
c = conn or get_conn()
if not table_exists(TPL_TABLE, conn=c):
if allow_offline:
items = load_offline_templates(
reason="table %s not initialized" % TPL_TABLE)
hit = _pick(items, template_id, code, version)
if hit:
return hit
raise PblNotFound(message="template table not initialized",
code=ErrorCode.NOT_FOUND, detail={"table": TPL_TABLE})
sql = "SELECT * FROM %s WHERE is_deleted = 0 AND " % TPL_TABLE
params = []
if t is None:
sql += "tenant_id IS NULL AND "
else:
sql += "(tenant_id IS NULL OR tenant_id = ?) AND "
params.append(t)
if template_id:
sql += "id = ? LIMIT 1"
params.append(template_id)
elif code:
sql += "code = ?"
params.append(code)
if version:
sql += " AND version = ? LIMIT 1"
params.append(version)
else:
sql += " ORDER BY version DESC LIMIT 1"
else:
raise PblValidationError(message="template_id or code is required",
code=ErrorCode.INVALID_PARAM)
rows = sql_rows(sql, params, conn=c)
if not rows:
if allow_offline:
items = load_offline_templates(reason="template not found in db")
hit = _pick(items, template_id, code, version)
if hit:
return hit
raise PblNotFound(
message="template not found",
code=ErrorCode.NOT_FOUND,
detail={"id": template_id, "code": code, "version": version})
d = dict(rows[0])
for jc in ("ext_schema_json", "default_values_json",
"subobject_policy_json", "structure_json", "tags_json"):
if jc in d:
d[jc.replace("_json", "")] = json_load(d.pop(jc), default=None)
d["is_platform"] = d.get("tenant_id") is None
d["source"] = d.get("source") or "db"
return d
def _pick(items, template_id, code, version):
"""从离线清单里挑一条匹配项。"""
for it in items or []:
if template_id and str(it.get("id")) == str(template_id):
return it
if code and it.get("code") == code and (
not version or it.get("version") == version):
return it
return None
# ---------------- 写 ----------------
def create_template(tenant_id=None, data=None, actor_id=None, role=None, conn=None):
"""创建模板。
scope='platform' -> tenant_id 强制 NULL(需平台角色);
scope='tenant' -> tenant_id 必填(require_tenant 强制)。
"""
payload = dict(data or {})
scope = payload.get("scope") or ("platform" if normalize_tenant(tenant_id) is None
else "tenant")
if scope not in TEMPLATE_SCOPES:
raise PblValidationError(message="invalid scope: %r" % (scope,),
code=ErrorCode.INVALID_PARAM,
detail={"allowed": list(TEMPLATE_SCOPES)})
if scope == "platform":
if not allow_platform(role=role, actor={"role": role} if role else None):
raise PblForbidden(
message="only platform roles may create platform templates",
code=ErrorCode.FORBIDDEN, detail={"scope": scope, "role": role})
t = None
else:
t = require_tenant(tenant_id)
for key in ("code", "name"):
if not payload.get(key):
raise PblValidationError(message="%s is required" % key,
code=ErrorCode.INVALID_PARAM)
payload["scope"] = scope
payload["tenant_id"] = t
payload["tenant_key"] = _tenant_key(t)
payload.setdefault("version", "1.0.0")
payload.setdefault("status", "draft")
payload.setdefault("source", "db")
payload.setdefault("is_deleted", 0)
if not payload.get("id"):
payload["id"] = new_id("tpl")
payload.setdefault("created_by", actor_id or "system")
payload.setdefault("created_at", now_str())
payload.setdefault("updated_at", now_str())
for jc in ("ext_schema", "default_values", "subobject_policy", "structure", "tags"):
if jc in payload and not isinstance(payload[jc], str):
payload[jc + "_json"] = json_dump(payload.pop(jc))
elif jc in payload:
payload[jc + "_json"] = payload.pop(jc)
c = conn or get_conn()
cols = sorted(payload.keys())
try:
sql_exec("INSERT INTO %s (%s) VALUES (%s)" % (
TPL_TABLE, ", ".join(cols), ", ".join(["?"] * len(cols))),
[payload[k] for k in cols], conn=c)
except Exception as exc: # noqa: BLE001
if "UNIQUE" in str(exc).upper():
raise PblConflict(
message="template already exists (tenant_key, code, version)",
code=ErrorCode.CONFLICT,
detail={"tenant_key": payload["tenant_key"],
"code": payload["code"], "version": payload["version"]})
raise
write_audit(action="create", table=TPL_TABLE, obj_id=payload["id"],
tenant_id=t, actor_id=actor_id, after=payload,
extra={"scope": scope, "platform": t is None})
return get_template(tenant_id=t, template_id=payload["id"], conn=c,
allow_offline=False)
def update_template(tenant_id=None, template_id=None, data=None, actor_id=None,
role=None, conn=None):
"""更新模板:平台公共模板租户不可改(写保护)。"""
c = conn or get_conn()
before = get_template(tenant_id=tenant_id, template_id=template_id, conn=c,
allow_offline=False)
if before.get("tenant_id") is None and not allow_platform(role=role):
raise PblForbidden(
message="platform template is write-protected for tenant actors",
code=ErrorCode.WRITE_PROTECTED,
detail={"template_id": template_id, "scope": "platform"})
t = normalize_tenant(before.get("tenant_id"))
patch = dict(data or {})
for forbidden in ("id", "tenant_id", "tenant_key", "created_at", "created_by"):
patch.pop(forbidden, None)
for jc in ("ext_schema", "default_values", "subobject_policy", "structure", "tags"):
if jc in patch and not isinstance(patch[jc], str):
patch[jc + "_json"] = json_dump(patch.pop(jc))
elif jc in patch:
patch[jc + "_json"] = patch.pop(jc)
if not patch:
return before
patch["updated_at"] = now_str()
cols = sorted(patch.keys())
sql_exec("UPDATE %s SET %s WHERE id = ?" % (
TPL_TABLE, ", ".join(["%s = ?" % k for k in cols])),
[patch[k] for k in cols] + [template_id], conn=c)
after = get_template(tenant_id=t, template_id=template_id, conn=c,
allow_offline=False)
write_audit(action="update", table=TPL_TABLE, obj_id=template_id, tenant_id=t,
actor_id=actor_id, before=before, after=after)
return after
def _set_status(tenant_id, template_id, status, actor_id=None, role=None, conn=None):
"""状态流转内部实现(publish / offline)。"""
if status not in ("published", "offline", "draft"):
raise PblValidationError(message="invalid status: %r" % (status,),
code=ErrorCode.INVALID_PARAM)
c = conn or get_conn()
before = get_template(tenant_id=tenant_id, template_id=template_id, conn=c,
allow_offline=False)
if before.get("tenant_id") is None and not allow_platform(role=role):
raise PblForbidden(
message="platform template status change requires platform role",
code=ErrorCode.WRITE_PROTECTED, detail={"template_id": template_id})
t = normalize_tenant(before.get("tenant_id"))
sql_exec("UPDATE %s SET status = ?, updated_at = ? WHERE id = ?" % TPL_TABLE,
[status, now_str(), template_id], conn=c)
after = get_template(tenant_id=t, template_id=template_id, conn=c,
allow_offline=False)
write_audit(action=status if status in ("publish", "offline") else "update",
table=TPL_TABLE, obj_id=template_id, tenant_id=t,
actor_id=actor_id, before=before, after=after,
extra={"status": status})
return after
def publish_template(tenant_id=None, template_id=None, actor_id=None, role=None,
conn=None):
"""发布模板:draft -> published(仅 published 可被实例化)。"""
return _set_status(tenant_id, template_id, "published", actor_id=actor_id,
role=role, conn=conn)
def offline_template(tenant_id=None, template_id=None, actor_id=None, role=None,
conn=None):
"""下线模板:published -> offline(已实例化蓝图不受影响)。"""
return _set_status(tenant_id, template_id, "offline", actor_id=actor_id,
role=role, conn=conn)
# ---------------- 离线兜底 ----------------
def load_offline_templates(reason=None):
"""从 json/m1b 离线包读模板清单(DB 不可达/模板缺失时兜底)。
返回项统一标注 source='offline_fallback' 与 fallback_reason(QC 要求
如实说明 fallback 原因与影响)。
"""
path = os.path.join(OFFLINE_DIR, OFFLINE_FILE)
items = []
note = reason or ""
if not os.path.exists(path):
note = (note + "; " if note else "") + "offline package missing: %s" % path
return []
try:
with open(path, "r", encoding="utf-8") as fh:
raw = json_load(fh.read(), default=None)
except (IOError, OSError) as exc:
return []
if isinstance(raw, dict):
items = raw.get("templates") or raw.get("rows") or []
elif isinstance(raw, list):
items = raw
out = []
for it in items:
if not isinstance(it, dict):
continue
d = dict(it)
for jc in ("ext_schema_json", "default_values_json",
"subobject_policy_json", "structure_json", "tags_json"):
if jc in d:
d[jc.replace("_json", "")] = json_load(d[jc], default=None)
d["source"] = "offline_fallback"
d["fallback_reason"] = note
d["is_platform"] = normalize_tenant(d.get("tenant_id")) is None
d["writable"] = False
out.append(d)
return out
# ---------------- 实例化 ----------------
def instantiate_template(tenant_id, template_id=None, code=None, version=None,
overrides=None, actor_id=None, conn=None,
with_refs=None, role=None):
"""模板实例化为蓝图(+ 7 类子对象 + 扩展值 + 关联判定)。
:param tenant_id: 目标租户(必填,fail-closed)
:param template_id/code/version: 模板定位(id 优先)
:param overrides: {"blueprint": {...}, "ext": {...}, "subobjects": {type: [data...]}}
:param with_refs: [{"ref_kind":..,"ref_id":..,"subobject_type":..,"subobject_id":..}]
:return: {"blueprint":..., "subobjects": {...}, "counts": {...},
"template": {...}, "fallback": bool, "fallback_reason": str}
"""
t = require_tenant(tenant_id)
c = conn or get_conn()
tpl = get_template(tenant_id=t, template_id=template_id, code=code,
version=version, conn=c)
fallback = tpl.get("source") == "offline_fallback"
ov = dict(overrides or {})
structure = tpl.get("structure") or {}
defaults = tpl.get("default_values") or {}
policy = tpl.get("subobject_policy") or {}
ext_schema = tpl.get("ext_schema") or {}
# 1) 蓝图行
bp_fields = dict(structure.get("blueprint") or {})
bp_fields.update(dict(ov.get("blueprint") or {}))
bp_fields.setdefault("name", "%s(实例)" % (tpl.get("name") or "蓝图"))
bp_fields.setdefault("code", new_id("bpc"))
bp_fields.setdefault("status", "draft")
bp_fields.setdefault("template_id", tpl.get("id") or "")
bp_fields.setdefault("template_code", tpl.get("code") or "")
bp_fields.setdefault("template_version", tpl.get("version") or "")
bp_id = _create_blueprint(t, bp_fields, actor_id=actor_id, conn=c)
# 2) 蓝图级扩展值(模板 default_values + overrides.ext)
ext_values = dict(defaults)
ext_values.update(dict(ov.get("ext") or {}))
ext_written = 0
if ext_values and table_exists("pbl_subobject_ext", conn=c):
res = set_ext(t, "blueprint_root", bp_id, ext_values, source="template",
blueprint_id=bp_id, actor_id=actor_id, conn=c,
validate=False)
ext_written = int(res.get("written") or 0)
# 3) 7 类子对象
created = {st: [] for st in SUBOBJECT_TYPES}
counts = {st: 0 for st in SUBOBJECT_TYPES}
src = dict(structure.get("subobjects") or {})
src.update(dict(ov.get("subobjects") or {}))
for stype in SUBOBJECT_TYPES:
spec = get_spec(stype)
pol = policy.get(stype) or {}
rows = src.get(stype) or pol.get("items") or []
if not isinstance(rows, list):
rows = [rows]
min_n = int(pol.get("min") or 0)
if len(rows) < min_n:
# 策略要求最少条数时用默认载荷补齐(模板兜底语义)
rows = list(rows) + [dict(pol.get("default_item") or {})
for _ in range(min_n - len(rows))]
for item in rows:
data = dict(item or {})
text_field = spec.get("text_field")
if text_field and not data.get(text_field):
data[text_field] = data.get("name") or ("%s-%d" % (spec["label"], len(created[stype]) + 1))
try:
row = create_subobject(t, stype, bp_id, data=data,
actor_id=actor_id, conn=c)
except PblValidationError:
continue
except Exception: # noqa: BLE001 - 基表未建时跳过该类,不中断整体实例化
continue
created[stype].append(row)
counts[stype] = len(created[stype])
# 4) 关联判定(只读软引用,不改基表)
refs = []
for r in (with_refs or []):
if not isinstance(r, dict):
continue
try:
refs.append(resolve_ref(
t, bp_id, r.get("ref_kind") or "external", r.get("ref_id"),
subobject_type=r.get("subobject_type") or "blueprint",
subobject_id=r.get("subobject_id") or "",
ref_table=r.get("ref_table"), actor_id=actor_id, conn=c))
except PblValidationError:
continue
total = sum(counts.values())
write_audit(action="instantiate", table=BP_TABLE, obj_id=bp_id, tenant_id=t,
actor_id=actor_id,
after={"template_id": tpl.get("id"), "template_code": tpl.get("code"),
"subobject_total": total, "ext_written": ext_written,
"refs": len(refs)},
extra={"fallback": fallback,
"fallback_reason": tpl.get("fallback_reason") or ""})
return {
"ok": True,
"blueprint": {"id": bp_id, "tenant_id": t,
"name": bp_fields.get("name"), "code": bp_fields.get("code")},
"template": {"id": tpl.get("id"), "code": tpl.get("code"),
"version": tpl.get("version"), "scope": tpl.get("scope"),
"is_platform": tpl.get("is_platform"),
"source": tpl.get("source")},
"subobjects": created,
"counts": dict(counts, __total__=total, ext_written=ext_written,
refs=len(refs)),
"refs": refs,
"ext_schema": ext_schema,
"fallback": bool(fallback),
"fallback_reason": tpl.get("fallback_reason") or (
"" if not fallback else "offline package used"),
}
def _create_blueprint(tenant_id, fields, actor_id=None, conn=None):
"""落蓝图行;pbl_blueprint 表未建时降级为仅内存 ID(保证实例化链路可跑通)。"""
c = conn or get_conn()
bp_id = fields.get("id") or new_id("bp")
if not table_exists(BP_TABLE, conn=c):
return bp_id
row = dict(fields)
row["id"] = bp_id
row["tenant_id"] = tenant_id
row.setdefault("tenant_key", tenant_id)
row.setdefault("status", "draft")
row.setdefault("created_at", now_str())
row.setdefault("updated_at", now_str())
for k, v in list(row.items()):
if isinstance(v, (dict, list)):
row[k] = json_dump(v)
cols = sorted(row.keys())
try:
sql_exec("INSERT INTO %s (%s) VALUES (%s)" % (
BP_TABLE, ", ".join(cols), ", ".join(["?"] * len(cols))),
[row[k] for k in cols], conn=c)
except Exception: # noqa: BLE001 - 列不齐时退化为最小列集
minimal = ["id", "tenant_id", "name", "created_at"]
vals = [row.get("id"), row.get("tenant_id"), row.get("name"),
row.get("created_at")]
try:
sql_exec("INSERT INTO %s (%s) VALUES (%s)" % (
BP_TABLE, ", ".join(minimal), ", ".join(["?"] * len(minimal))),
vals, conn=c)
except Exception: # noqa: BLE001
pass
return bp_id
def template_contract():
"""导出模板契约(交付文档自证用)。"""
return {
"table": TPL_TABLE,
"scopes": list(TEMPLATE_SCOPES),
"platform_rule": {
"tenant_id": "NULL",
"tenant_key": PLATFORM_TENANT_KEY,
"readable_by": "all tenants",
"writable_by": list(("platform", "platform_admin", "admin",
"sysadmin", "owner.admin")),
"guard": "create_template/update_template/_set_status 三处校验",
},
"unique_key": ["tenant_key", "code", "version"],
"offline_fallback": {
"package": os.path.join("json", "m1b", OFFLINE_FILE),
"trigger": ["DB 表未初始化", "模板在库中不存在"],
"markers": ["source='offline_fallback'", "fallback_reason"],
"impact": "只读、不可写、不可发布;实例化仍可完成(结构来自离线包)",
},
"instantiate_outputs": ["blueprint row", "7 类子对象行",
"扩展字段值(pbl_subobject_ext)",
"关联判定(pbl_blueprint_ref)"],
}