From 73f96255af120a0d9558a8098192f2472a9fe574 Mon Sep 17 00:00:00 2001 From: zizi Date: Thu, 10 Sep 2026 19:25:40 +0800 Subject: [PATCH] =?UTF-8?q?W03=20=E6=95=B0=E6=8D=AE=E5=BA=93=E3=80=81?= =?UTF-8?q?=E7=94=A8=E9=80=94=E9=9A=94=E7=A6=BB=E4=B8=8E=E8=BF=81=E7=A7=BB?= =?UTF-8?q?=E6=9C=BA=E5=88=B6=EF=BC=9A=E5=8F=97=E6=8E=A7=E8=BF=9E=E6=8E=A5?= =?UTF-8?q?=E3=80=81=E7=94=A8=E9=80=94=E8=A7=92=E8=89=B2=E9=9A=94=E7=A6=BB?= =?UTF-8?q?=E3=80=81Flyway=20=E5=BC=8F=E8=BF=81=E7=A7=BB=E6=89=A7=E8=A1=8C?= =?UTF-8?q?=E4=B8=8E=E9=9A=94=E7=A6=BB=E5=BA=93=E5=A4=B9=E5=85=B7=E3=80=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 按 R2 串行阶段整理提交;包内文件为该阶段交付(含后续小增量),状态以工作包清单为准。 --- src/muse/基础设施/凭据读取.py | 50 +++++ src/muse/基础设施/数据库/__init__.py | 4 + src/muse/基础设施/数据库/事务.py | 44 +++++ src/muse/基础设施/数据库/用途隔离.py | 36 ++++ src/muse/基础设施/数据库/迁移.py | 141 ++++++++++++++ src/muse/基础设施/数据库/连接.py | 99 ++++++++++ tests/集成/test_数据库连接与迁移.py | 267 ++++++++++++++++++++++++++ 数据库/初始化/用途角色.sql | 16 ++ 数据库/迁移/V0001__共享标识与版本.sql | 64 ++++++ 9 files changed, 721 insertions(+) create mode 100644 src/muse/基础设施/凭据读取.py create mode 100644 src/muse/基础设施/数据库/__init__.py create mode 100644 src/muse/基础设施/数据库/事务.py create mode 100644 src/muse/基础设施/数据库/用途隔离.py create mode 100644 src/muse/基础设施/数据库/迁移.py create mode 100644 src/muse/基础设施/数据库/连接.py create mode 100644 tests/集成/test_数据库连接与迁移.py create mode 100644 数据库/初始化/用途角色.sql create mode 100644 数据库/迁移/V0001__共享标识与版本.sql diff --git a/src/muse/基础设施/凭据读取.py b/src/muse/基础设施/凭据读取.py new file mode 100644 index 0000000..b530ea0 --- /dev/null +++ b/src/muse/基础设施/凭据读取.py @@ -0,0 +1,50 @@ +"""受控凭据读取。 + +合同依据:.agent/约束/凭据与外部调用.md、文件设计/后端-基础设施.md。 + +- 凭据值只从环境变量或受控存储位置解析;不写入源码、配置对象或日志。 +- 缺引用即失败(环境缺失错误),不猜测、不回退示例值。 +""" + +from __future__ import annotations + +import os +from pathlib import Path + +from muse.共享.错误 import 环境缺失错误 + +_允许来源 = frozenset({"环境变量", "受控存储"}) + + +def 读取凭据(来源: str, 位置: str, *, environ: dict[str, str] | None = None) -> str: + """按声明的来源与位置解析凭据值;只返回值,不记录值。 + + :param 来源: ``环境变量`` 或 ``受控存储``。 + :param 位置: 环境变量名,或受控存储内的绝对路径。 + """ + if 来源 not in _允许来源: + raise 环境缺失错误( + f"未支持的凭据来源:{来源}", + 上下文={"允许": sorted(_允许来源)}, + ) + 环 = os.environ if environ is None else environ + if 来源 == "环境变量": + 值 = 环.get(位置, "") + if 值.strip(): + return 值 + raise 环境缺失错误( + "凭据环境变量未设置", + 上下文={"获取位置": 位置}, + ) + 路径 = Path(位置) + if not 路径.is_absolute(): + raise 环境缺失错误("受控存储位置必须是绝对路径", 上下文={"位置前缀": "已省略"}) + if not 路径.is_file(): + raise 环境缺失错误("受控存储文件不存在", 上下文={"获取位置": 位置}) + 值 = 路径.read_text(encoding="utf-8").strip() + if 值: + return 值 + raise 环境缺失错误("受控存储文件为空", 上下文={"获取位置": 位置}) + + +__all__ = ["读取凭据"] diff --git a/src/muse/基础设施/数据库/__init__.py b/src/muse/基础设施/数据库/__init__.py new file mode 100644 index 0000000..7b905e6 --- /dev/null +++ b/src/muse/基础设施/数据库/__init__.py @@ -0,0 +1,4 @@ +"""数据库基础设施:连接、事务、用途隔离与迁移。 + +边界:不内置任何正式库连接串;缺配置即失败,不回退默认库。 +""" diff --git a/src/muse/基础设施/数据库/事务.py b/src/muse/基础设施/数据库/事务.py new file mode 100644 index 0000000..5eef680 --- /dev/null +++ b/src/muse/基础设施/数据库/事务.py @@ -0,0 +1,44 @@ +"""事务 owner:使用驱动事务上下文统一处理提交、回滚及嵌套保存点。""" + +from __future__ import annotations + +from collections.abc import Iterator +from contextlib import contextmanager + +import psycopg + +from muse.共享.错误 import Muse错误 + + +class 事务失败(Muse错误): + """当前事务或嵌套范围已回滚;错误原文不进入对外说明。""" + + 错误码 = "MUSE_TX_FAILED" + + +@contextmanager +def 事务(连接: psycopg.Connection) -> Iterator[psycopg.Connection]: + """空闲连接建立事务,已有事务则建立保存点,不提交调用方外层事务。""" + try: + with 连接.transaction(): + yield 连接 + except Muse错误: + raise + except Exception as exc: + raise 事务失败( + "本次事务范围已回滚", + 上下文={ + "原因类型": type(exc).__name__, + "sqlstate": exc.sqlstate if isinstance(exc, psycopg.Error) else None, + }, + ) from None + + +@contextmanager +def 保存点(连接: psycopg.Connection, 名称: str) -> Iterator[psycopg.Connection]: + """局部失败只回滚当前范围;名称引用与嵌套边界交给驱动。""" + with 连接.transaction(savepoint_name=名称): + yield 连接 + + +__all__ = ["事务", "保存点", "事务失败"] diff --git a/src/muse/基础设施/数据库/用途隔离.py b/src/muse/基础设施/数据库/用途隔离.py new file mode 100644 index 0000000..7244c6d --- /dev/null +++ b/src/muse/基础设施/数据库/用途隔离.py @@ -0,0 +1,36 @@ +"""固定用途与登录角色;数据可见范围由迁移中的 PostgreSQL 授权控制。""" + +from __future__ import annotations + +import psycopg + +from muse.共享.调用身份 import 用途 +from muse.共享.错误 import Muse错误 + +_用途角色 = {用途.生产: "muse_app", 用途.评测: "muse_eval", 用途.维护: "muse_maint"} + + +class 用途越权(Muse错误): + """登录角色与声明用途不匹配。""" + + 错误码 = "MUSE_PURPOSE_VIOLATION" + + +def 校验用途角色(连接: psycopg.Connection, 声明用途: 用途) -> None: + """登录身份与当前角色都必须符合用途,不接受管理员切换角色代入。""" + 行 = 连接.execute("SELECT session_user, current_user").fetchone() + 预期 = _用途角色[声明用途] + if 行 is None or tuple(行) != (预期, 预期): + raise 用途越权( + f"连接身份不允许 {声明用途.value} 用途", + 上下文={"要求角色": 预期}, + ) + + +def 读取用途(连接: psycopg.Connection) -> str: + """读取会话用途供审计;此参数不产生数据库权限。""" + 行 = 连接.execute("SELECT current_setting('app.purpose', true)").fetchone() + return str(行[0] or "") if 行 else "" + + +__all__ = ["校验用途角色", "读取用途", "用途越权"] diff --git a/src/muse/基础设施/数据库/迁移.py b/src/muse/基础设施/数据库/迁移.py new file mode 100644 index 0000000..0a85e20 --- /dev/null +++ b/src/muse/基础设施/数据库/迁移.py @@ -0,0 +1,141 @@ +"""显式维护迁移:包内 SQL、版本校验、会话锁和逐版本事务。""" + +from __future__ import annotations + +import hashlib +import re +from dataclasses import dataclass +from importlib import resources +from importlib.resources.abc import Traversable + +import psycopg +from psycopg.pq import TransactionStatus + +from muse.共享.调用身份 import 用途 +from muse.共享.错误 import Muse错误 +from muse.基础设施.数据库.用途隔离 import 校验用途角色 +from muse.基础设施.数据库.连接 import 数据库工厂 +from muse.配置 import 应用配置 + +_锁键 = 0x6D757365 +_版本模式 = re.compile(r"^V(\d{4})__.+\.sql$") + + +class 迁移错误(Muse错误): + """迁移文件、版本、锁或执行失败。""" + + 错误码 = "MUSE_MIGRATION" + + +@dataclass(frozen=True, slots=True) +class 迁移项: + 版本: int + 文件名: str + 校验和: str + 内容: str + + +def 列出迁移(迁移目录: Traversable | None = None) -> list[迁移项]: + """默认读取包内迁移资源;执行内容与校验和来自同一份字节。""" + 目录 = 迁移目录 if 迁移目录 is not None else resources.files("muse").joinpath("资源", "迁移") + if not 目录.is_dir(): + raise 迁移错误("迁移资源目录不存在;需要构建并安装包含 SQL 的应用包") + 结果: list[迁移项] = [] + 已见: set[int] = set() + for 文件 in sorted(目录.iterdir(), key=lambda 项: 项.name): + if not 文件.is_file() or not 文件.name.endswith(".sql"): + continue + 匹配 = _版本模式.fullmatch(文件.name) + if not 匹配 or int(匹配.group(1)) < 1: + raise 迁移错误(f"迁移文件命名不合法:{文件.name}") + 版本 = int(匹配.group(1)) + if 版本 in 已见: + raise 迁移错误(f"迁移版本重复:V{版本:04d}") + 已见.add(版本) + 内容 = 文件.read_bytes() + 结果.append(迁移项(版本, 文件.name, hashlib.sha256(内容).hexdigest(), 内容.decode("utf-8"))) + if not 结果: + raise 迁移错误("迁移目录没有可执行的 SQL 文件") + return 结果 + + +def 已应用版本(连接: psycopg.Connection) -> dict[int, str]: + """读取当前库的迁移账本,查询本身不留下隐式事务。""" + with 连接.transaction(): + 行 = 连接.execute("SELECT to_regclass('public.muse_migration')").fetchone() + if 行 is None or 行[0] is None: + return {} + 账本 = 连接.execute( + "SELECT version, checksum FROM public.muse_migration ORDER BY version" + ).fetchall() + return {int(版本): str(校验) for 版本, 校验 in 账本} + + +def 执行迁移( + 连接: psycopg.Connection, + 迁移目录: Traversable | None = None, + *, + 目标版本: int | None = None, +) -> list[迁移项]: + """维护连接取得锁后读取账本;每版 DDL 和登记共同提交,失败可重跑。""" + 全部 = 列出迁移(迁移目录) + if 目标版本 is not None: + if 目标版本 not in {项.版本 for 项 in 全部}: + raise 迁移错误("目标版本不在迁移资源中") + 全部 = [项 for 项 in 全部 if 项.版本 <= 目标版本] + if 连接.info.transaction_status != TransactionStatus.IDLE: + raise 迁移错误("迁移需要独立的空闲连接,不能接管调用方事务") + with 连接.transaction(): + 校验用途角色(连接, 用途.维护) + 锁行 = 连接.execute("SELECT pg_try_advisory_lock(%s)", (_锁键,)).fetchone() + if not 锁行 or not 锁行[0]: + raise 迁移错误("其他迁移进程持有锁;拒绝并发执行") + try: + 已应用 = 已应用版本(连接) + 待执行 = [] + for 项 in 全部: + if 项.版本 in 已应用: + if 已应用[项.版本] != 项.校验和: + raise 迁移错误( + f"迁移 V{项.版本:04d} 校验和与账本不符", + 上下文={"文件": 项.文件名}, + ) + elif any(版本 > 项.版本 for 版本 in 已应用): + raise 迁移错误("不允许向已发布库补录低版本迁移") + else: + 待执行.append(项) + for 项 in 待执行: + try: + with 连接.transaction(): + 连接.execute(项.内容) + 连接.execute( + "INSERT INTO public.muse_migration (version, name, checksum) " + "VALUES (%s, %s, %s)", + (项.版本, 项.文件名, 项.校验和), + ) + except Exception as exc: + raise 迁移错误( + f"迁移 V{项.版本:04d} 执行失败,本版本未登记成功", + 上下文={ + "文件": 项.文件名, + "原因类型": type(exc).__name__, + "sqlstate": exc.sqlstate if isinstance(exc, psycopg.Error) else None, + }, + ) from None + return 待执行 + finally: + if not 连接.closed: + with 连接.transaction(): + 连接.execute("SELECT pg_advisory_unlock(%s)", (_锁键,)) + + +def 迁移目标库(配置: 应用配置, *, 目标版本: int | None = None) -> list[迁移项]: + """维护用例经统一用途工厂打开连接;接入只传配置与目标版本。""" + if 配置.运行用途 is not 用途.维护: + raise 迁移错误("迁移命令必须使用 maintenance 用途配置") + 工厂 = 数据库工厂(配置.数据库, 配置.运行用途) + with 工厂.连接() as 连: + return 执行迁移(连, 目标版本=目标版本) + + +__all__ = ["列出迁移", "已应用版本", "执行迁移", "迁移目标库", "迁移项", "迁移错误"] diff --git a/src/muse/基础设施/数据库/连接.py b/src/muse/基础设施/数据库/连接.py new file mode 100644 index 0000000..aa28b65 --- /dev/null +++ b/src/muse/基础设施/数据库/连接.py @@ -0,0 +1,99 @@ +"""受控连接与生命周期:凭据引用、固定用途和数据库只读设置。""" + +from __future__ import annotations + +from collections.abc import Iterator +from contextlib import contextmanager +from dataclasses import dataclass +from typing import Any + +import psycopg + +from muse.共享.调用身份 import 用途 +from muse.共享.错误 import Muse错误, 环境缺失错误 +from muse.基础设施.凭据读取 import 读取凭据 +from muse.基础设施.数据库.用途隔离 import 校验用途角色, 用途越权 +from muse.配置 import 数据库引用 + + +class 数据库连接失败(Muse错误): + """连接或会话设置失败;错误不携带连接串和驱动原文。""" + + 错误码 = "MUSE_DB_CONNECTION" + 可重试 = True + + +def 解析连接串(引用: 数据库引用, *, environ: dict[str, str] | None = None) -> str: + """连接值只由凭据 owner 解析,缺失时拒绝任何隐式默认库。""" + try: + return 读取凭据(引用.取值方式, 引用.位置, environ=environ) + except 环境缺失错误 as exc: + raise 环境缺失错误( + "数据库连接串未提供;拒绝回退任何默认库", + 上下文={"获取位置": 引用.位置, "原说明": exc.说明}, + ) from None + + +def 连接( + 引用: 数据库引用, + *, + 只读: bool = False, + 用途标记: 用途 | str = 用途.生产, + environ: dict[str, str] | None = None, + **参数: Any, +) -> psycopg.Connection: + """返回用途已校验且无活动事务的连接;只读设置不被普通 options 覆盖。""" + try: + 声明用途 = 用途(用途标记) + except ValueError: + raise 用途越权("连接必须声明 production、evaluation 或 maintenance 用途") from None + 串 = 解析连接串(引用, environ=environ) + 自动提交 = 参数.pop("autocommit", False) + 参数.setdefault("application_name", f"muse:{声明用途.value}") + 参数.setdefault("connect_timeout", 10) + 连 = None + try: + # 初始化不留下隐式事务;随后交给事务 owner 建立业务边界。 + 连 = psycopg.connect(串, autocommit=True, **参数) + 校验用途角色(连, 声明用途) + 连.execute("SELECT set_config('app.purpose', %s, false)", (声明用途.value,)) + if 只读: + 连.execute("SELECT set_config('default_transaction_read_only', 'on', false)") + 连.read_only = True + 连.autocommit = 自动提交 + return 连 + except BaseException as exc: + if 连 is not None: + 连.close() + if isinstance(exc, psycopg.Error): + raise 数据库连接失败( + "数据库连接或会话初始化失败", + 上下文={"原因类型": type(exc).__name__, "sqlstate": exc.sqlstate}, + ) from None + raise + + +@dataclass(frozen=True, slots=True) +class 数据库工厂: + """装配时冻结连接引用与用途;取得连接时才读取凭据并连接。""" + + 引用: 数据库引用 + 用途: 用途 + + def 连接(self, *, 只读: bool = False, autocommit: bool = False) -> psycopg.Connection: + return 连接(self.引用, 只读=只读, 用途标记=self.用途, autocommit=autocommit) + + def 检查连接(self) -> None: + """只读探针在数据库 owner 内执行,不把 SQL 散入接入层。""" + with self.连接(只读=True) as 连: + 连.execute("SELECT 1").fetchone() + + +@contextmanager +def 连接范围(引用: 数据库引用, **关键字: Any) -> Iterator[psycopg.Connection]: + """委托驱动完成提交或回滚及连接释放。""" + with 连接(引用, **关键字) as 连: + yield 连 + + +__all__ = ["解析连接串", "连接", "连接范围", "数据库工厂", "数据库连接失败"] diff --git a/tests/集成/test_数据库连接与迁移.py b/tests/集成/test_数据库连接与迁移.py new file mode 100644 index 0000000..74ac9a5 --- /dev/null +++ b/tests/集成/test_数据库连接与迁移.py @@ -0,0 +1,267 @@ +"""W03 的真实连接、权限、事务与迁移验证;仅使用显式隔离 PostgreSQL。""" + +from __future__ import annotations + +import hashlib +import uuid +from collections.abc import Iterator +from pathlib import Path + +import psycopg +import pytest +from psycopg import sql + +from muse.__main__ import main +from muse.共享.调用身份 import 用途 +from muse.共享.错误 import 环境缺失错误 +from muse.启动 import 构建 +from muse.基础设施.数据库.事务 import 事务, 事务失败, 保存点 +from muse.基础设施.数据库.用途隔离 import 用途越权 +from muse.基础设施.数据库.迁移 import 已应用版本, 执行迁移, 迁移错误 +from muse.基础设施.数据库.连接 import 解析连接串, 连接 +from muse.配置 import 应用配置, 数据库引用 + +迁移目录 = Path(__file__).resolve().parents[2] / "数据库" / "迁移" +角色表 = {用途.生产: "muse_app", 用途.评测: "muse_eval", 用途.维护: "muse_maint"} + + +@pytest.fixture +def 用途角色(隔离数据库URL: str) -> None: + """隔离集群管理员只创建三个普通角色,不把管理员作为应用角色。""" + with psycopg.connect(隔离数据库URL, autocommit=True) as 管理员: + 初始化文件 = 迁移目录.parent / "初始化" / "用途角色.sql" + 管理员.execute(初始化文件.read_text(encoding="utf-8")) + + +@pytest.fixture +def 临时库( + 隔离数据库URL: str, 用途角色: None, monkeypatch: pytest.MonkeyPatch +) -> Iterator[dict[用途, 数据库引用]]: + """每例独占一库;维护角色拥有目标库,退出后由隔离管理员销毁。""" + 库名 = f"muse_case_{uuid.uuid4().hex[:12]}" + with psycopg.connect(隔离数据库URL, autocommit=True) as 管理员: + 管理员.execute( + sql.SQL("CREATE DATABASE {} OWNER muse_maint TEMPLATE template0").format( + sql.Identifier(库名) + ) + ) + 引用表 = {} + try: + for 声明用途, 角色 in 角色表.items(): + 参数 = psycopg.conninfo.conninfo_to_dict(隔离数据库URL) + 参数.update(dbname=库名, user=角色) + 环境名 = f"MUSE_CASE_{声明用途.name}_URL" + monkeypatch.setenv(环境名, psycopg.conninfo.make_conninfo(**参数)) + 引用表[声明用途] = 数据库引用("环境变量", 环境名) + yield 引用表 + finally: + with psycopg.connect(隔离数据库URL, autocommit=True) as 管理员: + 管理员.execute(sql.SQL("DROP DATABASE {} WITH (FORCE)").format(sql.Identifier(库名))) + + +def 维护连接(临时库: dict[用途, 数据库引用], **参数: object) -> psycopg.Connection: + return 连接(临时库[用途.维护], 用途标记=用途.维护, **参数) + + +def test_缺配置拒绝回退默认库__d10001(monkeypatch: pytest.MonkeyPatch) -> None: + """given 缺少连接引用;when 装配连接;then 拒绝默认库回退。""" + monkeypatch.delenv("MUSE_MISSING_DATABASE_URL", raising=False) + 配置 = 应用配置(数据库引用("环境变量", "MUSE_MISSING_DATABASE_URL"), "test") + with pytest.raises(环境缺失错误, match="拒绝回退"): + 构建(配置).要求数据库().连接() + + +@pytest.mark.数据库 +@pytest.mark.parametrize( + "options", ["", "-c application_name=readonly-probe"], ids=["default", "options"] +) +def test_只读连接拒绝写__d50005(临时库: dict[用途, 数据库引用], options: str) -> None: + """given 普通连接选项;when 请求只读;then 数据库拒绝写入。""" + with 维护连接(临时库, 只读=True, options=options) as 连: + with pytest.raises(psycopg.errors.ReadOnlySqlTransaction): + 连.execute("CREATE TABLE readonly_probe (id int)") + 连.rollback() + + +@pytest.mark.数据库 +@pytest.mark.parametrize("autocommit", [False, True], ids=["manual", "auto"]) +def test_事务回滚不留半成品__d60006(临时库: dict[用途, 数据库引用], autocommit: bool) -> None: + """given 两种提交模式;when 第二参与者失败;then 全部写入回滚。""" + with 维护连接(临时库, autocommit=autocommit) as 连: + with pytest.raises(事务失败) as 记录: + with 事务(连): + 连.execute("CREATE TABLE participant_one (id int)") + 连.execute("CREATE TABLE participant_two (id int)") + raise RuntimeError("第二参与者失败,正文与凭据不得被拼入错误") + assert "正文与凭据" not in str(记录.value.呈现()) + assert ( + 连.execute( + "SELECT count(*) FROM information_schema.tables " + "WHERE table_name IN ('participant_one', 'participant_two')" + ).fetchone()[0] + == 0 + ) + + +@pytest.mark.数据库 +@pytest.mark.parametrize( + ("登录用途", "声明用途"), + [(角色, 声明) for 角色 in 用途 for 声明 in 用途 if 角色 != 声明], + ids=[f"{角色.value}-as-{声明.value}" for 角色 in 用途 for 声明 in 用途 if 角色 != 声明], +) +def test_用途角色不匹配被拒__d70007( + 临时库: dict[用途, 数据库引用], 登录用途: 用途, 声明用途: 用途 +) -> None: + """given 固定普通角色;when 正式连接入口声明其他用途;then 连接被拒。""" + with pytest.raises(用途越权): + 连接(临时库[登录用途], 用途标记=声明用途) + with 连接(临时库[登录用途], 用途标记=登录用途) as 连: + assert 连.execute("SELECT current_user").fetchone()[0] == 角色表[登录用途] + + +@pytest.mark.数据库 +@pytest.mark.parametrize("目标表", ["answer", "blind_assignment"], ids=["oracle", "blind-map"]) +def test_生产用途拒绝oracle及匿名映射读取__d80008( + 临时库: dict[用途, 数据库引用], 目标表: str +) -> None: + """given 维护角色建立答案或匿名映射;when 生产连接查询;then PG 拒权。""" + 目标 = sql.Identifier("oracle", 目标表) + with 维护连接(临时库) as 连: + 执行迁移(连, 迁移目录, 目标版本=1) + 连.execute(sql.SQL("CREATE TABLE {} (id int)").format(目标)) + 连.execute(sql.SQL("INSERT INTO {} VALUES (1)").format(目标)) + with 连接(临时库[用途.生产], 用途标记=用途.生产) as 连: + with pytest.raises(psycopg.errors.InsufficientPrivilege): + 连.execute(sql.SQL("SELECT * FROM {}").format(目标)) + 连.rollback() + with 连接(临时库[用途.评测], 用途标记=用途.评测) as 连: + assert 连.execute(sql.SQL("SELECT id FROM {}").format(目标)).fetchone()[0] == 1 + + +@pytest.mark.数据库 +@pytest.mark.parametrize("autocommit", [False, True], ids=["manual", "auto"]) +def test_保存点局部失败保留外层事务__d90009( + 临时库: dict[用途, 数据库引用], autocommit: bool +) -> None: + """given 外层事务已有写入;when 保存点失败;then 局部回滚且外层可提交。""" + with 维护连接(临时库, autocommit=autocommit) as 连: + with 事务(连): + 连.execute("CREATE TABLE savepoint_probe (id int)") + 连.execute("INSERT INTO savepoint_probe VALUES (1)") + with pytest.raises(RuntimeError): + with 保存点(连, '中文保存点"'): + 连.execute("INSERT INTO savepoint_probe VALUES (2)") + raise RuntimeError("局部失败") + 连.execute("INSERT INTO savepoint_probe VALUES (3)") + assert 连.execute("SELECT id FROM savepoint_probe ORDER BY id").fetchall() == [(1,), (3,)] + with 维护连接(临时库) as 再读: + assert 再读.execute("SELECT count(*) FROM savepoint_probe").fetchone()[0] == 2 + + +@pytest.mark.数据库 +def test_评测角色仅可写评测对象__da000a(临时库: dict[用途, 数据库引用]) -> None: + """given 生产与评测对象;when 评测连接写入;then 仅评测对象可写。""" + with 维护连接(临时库) as 连: + 执行迁移(连, 迁移目录, 目标版本=1) + 连.execute("CREATE TABLE public.production_probe (id int)") + 连.execute("CREATE TABLE evaluation.sample_probe (id int)") + with 连接(临时库[用途.生产], 用途标记=用途.生产) as 连: + 连.execute("INSERT INTO public.production_probe VALUES (1)") + with 连接(临时库[用途.评测], 用途标记=用途.评测) as 连: + with pytest.raises(psycopg.errors.InsufficientPrivilege): + 连.execute("INSERT INTO public.production_probe VALUES (2)") + 连.rollback() + 连.execute("INSERT INTO evaluation.sample_probe VALUES (3)") + assert 连.execute("SELECT id FROM evaluation.sample_probe").fetchone()[0] == 3 + with 维护连接(临时库) as 连: + assert 连.execute("SELECT id FROM public.production_probe").fetchall() == [(1,)] + + +@pytest.mark.数据库 +def test_迁移幂等与篡改拦截__d20002(临时库: dict[用途, 数据库引用], tmp_path: Path) -> None: + """given 已登记迁移;when 重跑或篡改同版本;then 幂等跳过或拒绝。""" + with 维护连接(临时库) as 连: + assert [项.版本 for 项 in 执行迁移(连, 迁移目录, 目标版本=1)] == [1] + assert len(已应用版本(连)[1]) == 64 + assert 执行迁移(连, 迁移目录, 目标版本=1) == [] + (tmp_path / "V0001__共享标识与版本.sql").write_text("-- 改动\n", encoding="utf-8") + with pytest.raises(迁移错误, match="校验和"): + 执行迁移(连, tmp_path) + + +@pytest.mark.数据库 +def test_迁移失败不登记且可重跑__d30003(临时库: dict[用途, 数据库引用], tmp_path: Path) -> None: + """given 两个版本;when 第二版本失败;then 仅第一版入账且修复可重跑。""" + (tmp_path / "V0001__共享标识与版本.sql").write_bytes( + (迁移目录 / "V0001__共享标识与版本.sql").read_bytes() + ) + V2 = tmp_path / "V0002__迁移.sql" + V2.write_text("CREATE TABLE migration_probe (id int);\n这不是合法 SQL;\n", encoding="utf-8") + with 维护连接(临时库) as 连: + with pytest.raises(迁移错误, match="未登记成功"): + 执行迁移(连, tmp_path) + assert set(已应用版本(连)) == {1} + with 连.transaction(): + assert 连.execute("SELECT to_regclass('public.migration_probe')").fetchone()[0] is None + V2.write_text("CREATE TABLE migration_probe (id int);\n", encoding="utf-8") + assert [项.版本 for 项 in 执行迁移(连, tmp_path)] == [2] + assert set(已应用版本(连)) == {1, 2} + + +@pytest.mark.数据库 +def test_低版本补录被拒绝__d40004(临时库: dict[用途, 数据库引用], tmp_path: Path) -> None: + """given 已登记较高版本;when 提供未应用低版本;then 拒绝补录。""" + with 维护连接(临时库) as 连: + 执行迁移(连, 迁移目录, 目标版本=1) + with 连.transaction(): + 连.execute( + "INSERT INTO public.muse_migration (version, name, checksum) VALUES (%s, %s, %s)", + (3, "V0003__既往.sql", hashlib.sha256(b"past").hexdigest()), + ) + (tmp_path / "V0002__补录.sql").write_text( + "CREATE TABLE backfill (id int);", encoding="utf-8" + ) + with pytest.raises(迁移错误, match="低版本"): + 执行迁移(连, tmp_path) + + +@pytest.mark.数据库 +def test_迁移先取锁再判定账本__db000b(临时库: dict[用途, 数据库引用]) -> None: + """given 一个进程持有迁移锁;when 另一连接重跑;then 不读取未保护账本。""" + with 维护连接(临时库, autocommit=True) as 持有者: + 执行迁移(持有者, 迁移目录, 目标版本=1) + 持有者.execute("SELECT pg_advisory_lock(%s)", (0x6D757365,)) + with 维护连接(临时库) as 竞争者: + with pytest.raises(迁移错误, match="持有锁"): + 执行迁移(竞争者, 迁移目录, 目标版本=1) + 持有者.execute("SELECT pg_advisory_unlock(%s)", (0x6D757365,)) + assert 执行迁移(竞争者, 迁移目录, 目标版本=1) == [] + + +@pytest.mark.数据库 +def test_迁移命令使用维护装配与包资源__dc000c( + 临时库: dict[用途, 数据库引用], tmp_path: Path, capsys: pytest.CaptureFixture[str] +) -> None: + """given 用途配置;when CLI 迁移;then 生产被拒且维护从安装资源执行。""" + 配置文件 = tmp_path / "连接.toml" + for 声明用途, 预期码 in [(用途.生产, 1), (用途.维护, 0)]: + 配置文件.write_text( + '["数据库"]\n"取值方式" = "环境变量"\n' + f'"位置" = "{临时库[声明用途].位置}"\n' + '["资源"]\n"发布身份" = "test"\n' + f'["运行"]\n"用途" = "{声明用途.value}"\n', + encoding="utf-8", + ) + assert main(["迁移", str(配置文件), "--目标版本", "1"]) == 预期码 + assert "已应用 V0001" in capsys.readouterr().out + with 维护连接(临时库) as 连: + assert set(已应用版本(连)) == {1} + + +def test_数据库引用通过唯一凭据入口读取受控文件__dd000d(tmp_path: Path) -> None: + """given 受控文件引用;when 解析连接;then 与环境引用共用读取入口。""" + 文件 = tmp_path / "合成连接.txt" + 文件.write_text("dbname=isolated_example user=muse_app\n", encoding="utf-8") + assert 解析连接串(数据库引用("受控存储", str(文件))) == ( + "dbname=isolated_example user=muse_app" + ) diff --git a/数据库/初始化/用途角色.sql b/数据库/初始化/用途角色.sql new file mode 100644 index 0000000..f498252 --- /dev/null +++ b/数据库/初始化/用途角色.sql @@ -0,0 +1,16 @@ +-- 仅由新目标环境的管理员显式运行;不在应用启动或业务迁移中创建角色。 +-- 登录认证由部署环境设置,本脚本不包含密码或其他凭据。 +DO $$ +DECLARE + purpose_role text; +BEGIN + FOREACH purpose_role IN ARRAY ARRAY['muse_app', 'muse_eval', 'muse_maint'] LOOP + IF NOT EXISTS (SELECT FROM pg_roles WHERE rolname = purpose_role) THEN + EXECUTE format( + 'CREATE ROLE %I LOGIN NOSUPERUSER NOCREATEDB NOCREATEROLE NOINHERIT NOBYPASSRLS', + purpose_role + ); + END IF; + END LOOP; +END +$$; diff --git a/数据库/迁移/V0001__共享标识与版本.sql b/数据库/迁移/V0001__共享标识与版本.sql new file mode 100644 index 0000000..2f9342c --- /dev/null +++ b/数据库/迁移/V0001__共享标识与版本.sql @@ -0,0 +1,64 @@ +-- V0001 共享标识与版本:新版数据库的共享基础。 +-- 内容:迁移账本 + 跨模块共用的身份域(用途、内容用途、候选状态、审阅动作)。 +-- 边界:只建立共享基础;业务对象表由 V0005 起各模块迁移建立, +-- 不在本版本预建业务列,也不复用旧库 DDL。 + +-- 迁移账本:版本、校验和与应用时间是数据演进的唯一事实(迁移执行器维护) +CREATE TABLE muse_migration ( + version integer PRIMARY KEY CHECK (version >= 1), + name text NOT NULL, + checksum text NOT NULL CHECK (length(checksum) = 64), + applied_at timestamptz NOT NULL DEFAULT now() +); + +-- 执行隔离用途:与 src/muse/共享/调用身份.py 的 StrEnum 一一对应 +CREATE TYPE muse_purpose AS ENUM ( + 'production', + 'evaluation', + 'maintenance' +); + +-- 内容消费用途:任务与工具按此裁剪可见字段 +CREATE TYPE muse_content_purpose AS ENUM ( + 'planning', + 'generation', + 'detection', + 'extraction' +); + +-- 候选状态:候选先审后入的状态轴(各业务模块实例共用同一解释) +CREATE TYPE muse_candidate_state AS ENUM ( + 'draft', + 'candidate', + 'adopted', + 'rejected', + 'retired' +); + +-- 审阅动作:作者决策的闭合集合 +CREATE TYPE muse_review_action AS ENUM ( + 'adopt', + 'reject', + 'revise' +); + +-- 固定角色在目标环境初始化,迁移必须由 muse_maint 执行。 +-- schema 负责用途边界;业务模块后续迁移在相应用途中建立自己的对象。 +REVOKE ALL ON SCHEMA public FROM PUBLIC; +GRANT USAGE ON SCHEMA public TO muse_app, muse_eval; +CREATE SCHEMA evaluation AUTHORIZATION muse_maint; +CREATE SCHEMA oracle AUTHORIZATION muse_maint; +REVOKE ALL ON SCHEMA evaluation, oracle FROM PUBLIC; +GRANT USAGE ON SCHEMA evaluation, oracle TO muse_eval; + +-- 维护账本不授予运行角色写权限;后续生产表只允许生产角色读写。 +ALTER DEFAULT PRIVILEGES FOR ROLE muse_maint IN SCHEMA public + GRANT SELECT, INSERT, UPDATE, DELETE ON TABLES TO muse_app; +ALTER DEFAULT PRIVILEGES FOR ROLE muse_maint IN SCHEMA public + GRANT USAGE, SELECT ON SEQUENCES TO muse_app; +ALTER DEFAULT PRIVILEGES FOR ROLE muse_maint IN SCHEMA evaluation + GRANT SELECT, INSERT, UPDATE, DELETE ON TABLES TO muse_eval; +ALTER DEFAULT PRIVILEGES FOR ROLE muse_maint IN SCHEMA evaluation + GRANT USAGE, SELECT ON SEQUENCES TO muse_eval; +ALTER DEFAULT PRIVILEGES FOR ROLE muse_maint IN SCHEMA oracle + GRANT SELECT ON TABLES TO muse_eval;