zizi c154ca9085 配置与数据库:迁移 V0056–V0060 与流程模板
- 迁移:调用结算与配置验证(V0056)、来源当前授权(V0057)、成品补充验证回执(V0058)、检索索引代次(V0059)、
  评测单元复用来源(V0060);新表按既有约定加只追加守卫与角色授权。
- 流程模板与载入口径同步(内联保护字段改为登记类型约束);旧库迁移工具链按新表结构对齐。
- 配置:提供方模板与运行配置同步角色策略版本;角色策略白名单新增 qwen3.8-flash(见收尾报告待裁决项:
  该模型精确身份与独立性尚未核验,且策略版本号未随白名单升版)。
2026-09-18 01:15:25 +08:00

567 lines
25 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.

"""隔离目标初始化与身份门禁;目标身份不是来源声明或业务写入授权。"""
from __future__ import annotations
import hashlib
import json
import os
import re
from collections import Counter
from ipaddress import ip_address
from pathlib import Path
from uuid import UUID, uuid4
import psycopg
from psycopg import sql
from psycopg.conninfo import conninfo_to_dict, make_conninfo
from pydantic import ValidationError
from muse.共享.调用身份 import 用途, 调用身份
from muse.共享.错误 import Muse错误
from muse.启动 import 应用装配
from muse.基础设施.数据库.迁移 import 执行迁移
from muse.基础设施.数据库.连接 import 数据库工厂, 解析连接串
from muse.正式变更.接口 import 变更错误
from muse.配置 import 数据库引用
from 关联已确认 import 关联源指针, 冻结模式, 正式旧表
from 建立映射 import 稳定JSON, 读取JSON, 旧快照, 映射配置, 检查目标快照, 迁移错误
from 映射台账 import 迁移台账
from 知识记录分流 import 分流批次
from 转换业务对象 import 准备目标请求, 提交目标请求, 核对目标产物
回执字段 = {
"version",
"target_id",
"system_identifier",
"data_directory",
"host",
"port",
"database",
"database_oid",
"role",
"run_purpose",
"created_at",
}
def _配置文件(路径: Path) -> dict:
with 路径.open("rb") as 文件:
文本 = 文件.read(16385)
if len(文本) > 16384:
raise 迁移错误("target_denied", "隔离配置或回执超过16KiB")
值 = 读取JSON(文本.decode())
if not isinstance(值, dict):
raise 迁移错误("target_denied", "隔离配置或回执必须是对象")
return 值
def _端点配置(值: dict) -> dict:
# 旧测试夹具只登记socket;该兼容分支不能解释为TCP授权。
if "endpoint_type" not in 值 and "socket" in 值 and "host" not in 值:
值 = {**值, "endpoint_type": "unix", "host": 值["socket"]}
类型, 主机, 端口 = 值.get("endpoint_type"), 值.get("host"), 值.get("port")
if (
类型 not in ("unix", "tcp")
or not isinstance(主机, str)
or not 主机
or "," in 主机
or any(c.isspace() for c in 主机)
or not isinstance(端口, int)
or isinstance(端口, bool)
or not 1 <= 端口 <= 65535
or ("socket" in 值 and (类型 != "unix" or 值["socket"] != 主机))
):
raise 迁移错误("target_denied", "实例记录必须明确登记单一端点类型、地址和端口")
if 类型 == "unix":
if not 主机.startswith("/"):
raise 迁移错误("target_denied", "unix端点必须是绝对socket目录")
else:
try:
地址 = ip_address(主机)
if 地址.is_unspecified or 地址.is_multicast or "%" in 主机:
raise ValueError("不是单一实例地址")
except ValueError:
raise 迁移错误("target_denied", "tcp端点必须是明确批准的单一IP地址") from None
return 值
def _允许配置(路径: Path) -> dict:
值 = _端点配置(_配置文件(路径))
if (
not isinstance(值.get("data_directory"), str)
or not 值["data_directory"].startswith("/")
or not isinstance(值.get("system_identifier"), str)
or not re.fullmatch(r"[0-9]+", 值["system_identifier"])
):
raise 迁移错误("target_denied", "受控实例记录缺少实测系统标识或绝对数据目录")
return 值
# 连接参数不提供第二条寻址、会话角色或服务配置通道。
连接参数白名单 = {
"host",
"port",
"dbname",
"user",
"password",
"passfile",
"connect_timeout",
"application_name",
"sslmode",
"sslrootcert",
"sslcert",
"sslkey",
"sslpassword",
"sslcrl",
"sslcrldir",
"channel_binding",
"gssencmode",
"keepalives",
"keepalives_idle",
"keepalives_interval",
"keepalives_count",
}
禁止连接环境 = ("PGHOSTADDR", "PGSERVICE", "PGSERVICEFILE", "PGOPTIONS", "PGPORT")
def _连接配置(引用: 数据库引用, 允许: dict, *, 角色: str | None = None) -> str:
允许 = _端点配置(允许)
串 = 解析连接串(引用)
try:
参数 = conninfo_to_dict(串)
库名 = 参数.get("dbname", "")
if (
set(参数) - 连接参数白名单
or any(os.environ.get(k) for k in 禁止连接环境)
or 参数.get("host") != 允许["host"]
or 参数.get("port", "5432") != str(允许.get("port", "5432"))
or not isinstance(库名, str)
or not 库名
or "=" in 库名
or 库名.startswith(("postgresql://", "postgres://"))
or not 参数.get("user")
or (角色 is not None and 参数["user"] != 角色)
):
raise 迁移错误("target_denied", "连接端点、参数或角色不属于获准实例")
return 串
except psycopg.Error:
raise 迁移错误("target_denied", "目标连接配置无效") from None
def _实例指纹(连) -> dict:
return {
"data_directory": str(连.execute("SHOW data_directory").fetchone()[0]),
"system_identifier": str(
连.execute("SELECT system_identifier FROM pg_catalog.pg_control_system()").fetchone()[0]
),
}
def _实例一致(连, 允许: dict) -> None:
if any(值 != 允许[k] for k, 值 in _实例指纹(连).items()):
raise 迁移错误("target_denied", "实际实例与受控创建/核验记录不一致")
def 检查实例(*, 端点记录: Path, 管理引用: 数据库引用) -> dict:
"""只读实例和数据库目录身份;实测结果不自动成为迁移授权。"""
端点 = _端点配置(_配置文件(端点记录))
串 = _连接配置(管理引用, 端点)
try:
with psycopg.connect(串, connect_timeout=5) as 连:
连.read_only = True
指纹 = _实例指纹(连)
版本行 = 连.execute("SHOW server_version").fetchone()
当前行 = 连.execute("SELECT current_database(), current_user").fetchone()
if 版本行 is None or 当前行 is None:
raise 迁移错误("target_inspect_failed", "实例目录身份不完整")
版本 = 版本行[0]
当前库, 当前角色 = 当前行
数据库 = [
{"database": 名, "database_oid": int(oid), "owner": owner, "is_template": 模板}
for 名, oid, owner, 模板 in 连.execute(
"SELECT datname, oid, pg_catalog.pg_get_userbyid(datdba), datistemplate "
"FROM pg_catalog.pg_database ORDER BY datname"
).fetchall()
]
return {
"instance": {k: 端点[k] for k in ("endpoint_type", "host", "port")} | 指纹,
"server_version": 版本,
"current_database": 当前库,
"current_role": 当前角色,
"databases": 数据库,
"migration_authorized": False,
}
except psycopg.Error as 错:
raise 迁移错误(
"target_inspect_failed", "只读实例检查失败,SQLSTATE=" + (错.sqlstate or "unknown")
) from None
def _读取标记(连) -> dict:
行 = 连.execute("""SELECT jsonb_build_object(
'version',1,'target_id',target_id,'system_identifier',system_identifier,
'data_directory',data_directory,'host',host,'port',port,
'database',database_name,'database_oid',database_oid,
'role',role_name,'run_purpose',run_purpose,'created_at',created_at)
FROM public.muse_legacy_target WHERE singleton""").fetchone()
if 行 is None:
raise 迁移错误("target_denied", "目标没有初始化标记")
回执 = 行[0]
# PostgreSQL的OID转JSON时为字符串;接口统一成整数与目录查询比较。
回执["database_oid"] = int(回执["database_oid"])
return 回执
def 初始化目标(
*,
允许实例: Path,
管理引用: 数据库引用,
维护模板: 数据库引用,
应用模板: 数据库引用,
输出目录: Path,
) -> dict:
"""只在获准实例中新建随机目标;不接收源库或任意既有数据库名。"""
允许 = _允许配置(允许实例)
管理串 = _连接配置(管理引用, 允许)
维护串 = _连接配置(维护模板, 允许, 角色="muse_maint")
应用串 = _连接配置(应用模板, 允许, 角色="muse_app")
target_id = uuid4()
库名 = "muse_migration_" + target_id.hex
输出目录.mkdir(mode=0o700)
维护文件, 应用文件, 回执文件 = [
输出目录 / n for n in ("维护连接.txt", "应用连接.txt", "目标回执.json")
]
已建库, 成功 = False, False
新库OID = None
try:
# 先落私有连接引用,进程被终止时仍可定位本次拟创建的随机库;缺最终回执不得使用。
for 文件, 串 in ((维护文件, 维护串), (应用文件, 应用串)):
with 文件.open("x", encoding="utf-8") as f:
文件.chmod(0o600)
f.write(make_conninfo(串, dbname=库名))
with psycopg.connect(管理串, autocommit=True, connect_timeout=5) as 管理:
_实例一致(管理, 允许)
管理.execute(
sql.SQL("CREATE DATABASE {} OWNER muse_maint").format(sql.Identifier(库名))
)
已建库 = True
新库行 = 管理.execute(
"SELECT oid FROM pg_catalog.pg_database WHERE datname=%s", (库名,)
).fetchone()
if 新库行 is None:
raise 迁移错误("target_denied", "无法取得本次新库OID,保留定位文件")
新库OID = 新库行[0]
维护 = 数据库工厂(数据库引用("受控存储", str(维护文件)), 用途.维护)
with 维护.连接() as 连:
执行迁移(连, Path(__file__).resolve().parents[1] / "迁移")
with 维护.连接() as 连, 连.transaction():
目录行 = 连.execute(
"SELECT oid FROM pg_catalog.pg_database WHERE datname=current_database()"
).fetchone()
if 目录行 is None:
raise 迁移错误("target_denied", "无法核对新目标数据库身份")
oid = 目录行[0]
if oid != 新库OID:
raise 迁移错误("target_denied", "新目标数据库身份已改变")
连.execute(
"""INSERT INTO public.muse_legacy_target
(target_id,system_identifier,data_directory,host,port,database_name,database_oid,role_name,run_purpose)
VALUES (%s,%s,%s,%s,%s,%s,%s,'muse_app','production')""",
(
target_id,
允许["system_identifier"],
允许["data_directory"],
允许["host"],
允许["port"],
库名,
oid,
),
)
回执 = _读取标记(连)
with 回执文件.open("x", encoding="utf-8") as f:
回执文件.chmod(0o600)
json.dump(回执, f, ensure_ascii=False, indent=2)
f.write("\n")
成功 = True
return 回执
except psycopg.Error as 错:
raise 迁移错误(
"target_init_failed", "隔离目标初始化失败,SQLSTATE=" + (错.sqlstate or "unknown")
) from None
finally:
if not 成功:
if 已建库:
with psycopg.connect(管理串, autocommit=True, connect_timeout=5) as 管理:
_实例一致(管理, 允许)
实际 = 管理.execute(
"SELECT oid FROM pg_catalog.pg_database WHERE datname=%s", (库名,)
).fetchone()
if (
新库OID is None
or 实际 is None
or 实际[0] != 新库OID
or 库名 != "muse_migration_" + target_id.hex
):
raise 迁移错误(
"target_denied", "本次新库身份缺失或改变,保留定位文件不清理"
)
管理.execute(
sql.SQL("DROP DATABASE {} WITH (FORCE)").format(sql.Identifier(库名))
)
for 文件 in (维护文件, 应用文件, 回执文件):
文件.unlink(missing_ok=True)
输出目录.rmdir()
def _原回执(正式, 身份, 命令ID):
try:
return 正式.读取回执(身份, 命令ID)
except 变更错误 as 错:
if 错.错误码 == "RECEIPT_NOT_FOUND":
return None
raise
def 导入快照(
装配: 应用装配,
身份: 调用身份,
快照: list[旧快照],
映射: 映射配置,
*,
允许实例: Path,
管理引用: 数据库引用,
目标回执: Path,
文件回执: Path | None = None,
) -> dict:
"""重放真实S01命令补齐映射;终态台账不是业务对象存在的替代证明。"""
检查目标快照(快照)
数据库 = 数据库工厂(装配.配置.数据库, 装配.配置.运行用途)
核对隔离目标(数据库, 允许实例=允许实例, 管理引用=管理引用, 目标回执=目标回执)
if (
not 身份.允许写正式内容
or not 身份.作者
or any(r.get("author_id") != 身份.作者 for r in 映射.值["works"])
):
raise 迁移错误("target_denied", "目标作者必须来自已认证配置,不能采用映射中的其他作者")
from 读取旧文件 import 核对文件保全
文件检查 = 核对文件保全(文件回执, 快照)
if not 文件检查["passed"]:
raise 迁移错误("file_copy_unverified", "文件保全尚未核对通过;未写入台账或业务")
正式 = 装配.要求知识方法().正式
台账 = 迁移台账(数据库, 身份, 正式, 装配.要求故事世界(), 装配.要求审校(), 装配)
mapping_hash = hashlib.sha256(稳定JSON(映射.值).encode()).hexdigest()
已存源 = {}
输入键 = {s.source.记录键 for s in 快照}
for s in 快照:
指针 = 关联源指针(s, 映射)
if 指针 and 指针["source_key"] not in 输入键:
旧行 = 台账.读取(指针["source_key"])
if 旧行 is not None:
已存源[指针["source_key"]] = 旧快照.从载荷(旧行["original"])
from PG作品映射 import PG依赖指针, PG排序, 是PG内容
from PG台账接线 import PG分片等待, 执行PG行
if 映射.值.get("pg_works") and 映射.值.get("pg_author_id") != 身份.作者:
raise 迁移错误("target_denied", "PG清单作者与应用配置不同")
指针表 = {}
for s in 快照:
if 是PG内容(s):
try:
for p in PG依赖指针(映射, s):
指针表[p["source_key"]] = p
except 迁移错误:
continue # 纯分流逐条保存缺映射/漂移原因。
for key in 指针表.keys() - 输入键 - 已存源.keys():
旧行 = 台账.读取(key)
if 旧行 is not None:
已存源[key] = 旧快照.从载荷(旧行["original"])
# 提前检查本批及显式依赖的既有映射,不能先写父源再发现子源漂移。
for s in [*快照, *已存源.values()]:
旧行 = 台账.读取(s.source.记录键)
if 旧行 is not None and (
旧行["source_hash"] != s.源哈希 or 旧行["mapping_hash"] != mapping_hash
):
raise 迁移错误("drift", "本批来源或明确依赖的映射已改变,整批未写入")
# 后续分片可推动先前reserved的PG父源与成员;不要求重传整书正文。
待处理 = [*快照, *(s for k, s in 已存源.items() if k not in 输入键 and 是PG内容(s))]
分流 = 分流批次(待处理, 映射, 已存正式源=已存源)
if any(r.reason_code == "drift" for r in 分流):
raise 迁移错误("drift", "同批同一来源身份存在不同内容,整批不写入目标")
可用源 = {**已存源, **{s.source.记录键: s for s in 快照}}
有序 = sorted(
zip(待处理, 分流, strict=True),
key=lambda x: (*PG排序(映射, x[0]), not 正式旧表(x[0])),
)
# 聚合正文必须先逐源留存,避免首块成功而其余块没有原行台账。
for 源, 判定 in 有序:
if 是PG内容(源):
台账.占用(源, mapping_hash, 判定)
for 源, 判定 in 有序:
行 = 台账.占用(源, mapping_hash, 判定)
key, 命令ID = 源.source.记录键, 行["command_id"]
if 行["state"] in {"historical", "quarantined"}:
continue
if 判定.type_id.startswith("pg_") and 判定.disposition == "待转换":
执行PG行(装配, 身份, 台账, 源, 判定, 映射, 可用源, 行)
continue
if 行["state"] == "mapped":
if 冻结模式(行["frozen_request"]) == "link_confirmed":
回执, 目标 = 台账.核对关联(key)
elif 行["frozen_request"].get("quality_kind") == "rule_version":
回执 = 行["receipt"]
目标 = 核对目标产物(装配, 身份, 行["frozen_request"], 回执)
else:
回执 = 正式.读取回执(身份, 命令ID)
目标 = 核对目标产物(装配, 身份, 行["frozen_request"], 回执)
if 稳定JSON(目标) != 稳定JSON(行["targets"]) or 稳定JSON(回执) != 稳定JSON(
行["receipt"]
):
raise 迁移错误("target_mismatch", "已有映射与真实目标/回执不一致,不重建或覆盖目标")
continue
if 判定.disposition != "待转换":
状态 = "historical" if 判定.disposition == "历史保留" else "quarantined"
台账.保存终态(key, 状态, 原因=判定.reason_code)
continue
关联 = 判定.reason_code == "linked_existing_confirmation"
if 关联:
指针 = 判定.content["link"]["target_source"]
if 指针["source_key"] != key:
根 = 台账.读取(指针["source_key"])
if 根 is not None and 根["state"] == "reserved":
raise 迁移错误(
"legacy_dependency_pending", "正式源尚在处理中,保留reserved供重试"
)
if 根 is None or 根["state"] != "mapped":
台账.保存终态(key, "quarantined", 原因="legacy_target_source_missing")
continue
回执 = None if 关联 else _原回执(正式, 身份, 命令ID)
冻结 = 行["frozen_request"]
if 回执 is not None and 冻结 is None:
raise 迁移错误("command_collision", "迁移未冻结请求但同命令已有业务回执")
if 冻结 is None:
try:
冻结 = 准备目标请求(装配, 身份, 源, 判定, 映射, 可用源=可用源)
except (迁移错误, Muse错误, ValidationError) as 错:
if isinstance(错, Muse错误) and 错.可重试:
raise
代码 = (
错.code
if isinstance(错, 迁移错误)
else getattr(错, "错误码", "target_request_invalid")
)
if 代码 == "target_owner_read_unavailable":
raise
台账.保存终态(key, "quarantined", 原因=代码)
continue
冻结 = 台账.冻结请求(key, 冻结)
if 冻结模式(冻结) == "link_confirmed":
回执, 目标 = 台账.核对关联(key)
台账.保存终态(key, "mapped", 目标=目标, 回执=回执)
continue
if 回执 is None:
try:
回执 = 提交目标请求(装配, 身份, 命令ID, 冻结)
except (迁移错误, Muse错误, ValidationError) as 错:
if isinstance(错, Muse错误) and 错.可重试:
raise
回执 = _原回执(正式, 身份, 命令ID)
if 回执 is None:
代码 = (
错.code
if isinstance(错, 迁移错误)
else getattr(错, "错误码", "target_request_invalid")
)
台账.保存终态(key, "quarantined", 原因=代码)
continue
# 目标存在性核对或台账保存失败时保留reserved,下一次读取真实回执恢复,不伪判未写。
目标 = 核对目标产物(装配, 身份, 冻结, 回执)
台账.保存终态(key, "mapped", 目标=目标, 回执=回执)
记录 = 台账.列出([x.source.记录键 for x in 快照])
统计 = dict(Counter(r["state"] for r in 记录))
错误 = [
{"source_key": r["source_key"], "reason_code": r["reason_code"]}
for r in 记录
if r["state"] == "quarantined"
]
return {
"source_count": len(快照),
"distinct_source_count": len({x.source.记录键 for x in 快照}),
"states": 统计,
"errors": 错误,
"fragment_wait": PG分片等待(台账, 映射, 可用源, 待处理),
"mapping_hash": mapping_hash,
"all_accounted": len(记录) == len({x.source.记录键 for x in 快照})
and not 统计.get("reserved"),
"production_switch_authorized": False,
"package_acceptance_claimed": False,
"file_checks": 文件检查,
"history_links": [
{
"source_key": r["source_key"],
"kind": "linked_existing_confirmation",
"legacy_draft_id": r["route"]["content"]["link"]["legacy_draft_id"],
"legacy_target_id": r["route"]["content"]["link"]["legacy_target_id"],
"decision_ref": r["frozen_request"]["payload"]["decision_ref"],
"target_ref": r["targets"][0]["id"],
"revision": r["targets"][0]["revision"],
"receipt_id": r["receipt"]["receipt_id"],
}
for r in 记录
if r["state"] == "mapped" and 冻结模式(r["frozen_request"]) == "link_confirmed"
],
}
def 核对隔离目标(
数据库: 数据库工厂, *, 允许实例: Path, 管理引用: 数据库引用, 目标回执: Path
) -> dict:
"""实际管理实例、应用端点/库OID/角色与只读标记必须同时吻合。"""
if 数据库.用途 is not 用途.生产:
raise 迁移错误("target_denied", "S01迁移候选仅使用隔离库的production用途")
允许, 回执 = _允许配置(允许实例), _配置文件(目标回执)
if (
set(回执) != 回执字段
or 回执["version"] != 1
or not re.fullmatch(r"muse_migration_[0-9a-f]{32}", str(回执["database"]))
):
raise 迁移错误("target_denied", "目标回执不完整或不是新建随机库")
try:
if "muse_migration_" + UUID(回执["target_id"]).hex != 回执["database"]:
raise ValueError("目标身份不一致")
except (ValueError, TypeError, AttributeError):
raise 迁移错误("target_denied", "目标身份与随机库名不一致") from None
管理串 = _连接配置(管理引用, 允许)
应用串 = _连接配置(数据库.引用, 允许, 角色="muse_app")
if conninfo_to_dict(应用串)["dbname"] != 回执["database"]:
raise 迁移错误("target_denied", "应用连接不属于登记的目标库")
try:
with psycopg.connect(管理串, connect_timeout=5) as 管理:
管理.read_only = True
_实例一致(管理, 允许)
with 数据库.连接(只读=True) as 连:
目录行 = 连.execute(
"SELECT oid,current_user FROM pg_catalog.pg_database "
"WHERE datname=current_database()"
).fetchone()
if 目录行 is None:
raise 迁移错误("target_denied", "无法核对目标数据库身份")
oid, role = 目录行
if oid != 回执["database_oid"] or role != "muse_app":
raise 迁移错误("target_denied", "实际目标库或角色已改变")
if _读取标记(连) != 回执 or any(
回执[k] != 允许[a]
for k, a in (
("system_identifier", "system_identifier"),
("data_directory", "data_directory"),
("host", "host"),
("port", "port"),
)
):
raise 迁移错误("target_denied", "目标回执、实例或库内标记不一致")
return 回执
except psycopg.Error as 错:
raise 迁移错误(
"target_denied", "隔离目标核对失败,SQLSTATE=" + (错.sqlstate or "unknown")
) from None