muse-agent-example/tests/集成/test_结构版本与并发保护.py
zizi b0623d9048 W04 元数据结构、扩展与字段用途:元数据结构、策略版本、字段用途投影与内置结构导入。
按 R2 串行阶段整理提交;包内文件为该阶段交付(含后续小增量),状态以工作包清单为准。
2026-09-10 19:25:40 +08:00

319 lines
14 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.

"""真实 PostgreSQL 验证元数据版本与调用方提交点;每例独占可销毁数据库。"""
from __future__ import annotations
import uuid
from collections.abc import Iterator
from pathlib import Path
import psycopg
import pytest
from psycopg import sql
from muse.元数据.接口 import (
元数据服务,
元数据错误,
导入内置结构,
读取内置种子,
)
from muse.共享.调用身份 import 用途
from muse.基础设施.数据库.迁移 import 执行迁移
from muse.基础设施.数据库.连接 import 连接
from muse.配置 import 数据库引用
pytestmark = pytest.mark.数据库
迁移目录 = Path(__file__).resolve().parents[2] / "数据库" / "迁移"
@pytest.fixture
def 元数据库(隔离数据库URL: str, monkeypatch: pytest.MonkeyPatch) -> Iterator[数据库引用]:
库名 = f"muse_meta_{uuid.uuid4().hex[:12]}"
with psycopg.connect(隔离数据库URL, autocommit=True) as 管理员:
管理员.execute((迁移目录.parent / "初始化" / "用途角色.sql").read_text(encoding="utf-8"))
管理员.execute(
sql.SQL("CREATE DATABASE {} OWNER muse_maint TEMPLATE template0").format(
sql.Identifier(库名)
)
)
参数 = psycopg.conninfo.conninfo_to_dict(隔离数据库URL)
参数.update(dbname=库名, user="muse_maint")
monkeypatch.setenv("MUSE_META_CASE_URL", psycopg.conninfo.make_conninfo(**参数))
引用 = 数据库引用("环境变量", "MUSE_META_CASE_URL")
try:
with 连接(引用, 用途标记=用途.维护) as 连:
执行迁移(连, 迁移目录, 目标版本=2)
yield 引用
finally:
with psycopg.connect(隔离数据库URL, autocommit=True) as 管理员:
管理员.execute(sql.SQL("DROP DATABASE {} WITH (FORCE)").format(sql.Identifier(库名)))
def test_内置发布可重复且不自动启用__a42001(元数据库: 数据库引用) -> None:
"""given 隔离空库;when 两次导入种子;then 24 个固定版本可回读且没有隐式绑定。"""
with 连接(元数据库, 用途标记=用途.维护) as 连:
服务 = 元数据服务(连)
with 连.transaction():
初次 = 导入内置结构(服务)
再次 = 导入内置结构(服务)
assert len(初次) == len(再次) == 24
assert sum(not f.field_id.startswith("common.") for 项 in 初次 for f in 项.字段) == 216
assert sum(f.field_id.startswith("common.") for 项 in 初次 for f in 项.字段) == 88
for 种子 in 读取内置种子():
assert 服务.读取结构(种子.结构.schema_id, 1) == 种子.结构
with pytest.raises(元数据错误) as 错误:
服务.读取绑定("work-a:character")
assert 错误.value.错误码 == "SCHEMA_UNKNOWN"
def 安装动态结构(连: psycopg.Connection):
from muse.元数据.接口 import (
值规则,
合成结构,
启用命令,
字段定义,
字段用途,
类型定义,
结构定义,
)
服务 = 元数据服务(连)
类型 = 类型定义("weather_spirit", "天气精灵", "world", "entity", ("entity",))
基础 = 结构定义(
"weather_spirit",
"weather_spirit",
1,
(
字段定义(
"weather.name",
"名称",
"名称",
值规则("text"),
必填=True,
用途=字段用途(aiContext=True),
),
),
)
服务.登记类型(类型)
服务.登记结构候选(基础)
服务.发布结构(基础.schema_id, 1, 基础.内容哈希)
有效 = 合成结构(基础)
命令 = 启用命令("activate-initial", 有效.effective_schema_hash, 0, "author-a", "work-a:spirit")
绑定 = 服务.启用绑定("work-a:spirit", "work-a", 有效, 命令)
return 服务, 基础, 有效, 绑定, 命令
def test_扩展发布保留历史绑定且回滚不留半成品__a42002(元数据库: 数据库引用) -> None:
"""given 已启用结构;when 发布作品扩展、回滚与重复内置导入;then 版本固定且旧值可回读。"""
from muse.元数据.接口 import 作品扩展, 值规则, 启用命令, 字段定义
with 连接(元数据库, 用途标记=用途.维护) as 连:
with 连.transaction():
服务, 基础, _, 旧绑定, 原命令 = 安装动态结构(连)
扩展 = 作品扩展(
"work-a",
基础.schema_id,
1,
1,
(
字段定义(
"work-a.affinity", "天气亲和", "天气亲和", 值规则("enum", 枚举=("雨", "雪"))
),
),
)
服务.发布扩展(扩展, 扩展.内容哈希)
服务.发布扩展(扩展, 扩展.内容哈希)
新结构 = 服务.读取有效结构(基础.schema_id, 1, work_id="work-a", extension_version=1)
assert 服务.读取绑定("work-a:spirit") == 旧绑定
with pytest.raises(RuntimeError), 连.transaction():
服务.启用绑定(
"work-a:spirit",
"work-a",
新结构,
启用命令(
"activate-fail", 新结构.effective_schema_hash, 1, "author-a", "work-a:spirit"
),
)
raise RuntimeError("第二参与者失败")
with 连.transaction():
assert 服务.读取绑定("work-a:spirit") == 旧绑定
新绑定 = 服务.启用绑定(
"work-a:spirit",
"work-a",
新结构,
启用命令(
"activate-ok", 新结构.effective_schema_hash, 1, "author-a", "work-a:spirit"
),
)
assert 新绑定.extension_version == 1
导入内置结构(服务)
assert 服务.读取绑定("work-a:spirit") == 新绑定
assert 服务.读取结构(基础.schema_id, 1) == 基础
# 原命令重放返回当时回执,不把当前绑定降回原版本。
assert (
服务.启用绑定(
"work-a:spirit", "work-a", 服务.读取有效结构(基础.schema_id, 1), 原命令
)
== 旧绑定
)
assert 服务.读取绑定("work-a:spirit") == 新绑定
def test_候选未发布与过期绑定均不可启用__a42003(元数据库: 数据库引用) -> None:
"""given 已启用结构;when 候选未发布或预期版本过期;then 不改变当前绑定。"""
from dataclasses import replace
from muse.元数据.接口 import 合成结构, 启用命令
with 连接(元数据库, 用途标记=用途.维护) as 连:
with 连.transaction():
服务, 基础, _, 绑定, _ = 安装动态结构(连)
新版 = replace(基础, schema_version=2, 父版本=1)
服务.登记结构候选(新版)
新 = 合成结构(新版)
with pytest.raises(元数据错误) as 错误:
服务.启用绑定(
绑定.target_ref,
绑定.work_id,
新,
启用命令("draft", 新.effective_schema_hash, 1, "author-a", 绑定.target_ref),
)
assert 错误.value.错误码 == "SCHEMA_UNKNOWN"
服务.发布结构(新版.schema_id, 2, 新版.内容哈希)
with pytest.raises(元数据错误) as 错误:
服务.启用绑定(
绑定.target_ref,
绑定.work_id,
新,
启用命令("stale", 新.effective_schema_hash, 0, "author-a", 绑定.target_ref),
)
assert 错误.value.错误码 == "SCHEMA_STALE"
assert 服务.读取绑定(绑定.target_ref) == 绑定
@pytest.mark.parametrize(
"变更类别", ["structure", "policy"], ids=["schema-switch", "policy-revoke"]
)
def test_结构策略切换与业务提交共享锁__a42004(元数据库: 数据库引用, 变更类别: str) -> None:
"""given 已预览候选;when 提交持有结构锁且另连接切换;then 切换等待且旧请求在切换后拒绝。"""
from concurrent.futures import ThreadPoolExecutor
from dataclasses import asdict, replace
from muse.元数据.接口 import (
合成结构,
启用命令,
字段限制,
定义哈希,
投影字段,
授权摘要,
策略快照,
)
授权 = 授权摘要("author-v1", "source-v1")
with 连接(元数据库, 用途标记=用途.维护) as 连:
with 连.transaction():
服务, 基础, 原结构, 绑定, _ = 安装动态结构(连)
新版 = replace(基础, schema_version=2, 父版本=1)
服务.登记结构候选(新版)
服务.发布结构(新版.schema_id, 2, 新版.内容哈希)
新结构 = 合成结构(新版)
# 仅夹具业务 owner 拥有此表;元数据模块从不访问它。
连.execute("CREATE TABLE public.meta_owner_probe (id int PRIMARY KEY, body text)")
预览 = 投影字段(
原结构,
服务.当前策略(基础.type_id),
授权,
字段用途名="userEditable",
内容用途="planning",
运行用途="production",
)
def 切换(限时: bool) -> None:
with 连接(元数据库, 用途标记=用途.维护) as 第二连, 第二连.transaction():
if 限时:
第二连.execute("SET LOCAL lock_timeout='250ms'")
第二服务 = 元数据服务(第二连)
if 变更类别 == "structure":
第二服务.启用绑定(
绑定.target_ref,
绑定.work_id,
新结构,
启用命令(
"switch", 新结构.effective_schema_hash, 1, "author-a", 绑定.target_ref
),
)
else:
新策略 = 策略快照(
基础.type_id, 1, (字段限制("weather.name", ("userEditable",)),)
)
第二服务.更新策略(
新策略,
启用命令(
"revoke",
定义哈希(asdict(新策略)),
0,
"author-a",
f"type:{基础.type_id}",
),
)
with 连.transaction():
服务.保护提交(
绑定.target_ref,
预览.effective_schema_hash,
预览.projection_version,
授权,
字段用途名="userEditable",
内容用途="planning",
运行用途="production",
)
with ThreadPoolExecutor(max_workers=1) as 池:
with pytest.raises(psycopg.errors.LockNotAvailable):
池.submit(切换, True).result(timeout=5)
连.execute("INSERT INTO public.meta_owner_probe VALUES (1, '作者确认的天气精灵')")
切换(False)
with 连.transaction():
with pytest.raises(元数据错误) as 错误:
服务.保护提交(
绑定.target_ref,
预览.effective_schema_hash,
预览.projection_version,
授权,
字段用途名="userEditable",
内容用途="planning",
运行用途="production",
)
assert 错误.value.错误码 == (
"SCHEMA_STALE" if 变更类别 == "structure" else "PROJECTION_STALE"
)
assert 连.execute("SELECT id FROM public.meta_owner_probe").fetchall() == [(1,)]
def test_已发布定义与扩展追加合同不可绕过__a42005(元数据库: 数据库引用) -> None:
"""given 已发布结构与作品扩展;when 覆盖同版本或删除既有扩展;then 历史定义保持不变。"""
from dataclasses import replace
from muse.元数据.接口 import 作品扩展, 值规则, 字段定义
with 连接(元数据库, 用途标记=用途.维护) as 连, 连.transaction():
服务, 基础, _, _, _ = 安装动态结构(连)
assert 服务.读取类型(基础.type_id).实例族 == ("entity",)
with pytest.raises(元数据错误) as 错误:
服务.登记结构候选(replace(基础, 字段=()))
assert 错误.value.错误码 == "SCHEMA_CONFLICT"
扩展 = 作品扩展(
"work-a",
基础.schema_id,
1,
1,
(字段定义("work-a.temp", "体温", "体温", 值规则("number")),),
)
服务.发布扩展(扩展, 扩展.内容哈希)
丢字段 = replace(扩展, extension_version=2, 字段=())
with pytest.raises(元数据错误) as 错误:
服务.发布扩展(丢字段, 丢字段.内容哈希)
assert 错误.value.错误码 == "FIELD_NOT_ALLOWED"
assert 服务.读取结构(基础.schema_id, 1) == 基础
assert (
服务.读取有效结构(基础.schema_id, 1, work_id="work-a", extension_version=1).扩展 == 扩展
)