实现侧: - 上下文:任务范围拆分为 范围校验/范围授权;索引按可发现口径重建、索引新鲜度改对称差;依赖校验统一快照漂移说明。 - 知识方法:方法与材料读取口径统一;超限方法材料按可选省略,核对路径不再二次计费;删除无合同的读时重算。 - 任务运行:新增 context.usage/tool.denied 事件类型;连接池常驻并在装配生命周期内开关;调用结算与核对分列。 - 效果评测/审校修订/交付连载/作者经验/作品规划:凭据冻结、标定消费、导出补证、事实引文核对等收尾修复。 - 资源加载:能力正文不再夹带索引用的导航注记(该注记此前进入角色与技能的模型提示)。 - 元数据:受保护骨架与代码保护属性对齐;字段校验与内置结构口径同步。 - 基础设施:环境预检进入装配生命周期;数据库连接运行期字段不参与相等比较;索引指纹归一化 jsonb 浮点。 - 删除被替代实现:7 份旧提示词模板与空壳 资料来源 读取器。 用例侧: - 用例身份与导航元信息迁移;夹具补生命周期、同库暴露与模板封存; - 本轮定向修复:方法材料省略、事实引文、迁移回执、额度与暂停用例、慢用例超时预算等。
215 lines
8.0 KiB
Python
215 lines
8.0 KiB
Python
"""受控连接与生命周期:凭据引用、固定用途和数据库只读设置。"""
|
||
|
||
from __future__ import annotations
|
||
|
||
from collections.abc import Iterator
|
||
from contextlib import contextmanager
|
||
from dataclasses import dataclass, field
|
||
from importlib import import_module
|
||
from threading import RLock
|
||
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
|
||
|
||
|
||
class 连接池未就绪(Muse错误):
|
||
错误码 = "MUSE_POOL_NOT_OPEN"
|
||
|
||
|
||
class 连接池耗尽(Muse错误):
|
||
错误码 = "POOL_EXHAUSTED"
|
||
可重试 = True
|
||
|
||
|
||
@dataclass(frozen=True, slots=True)
|
||
class 数据库工厂:
|
||
"""同用途应用唯一借还入口;维护工具默认直连,常驻应用显式打开池。"""
|
||
|
||
引用: 数据库引用
|
||
用途: 用途
|
||
|
||
# 池、连接串与锁是运行态资源,不参与相等:同一受控目标的两个工厂必须相等,
|
||
# 否则「同一数据库」的守卫会退化成对象同一性比较。
|
||
常驻: bool = field(default=False, compare=False)
|
||
最大连接: int = field(default=4, compare=False)
|
||
等待秒: float = field(default=2.0, compare=False)
|
||
_池: Any = field(default=None, init=False, repr=False, compare=False)
|
||
_连接串: str | None = field(default=None, init=False, repr=False, compare=False)
|
||
_锁: Any = field(default_factory=RLock, init=False, repr=False, compare=False)
|
||
|
||
def 同一目标(self, 其他: 数据库工厂) -> bool:
|
||
"""受控目标身份只看凭据引用与用途,不看池参数与运行态。"""
|
||
return self.引用 == 其他.引用 and self.用途 is 其他.用途
|
||
|
||
def 打开(self) -> None:
|
||
if not self.常驻:
|
||
return
|
||
with self._锁:
|
||
if self._池 is not None:
|
||
return
|
||
try:
|
||
ConnectionPool = import_module("psycopg_pool").ConnectionPool
|
||
except ImportError:
|
||
raise 连接池未就绪("常驻数据库需要已锁定的 psycopg-pool 依赖") from None
|
||
池 = ConnectionPool(
|
||
conninfo=self._取得连接串,
|
||
min_size=0,
|
||
max_size=self.最大连接,
|
||
timeout=self.等待秒,
|
||
kwargs={
|
||
"autocommit": True,
|
||
"connect_timeout": 10,
|
||
"prepare_threshold": None,
|
||
"application_name": f"muse:{self.用途.value}",
|
||
},
|
||
configure=self._初始化,
|
||
reset=self._复位,
|
||
open=False,
|
||
)
|
||
object.__setattr__(self, "_池", 池)
|
||
try:
|
||
池.open()
|
||
except BaseException:
|
||
池.close()
|
||
object.__setattr__(self, "_池", None)
|
||
raise
|
||
|
||
def 关闭(self) -> None:
|
||
with self._锁:
|
||
池 = self._池
|
||
object.__setattr__(self, "_池", None)
|
||
object.__setattr__(self, "_连接串", None)
|
||
if 池 is not None:
|
||
池.close()
|
||
|
||
def _取得连接串(self) -> str:
|
||
# min_size=0的打开只启动生命周期;真正借用时才需要数据库凭据。
|
||
with self._锁:
|
||
if self._池 is None:
|
||
raise 连接池未就绪("数据库 provider 尚未进入应用生命周期")
|
||
if self._连接串 is None:
|
||
object.__setattr__(self, "_连接串", 解析连接串(self.引用))
|
||
assert self._连接串 is not None
|
||
return self._连接串
|
||
|
||
def _初始化(self, 连) -> None:
|
||
校验用途角色(连, self.用途)
|
||
连.execute("SELECT set_config('app.purpose', %s, false)", (self.用途.value,))
|
||
|
||
def _复位(self, 连) -> None:
|
||
# 业务上下文先提交/回滚;任何会话状态残留都在再次出借前清理。
|
||
if 连.closed:
|
||
return
|
||
连.rollback()
|
||
连.autocommit = True
|
||
连.execute("DISCARD ALL")
|
||
连.read_only = False
|
||
连.isolation_level = None
|
||
连.deferrable = False
|
||
self._初始化(连)
|
||
|
||
@contextmanager
|
||
def 连接(self, *, 只读: bool = False, autocommit: bool = False) -> Iterator[psycopg.Connection]:
|
||
if not self.常驻:
|
||
with 连接(self.引用, 只读=只读, 用途标记=self.用途, autocommit=autocommit) as 连:
|
||
yield 连
|
||
return
|
||
池 = self._池
|
||
if 池 is None:
|
||
raise 连接池未就绪("数据库 provider 尚未进入应用生命周期")
|
||
self._取得连接串()
|
||
# 超时只翻译借用阶段,不能把业务异常误报为容量耗尽。
|
||
try:
|
||
连 = 池.getconn(timeout=self.等待秒)
|
||
except Exception as exc:
|
||
if type(exc).__name__ == "PoolTimeout":
|
||
raise 连接池耗尽("数据库并发容量已占满;请稍后重试") from None
|
||
raise
|
||
try:
|
||
连.autocommit = True
|
||
连.execute(
|
||
"SELECT set_config('default_transaction_read_only', %s, false)",
|
||
("on" if 只读 else "off",),
|
||
)
|
||
连.read_only = 只读
|
||
连.autocommit = autocommit
|
||
try:
|
||
yield 连
|
||
except BaseException:
|
||
连.rollback()
|
||
raise
|
||
else:
|
||
连.commit()
|
||
finally:
|
||
池.putconn(连)
|
||
|
||
def 检查连接(self) -> None:
|
||
with self.连接(只读=True) as 连:
|
||
连.execute("SELECT 1").fetchone()
|
||
|
||
|
||
__all__ = ["解析连接串", "连接", "数据库工厂", "数据库连接失败"]
|