接续 88dd570,保存 W20–W24 已实现的共享接口、业务入口、迁移、工作台、测试与文档。 W20/W22/W23 保持 in_progress,W21/W24 保持 verified;此提交不宣称方法或规则正式启用、多轮返修、真实角色评测完成。 W25 新增实验、标定、逐调用交付与角色执行及其迁移/测试/索引留在实施工作树,原有私人和旧实现保留项不纳入。 验证:离线 571、前端 33 通过;PG 469 项通过、2 项浏览器未启用,2 项误带入的 W25 用例已移出本提交;最终任务与交付边界 37 项通过。make 检查、最终类型、84 项资源及 diff 检查通过。未重跑浏览器或 Pi 宿主,不以合成调用认证外部模型效果。 独立整体审查四维通过;证据保存在 R2-20260909/提交W20-W24。
410 lines
18 KiB
Python
410 lines
18 KiB
Python
"""隔离目标初始化与身份门禁;目标身份不是来源声明或业务写入授权。"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import hashlib
|
||
import json
|
||
import re
|
||
from collections import Counter
|
||
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 _允许配置(路径: Path) -> dict:
|
||
值 = _配置文件(路径)
|
||
if not {"system_identifier", "data_directory", "socket", "port"} <= set(值):
|
||
raise 迁移错误("target_denied", "受控实例记录缺少指纹或端点")
|
||
if (
|
||
not isinstance(值["port"], int)
|
||
or isinstance(值["port"], bool)
|
||
or not 1 <= 值["port"] <= 65535
|
||
or not isinstance(值["socket"], str)
|
||
or not 值["socket"].startswith("/")
|
||
or not isinstance(值["data_directory"], str)
|
||
or not 值["data_directory"].startswith("/")
|
||
or not re.fullmatch(r"[0-9]+", str(值["system_identifier"]))
|
||
):
|
||
raise 迁移错误("target_denied", "实例记录不是明确的本地临时PG指纹")
|
||
return 值
|
||
|
||
|
||
def _连接配置(引用: 数据库引用, 允许: dict, *, 角色: str | None = None) -> str:
|
||
串 = 解析连接串(引用)
|
||
try:
|
||
参数 = conninfo_to_dict(串)
|
||
if (
|
||
参数.get("host") != 允许["socket"]
|
||
or str(参数.get("port", "5432")) != str(允许["port"])
|
||
or not str(参数.get("host", "")).startswith("/")
|
||
or not 参数.get("dbname")
|
||
or not 参数.get("user")
|
||
or (角色 is not None and 参数["user"] != 角色)
|
||
):
|
||
raise 迁移错误("target_denied", "连接端点或角色不属于获准临时实例")
|
||
return 串
|
||
except psycopg.Error:
|
||
raise 迁移错误("target_denied", "目标连接配置无效") from None
|
||
|
||
|
||
def _实例一致(连, 允许: dict) -> None:
|
||
目录 = 连.execute("SHOW data_directory").fetchone()[0]
|
||
系统号 = str(
|
||
连.execute("SELECT system_identifier FROM pg_catalog.pg_control_system()").fetchone()[0]
|
||
)
|
||
if 系统号 != 允许["system_identifier"] or str(目录) != 允许["data_directory"]:
|
||
raise 迁移错误("target_denied", "实际实例与受控创建/核验记录不一致")
|
||
|
||
|
||
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
|
||
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
|
||
维护 = 数据库工厂(数据库引用("受控存储", 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]
|
||
连.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"],
|
||
允许["socket"],
|
||
允许["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(
|
||
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"])
|
||
# 提前检查本批及显式依赖的既有映射,不能先写父源再发现子源漂移。
|
||
for s in [*快照, *已存源.values()]:
|
||
旧行 = 台账.读取(s.source.记录键)
|
||
if 旧行 is not None and (
|
||
旧行["source_hash"] != s.源哈希 or 旧行["mapping_hash"] != mapping_hash
|
||
):
|
||
raise 迁移错误("drift", "本批来源或明确依赖的映射已改变,整批未写入")
|
||
分流 = 分流批次(快照, 映射, 已存正式源=已存源)
|
||
if any(r.reason_code == "drift" for r in 分流):
|
||
raise 迁移错误("drift", "同批同一来源身份存在不同内容,整批不写入目标")
|
||
有序 = sorted(zip(快照, 分流, strict=True), key=lambda x: not 正式旧表(x[0]))
|
||
for 源, 判定 in 有序:
|
||
行 = 台账.占用(源, mapping_hash, 判定)
|
||
key, 命令ID = 源.source.记录键, 行["command_id"]
|
||
if 行["state"] in {"historical", "quarantined"}:
|
||
continue
|
||
if 行["state"] == "mapped":
|
||
if 冻结模式(行["frozen_request"]) == "link_confirmed":
|
||
回执, 目标 = 台账.核对关联(key)
|
||
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")
|
||
)
|
||
台账.保存终态(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": 错误,
|
||
"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", "socket"),
|
||
("port", "port"),
|
||
)
|
||
):
|
||
raise 迁移错误("target_denied", "目标回执、实例或库内标记不一致")
|
||
return 回执
|
||
except psycopg.Error as 错:
|
||
raise 迁移错误(
|
||
"target_denied", "隔离目标核对失败,SQLSTATE=" + (错.sqlstate or "unknown")
|
||
) from None
|