396 lines
17 KiB
Python

"""配置版本、验证回执、启用指针与任务冻结副本;不读取未审工作树或凭据值。"""
from __future__ import annotations
import json
import re
import uuid
from collections.abc import Iterator
from contextlib import contextmanager
from dataclasses import asdict, dataclass
from pathlib import Path
from typing import Any, LiteralString, Protocol
from urllib.parse import urlsplit
import psycopg
from psycopg import sql
from psycopg.rows import dict_row
from psycopg.types.json import Jsonb
from muse.任务运行.存储 import 任务存储
from muse.任务运行.模型 import 内容哈希
from muse.共享.调用身份 import 用途
from muse.共享.错误 import Muse错误
from muse.基础设施.数据库.连接 import 数据库工厂
_角色 = frozenset({"writer", "planner", "extractor", "detector", "judge"})
class 配置版本错误(Muse错误):
错误码 = "MUSE_RUNTIME_CONFIG"
@dataclass(frozen=True, slots=True)
class 凭据引用:
名称: str
来源: str
位置: str
def __post_init__(self) -> None:
if not self.名称 or self.来源 not in ("环境变量", "受控存储"):
raise 配置版本错误("配置凭据只接受明确的引用来源")
if self.来源 == "环境变量" and not re.fullmatch(r"[A-Za-z_][A-Za-z0-9_]*", self.位置):
raise 配置版本错误("凭据环境变量引用不合法")
if self.来源 == "受控存储" and not Path(self.位置).is_absolute():
raise 配置版本错误("受控存储引用必须是绝对位置")
@dataclass(frozen=True, slots=True)
class 提供方配置:
身份: str
协议: str
地址: str
凭据名称: str
def __post_init__(self) -> None:
if not self.身份 or self.协议 not in {"responses", "anthropic", "chat-completions"}:
raise 配置版本错误("提供方需要稳定身份和已支持协议")
地址 = urlsplit(self.地址)
if (
地址.scheme not in {"http", "https"}
or not 地址.hostname
or 地址.username
or 地址.password
or 地址.fragment
or 地址.query
):
raise 配置版本错误("提供方地址必须为不含认证信息或查询参数的HTTP端点")
if not self.凭据名称:
raise 配置版本错误("提供方必须绑定具名凭据引用")
@dataclass(frozen=True, slots=True)
class 运行配置内容:
宿主: str
宿主版本: str
角色策略版本: str
资源发布身份: str
预算策略引用: str
角色配置: dict[str, dict[str, Any]]
凭据: tuple[凭据引用, ...]
提供方: tuple[提供方配置, ...]
计价版本: str
Node路径: str | None = None
Pi包目录: str | None = None
def __post_init__(self) -> None:
if self.宿主 not in ("pi", "direct") or not all(
(self.宿主版本, self.角色策略版本, self.资源发布身份, self.预算策略引用, self.计价版本)
):
raise 配置版本错误("配置需要已知宿主以及宿主、角色、资源和预算版本引用")
if not self.角色配置 or not set(self.角色配置) <= _角色:
raise 配置版本错误("配置只允许显式登记的五类角色")
for 策略 in self.角色配置.values():
if set(策略) - {
"provider",
"model",
"thinking",
"tools",
"stage_tools",
"input_schema",
"output_schema",
}:
raise 配置版本错误("角色配置包含不支持的字段;凭据只能放在引用表")
if "stage_tools" in 策略 and (
not isinstance(策略["stage_tools"], dict)
or not all(isinstance(v, list) for v in 策略["stage_tools"].values())
):
raise 配置版本错误("stage_tools 必须是阶段到工具名单的映射")
for 字段 in ("provider", "model", "thinking"):
if not isinstance(策略.get(字段), str) or not 策略[字段].strip():
raise 配置版本错误("角色 provider、model 与 thinking 必须显式配置")
for 字段 in ("input_schema", "output_schema"):
if 字段 in 策略 and (not isinstance(策略[字段], str) or not 策略[字段]):
raise 配置版本错误("结构字段只保存明确的结构版本身份")
if "tools" in 策略 and not (
isinstance(策略["tools"], list) and all(isinstance(t, str) for t in 策略["tools"])
):
raise 配置版本错误("角色工具必须是具名身份列表")
if len({r.名称 for r in self.凭据}) != len(self.凭据):
raise 配置版本错误("凭据引用名称重复")
提供方 = {p.身份: p for p in self.提供方}
if len(提供方) != len(self.提供方) or not 提供方:
raise 配置版本错误("提供方必须非空且身份唯一")
if any(p.凭据名称 not in {r.名称 for r in self.凭据} for p in self.提供方):
raise 配置版本错误("提供方引用了未登记凭据")
if any(p["provider"] not in 提供方 for p in self.角色配置.values()):
raise 配置版本错误("角色引用了未登记提供方")
if self.宿主 == "pi" and not all(
isinstance(p, str) and Path(p).is_absolute() for p in (self.Node路径, self.Pi包目录)
):
raise 配置版本错误("Pi配置必须固定Node与包的绝对位置")
object.__setattr__(self, "角色配置", json.loads(json.dumps(self.角色配置)))
def 冻结(self) -> dict:
return json.loads(json.dumps(asdict(self), ensure_ascii=False))
@classmethod
def 从快照(cls, 内容: dict) -> 运行配置内容:
return cls(
**{
**内容,
"凭据": tuple(凭据引用(**r) for r in 内容["凭据"]),
"提供方": tuple(提供方配置(**p) for p in 内容["提供方"]),
}
)
@dataclass(frozen=True, slots=True)
class 配置验证证据:
配置哈希: str
角色策略版本: str
资源发布身份: str
执行用途: 用途
证据引用: tuple[str, ...]
模式: str
class 配置验证器(Protocol):
"""由装配登记的角色与资源验证实现;接入层不能自报验证成功。"""
@property
def 身份(self) -> str: ...
def 验证(self, 内容: 运行配置内容, 执行用途: 用途) -> 配置验证证据: ...
@dataclass(frozen=True, slots=True)
class 配置快照:
配置ID: str
版本: str
内容哈希: str
内容: 运行配置内容
class 配置版本管理:
def __init__(self, 数据库: 数据库工厂, 验证器: 配置验证器 | None = None) -> None:
self.数据库 = 数据库
self.验证器 = 验证器
self._schema = "evaluation" if 数据库.用途 is 用途.评测 else "public"
def _查询(self, 连: psycopg.Connection, 语句: LiteralString, 参数: tuple = ()) -> Any:
return 连.cursor(row_factory=dict_row).execute(
sql.SQL(语句).format(s=sql.Identifier(self._schema)), 参数
)
@contextmanager
def _事务(self) -> Iterator[psycopg.Connection]:
with self.数据库.连接() as 连, 连.transaction():
yield 连
def 保存草案(self, 配置ID: str, 版本: str, 内容: 运行配置内容) -> 配置快照:
if not 配置ID or not 版本:
raise 配置版本错误("配置身份与版本不能为空")
冻结 = 运行配置内容.从快照(内容.冻结()).冻结()
哈希 = 内容哈希(冻结)
with self._事务() as 连:
self._查询(
连,
"INSERT INTO {s}.muse_runtime_config_version "
"(run_purpose,config_id,version,content,content_hash) VALUES (%s,%s,%s,%s,%s) "
"ON CONFLICT DO NOTHING",
(self.数据库.用途.value, 配置ID, 版本, Jsonb(冻结), 哈希),
)
已存 = self._版本(连, 配置ID, 版本)
if 已存.内容哈希 != 哈希:
raise 配置版本错误("同一配置版本不可覆盖,修改应建立新版本")
return 已存
def 读取版本(self, 配置ID: str, 版本: str) -> 配置快照:
with self._事务() as 连:
return self._版本(连, 配置ID, 版本)
def _版本(self, 连: psycopg.Connection, 配置ID: str, 版本: str) -> 配置快照:
行 = self._查询(
连,
"SELECT * FROM {s}.muse_runtime_config_version "
"WHERE run_purpose=%s AND config_id=%s AND version=%s",
(self.数据库.用途.value, 配置ID, 版本),
).fetchone()
if 行 is None:
raise 配置版本错误("配置版本不存在或不属于当前用途")
if 内容哈希(行["content"]) != 行["content_hash"]:
raise 配置版本错误("配置内容与保存哈希不一致")
return 配置快照(配置ID, 版本, 行["content_hash"], 运行配置内容.从快照(行["content"]))
def 验证版本(self, 配置ID: str, 版本: str) -> str:
if self.验证器 is None:
raise 配置版本错误("配置验证器尚未装配,不能自行声明验证通过")
快照 = self.读取版本(配置ID, 版本)
# 探针与资源核验由对应 owner 执行,验证期间不保持数据库事务。
证据 = self.验证器.验证(快照.内容, self.数据库.用途)
if (
证据.配置哈希 != 快照.内容哈希
or 证据.执行用途 is not self.数据库.用途
or 证据.角色策略版本 != 快照.内容.角色策略版本
or 证据.资源发布身份 != 快照.内容.资源发布身份
):
raise 配置版本错误("验证证据没有绑定当前配置、角色、资源与执行用途")
if (
not self.验证器.身份
or not 证据.证据引用
or 证据.模式 not in ("offline_contract", "runtime")
):
raise 配置版本错误("配置验证缺少明确模式与可回查证据")
if self.数据库.用途 is 用途.生产 and 证据.模式 != "runtime":
raise 配置版本错误("离线协议证据不能启用生产配置")
身份 = str(uuid.uuid4())
with self._事务() as 连:
if self._版本(连, 配置ID, 版本).内容哈希 != 快照.内容哈希:
raise 配置版本错误("验证期间配置发生变化")
self._查询(
连,
"INSERT INTO {s}.muse_runtime_config_validation "
"(receipt_id,run_purpose,config_id,version,content_hash,validator_id,"
"evidence_refs) "
"VALUES (%s,%s,%s,%s,%s,%s,%s)",
(
身份,
self.数据库.用途.value,
配置ID,
版本,
快照.内容哈希,
self.验证器.身份,
Jsonb({"模式": 证据.模式, "引用": list(证据.证据引用)}),
),
)
return 身份
def _锁指针(self, 连: psycopg.Connection, 配置ID: str) -> dict:
self._查询(
连,
"INSERT INTO {s}.muse_runtime_config_active (run_purpose,config_id) "
"VALUES (%s,%s) ON CONFLICT DO NOTHING",
(self.数据库.用途.value, 配置ID),
)
return self._查询(
连,
"SELECT * FROM {s}.muse_runtime_config_active "
"WHERE run_purpose=%s AND config_id=%s FOR UPDATE",
(self.数据库.用途.value, 配置ID),
).fetchone()
def 启用(self, 配置ID: str, 版本: str, *, 验证回执: str, 批准引用: str, 预期代次: int) -> int:
if not 批准引用:
raise 配置版本错误("配置启用需要明确批准引用")
with self._事务() as 连:
当前 = self._锁指针(连, 配置ID)
if (
当前["version"] == 版本
and str(当前["validation_receipt"]) == 验证回执
and 当前["approval_ref"] == 批准引用
):
return 当前["generation"]
if 当前["generation"] != 预期代次:
raise 配置版本错误("配置启用指针已变化")
快照 = self._版本(连, 配置ID, 版本)
证据 = self._查询(
连,
"SELECT * FROM {s}.muse_runtime_config_validation "
"WHERE receipt_id=%s AND run_purpose=%s AND config_id=%s AND version=%s "
"AND content_hash=%s",
(验证回执, self.数据库.用途.value, 配置ID, 版本, 快照.内容哈希),
).fetchone()
if 证据 is None:
raise 配置版本错误("验证回执不属于本配置版本")
新代次 = 当前["generation"] + 1
self._查询(
连,
"UPDATE {s}.muse_runtime_config_active SET version=%s,validation_receipt=%s,"
"approval_ref=%s,generation=%s WHERE run_purpose=%s AND config_id=%s",
(版本, 验证回执, 批准引用, 新代次, self.数据库.用途.value, 配置ID),
)
return 新代次
def 停用(self, 配置ID: str, *, 预期代次: int) -> int:
with self._事务() as 连:
当前 = self._锁指针(连, 配置ID)
if 当前["generation"] != 预期代次:
raise 配置版本错误("配置停用指针已变化")
新代次 = 当前["generation"] + 1
self._查询(
连,
"UPDATE {s}.muse_runtime_config_active SET version=NULL,validation_receipt=NULL,"
"approval_ref=NULL,generation=%s WHERE run_purpose=%s AND config_id=%s",
(新代次, self.数据库.用途.value, 配置ID),
)
return 新代次
def 冻结到任务(self, 任务ID: str, 配置ID: str) -> 配置快照:
with self._事务() as 连:
任务 = 任务存储(连, self.数据库.用途).任务行(任务ID, 锁定=True)
已绑定 = self._查询(
连, "SELECT * FROM {s}.muse_task_config_binding WHERE task_id=%s", (任务ID,)
).fetchone()
if 已绑定:
if 已绑定["config_id"] != 配置ID:
raise 配置版本错误("任务已经绑定另一配置身份")
return 配置快照(
配置ID,
已绑定["version"],
已绑定["content_hash"],
运行配置内容.从快照(已绑定["content"]),
)
当前 = self._锁指针(连, 配置ID)
if 当前["version"] is None:
raise 配置版本错误("配置尚未验证启用或已停用")
快照 = self._版本(连, 配置ID, 当前["version"])
if (快照.内容.角色策略版本, 快照.内容.资源发布身份) != (
任务["frozen_input"]["角色策略版本"],
任务["frozen_input"]["资源发布身份"],
):
raise 配置版本错误("运行配置与任务冻结角色或资源版本不一致")
self._查询(
连,
"INSERT INTO {s}.muse_task_config_binding "
"(task_id,config_id,version,content_hash,content,validation_receipt) "
"VALUES (%s,%s,%s,%s,%s,%s)",
(
任务ID,
配置ID,
快照.版本,
快照.内容哈希,
Jsonb(快照.内容.冻结()),
当前["validation_receipt"],
),
)
return 快照
def 读取任务绑定(self, 任务ID: str) -> 配置快照:
with self._事务() as 连:
行 = self._查询(
连, "SELECT * FROM {s}.muse_task_config_binding WHERE task_id=%s", (任务ID,)
).fetchone()
if 行 is None:
raise 配置版本错误("任务尚未绑定验证启用的运行配置")
if 内容哈希(行["content"]) != 行["content_hash"]:
raise 配置版本错误("任务运行配置快照与保存哈希不一致")
return 配置快照(
行["config_id"],
行["version"],
行["content_hash"],
运行配置内容.从快照(行["content"]),
)
__all__ = [
"提供方配置",
"凭据引用",
"运行配置内容",
"配置验证证据",
"配置验证器",
"配置快照",
"配置版本管理",
"配置版本错误",
]