ErisPulse.Core.Bases.sql_base 模块
模块概述
ErisPulse SQL 存储共享基类
提供 SQLite / MySQL / PostgreSQL 三种内置后端共享的通用 SQL 逻辑:
嵌套键 KV 存取、事务编排、DDL、查询构建与 ALTER TABLE。
方言差异(占位符、标识符引号、UPSERT、类型翻译等)收敛到
:class:SQLDialect 及其子类。
提示
- 新增 SQL 后端只需:定义 Dialect 子类 + 继承 SQLStorageBase 实现连接管理与执行漏斗
- KV 值统一以 JSON 文本存储,三种后端行为完全一致
函数列表
_validate_identifier(name: str, context: str = 'identifier')
内部方法 验证 SQL 标识符(表名/列名)是否安全
- name (
标识符名称): - context: 上下文描述(用于错误消息) 异常:ValueError- 当标识符包含非法字符时
_validate_select_column(name: str, context: str = 'column')
内部方法 验证 SELECT/ORDER BY 列表达式是否安全
采用黑名单模式:仅拦截 SQL 注入危险字符(; ' " -- /* */ 换行), 允许任意合法 SQL 列表达式,包括简单列名、聚合函数、别名、表达式等。
- name (
列表达式): - context: 上下文描述(用于错误消息) 异常:ValueError- 当包含注入危险字符时
_validate_column_type(col_type: str)
内部方法 验证列类型定义是否安全(防止通过类型定义注入 SQL)
- col_type (
列类型定义): 异常:ValueError- 当列类型包含潜在危险内容时
类列表
class _SingletonMixin
内部方法 进程级单例混入(
__new__级去重,可__new__绕过用于测试)
class SQLDialect
SQL 方言基类
收敛三种内置后端的方言差异:占位符翻译、标识符引号、UPSERT 语法、 列类型翻译(含 AUTOINCREMENT 自增主键)、表存在性查询等。
提示
- 构建期统一使用
?占位符记号,由 :meth:translate_sql翻译为方言样式- 标识符一律经 :meth:
quote引用(列名/表名均先经合法性校验)
方法列表
translate_sql(sql: str)
将内部 ? 占位符记号翻译为方言占位符
- sql (
内部记号): SQL 返回值 (方言): SQL
quote(name: str)
引用 SQL 标识符
- name (
已通过合法性校验的标识符): 返回值: 带方言引号的标识符
kv_table_ddl(table: str, key_col: str, value_col: str)
生成 KV 表建表 DDL
- table (
已引用的表名): - key_col: 已引用的键列名 - value_col (
已引用的值列名): 返回值 (CREATE): TABLE 语句
upsert_kv_sql(table: str, key_col: str, value_col: str)
生成 KV UPSERT 语句(内部 ? 占位符)
- table (
已引用的表名): - key_col: 已引用的键列名 - value_col (
已引用的值列名): 返回值 (UPSERT): 语句
has_table_sql(table_name: str)
生成表存在性查询
- table_name (
表名(未引用)): 返回值 ((SQL,): 参数列表)
create_index_sql(table: str, index: str, column: str)
生成普通索引创建 DDL
基类语义为"不存在则创建"(IF NOT EXISTS);不支持该语法的方言
(MySQL)覆写本方法为裸 CREATE INDEX,幂等由 :meth:has_index_sql
存在性预检承接。
- table (
表名(未引用)): - index: 索引名(未引用) - column (
列名(未引用)): 返回值 (CREATE): INDEX 语句
has_index_sql(table: str, index: str)
生成索引存在性查询(内部 ? 占位符)
- table (
表名(未引用)): - index: 索引名(未引用) 返回值 ((SQL,): 参数列表)
autoincrement_column(base_type: str)
翻译自增主键定义("INTEGER PRIMARY KEY AUTOINCREMENT" 中的类型部分)
- base_type (
自增列基础类型(INTEGER/BIGINT): 等) 返回值: 方言等价定义
last_insert_id_sql()
生成"最近一次插入的自增主键"查询 SQL
供 ORM(Core/Bases/model.py)在 INSERT 后同一连接上回填自增主键。 须与 INSERT 在同一连接/事务内执行(MySQL LAST_INSERT_ID() 与 PostgreSQL lastval() 均为连接级语义)。
返回值 (查询): SQL(内部 ? 占位符记号,无参数)
_map_first_word(word: str)
内部方法 按首词映射列类型,保留括号参数与后缀
translate_column_type(col_type: str)
翻译列类型定义为方言等价写法
处理自增主键重排与首词类型映射;无法识别的类型原样保留。
- col_type (
用户传入的列类型定义): 返回值: 方言列类型定义
is_missing_table_error(exc: BaseException)
判断异常是否为"KV 表不存在"(用于自动重建后重试)
- exc (
捕获到的异常): 返回值: 是否为缺表错误
class SQLQueryBuilder(BaseQueryBuilder)
SQL 查询构建器(方言无关的共享实现)
链式构建标准 SQL(内部 ? 占位符),终止方法通过所属后端的
:meth:SQLStorageBase._execute_query 漏斗执行,由方言完成占位符翻译。
提示 使用方式:
- await storage.Table("users").Insert({"name": "Alice"}).aExecute()
- storage.Table("users").Select("name").Where("age > ?", 18).Execute()
方法列表
_rows_to_dict(rows: list[tuple], columns: list[str] | None)
内部方法 将 tuple 行列表转为字典列表(columns 不可用时原样返回)
async aExecute()
执行构建的查询
SELECT 返回 list[tuple](调用 ToDict() 后为 list[dict])
INSERT/UPDATE/DELETE 返回受影响行数 int
conn (
内部参数:事务连接路由(勿手动传入)): 返回值 (查询结果或受影响行数): 示例:
>>> rows = await storage.Table("users").Select("name", "age").aExecute()
>>> affected = await storage.Table("users").Delete().Where("age < ?", 18).aExecute()
async _aexecute_insert_multi()
内部方法 批量插入执行
async aExecuteOne()
执行查询并返回单条结果
- conn (
内部参数:事务连接路由(勿手动传入)): 返回值 (单行元组(调用): ToDict() 后为字典)或 None
示例:
>>> row = await storage.Table("users").Select("*").Where("id = ?", 1).aExecuteOne()
async aCount()
执行 COUNT 查询
- conn (
内部参数:事务连接路由(勿手动传入)): 返回值 (匹配的行数): 示例:
>>> total = await storage.Table("users").Where("age > ?", 18).aCount()
async aExists()
检查是否存在匹配的记录
- conn (
内部参数:事务连接路由(勿手动传入)): 返回值 (是否存在): 示例:
>>> if await storage.Table("users").Where("name = ?", "Alice").aExists():
class AlterTableBuilder
ALTER TABLE 构建器
链式收集表结构修改操作,aExecute() 原生异步执行,
Execute() 同步兼容桥接。
提示 使用方式:
- await storage.AlterTable("users").AddColumn("email", "TEXT").aExecute()
- storage.AlterTable("users").RenameTo("members").Execute()
方法列表
AddColumn(column_name: str, column_type: str)
添加列
- column_name (
列名): - column_type: 列类型(如 "TEXT", "INTEGER DEFAULT 0") 返回值 (self): 示例:
>>> storage.AlterTable("users").AddColumn("email", "TEXT").Execute()
RenameTo(new_name: str)
重命名表
- new_name (
新表名): 返回值 (self): 示例:
>>> storage.AlterTable("users").RenameTo("members").Execute()
async aExecute()
执行所有已收集的 ALTER TABLE 操作
返回值: 操作是否成功
Execute()
执行所有已收集的 ALTER TABLE 操作(同步兼容,桥接到 :meth:aExecute)
返回值: 操作是否成功
class SQLStorageBase(BaseStorage)
SQL 存储后端共享基类
实现三种内置 SQL 后端共享的 KV 存取(含嵌套键)、批量操作、DDL、 事务编排与 ALTER TABLE;子类只需提供连接管理与方言执行漏斗。
提示 子类必须实现:
_create_loop_resource/_destroy_loop_resource:每事件循环资源(连接池)_acquire_resource_conn/_release_resource_conn:非事务连接获取/归还_open_txn_conn/_close_txn_conn:事务专用连接获取/释放_exec_query_on:方言执行漏斗(占位符翻译 + 游标语义)
方法列表
_ensure_resource_state()
内部方法 确保实例级资源表存在(new 绕过构造时兜底)
_is_ready()
内部方法 检查存储后端是否已初始化完成
_finish_init()
内部方法 子类完成后调用:标记就绪并订阅配置热更新告警
_watch_storage_config(check: Callable[[], str | None])
内部方法 订阅存储配置热更新;check() 返回告警文案(None 表示无需告警)
_default_project_path(filename: str)
内部方法 项目目录下的默认数据库文件路径
async _create_loop_resource()
内部方法 为当前事件循环创建连接资源(连接池/共享连接)
async _destroy_loop_resource(resource: Any)
内部方法 销毁事件循环连接资源
async _acquire_resource_conn(resource: Any)
内部方法 从资源获取非事务连接
async _release_resource_conn(resource: Any, conn: Any)
内部方法 归还非事务连接
async _open_txn_conn()
内部方法 获取事务专用连接
async _close_txn_conn(conn: Any)
内部方法 释放事务专用连接
async _get_loop_resource()
内部方法 获取当前事件循环的连接资源(惰性创建,幂等,瞬时失败自动重试)
async _emit_storage_event(event: str)
内部方法 发出存储生命周期事件(
storage.ready/storage.unreachable/storage.recovered),供外部感知连接状态变化(如 Dashboard 告警)。
后台发射(fire):存储操作绝不等待观测者;lifecycle 未就绪或 处理异常时静默跳过,不影响存储操作本身。
async _acquire()
内部方法 获取连接的统一入口:优先异步事务上下文绑定连接,否则走事件循环资源
async _run_with_conn(impl: Callable[..., Any])
内部方法 执行
impl(conn, *args):显式连接优先,否则自动获取并在退出时归还
async _exec_query_on(kind: str, sql: str, params: Any, conn: Any)
内部方法 方言执行漏斗:翻译占位符后按 kind 执行
- kind (
"select"): / "one" / "count" / "dml" / "dml_multi" - sql (
内部):?占位符 SQL - params (
参数列表(dml_multi): 为参数行列表) - conn (
数据库连接): 返回值 (select/one): 返回 (行, 列名列表或None);其余返回受影响行数
async _execute_query(kind: str, sql: str, params: Any)
内部方法 查询执行统一入口(连接路由 → 执行漏斗)
async _acquire_txn_conn()
内部方法 获取事务专用连接
async _begin_txn(conn: Any)
内部方法 开启事务(SQL BEGIN,方言可覆写)
async _commit_txn(conn: Any, handle: Any = None)
内部方法 提交事务
async _rollback_txn(conn: Any, handle: Any = None)
内部方法 回滚事务
async _release_txn_conn(conn: Any)
内部方法 释放事务专用连接
async _run_alter(table_name: str, operations: 'list[tuple[str, tuple[Any, ...]]]')
内部方法 顺序执行 ALTER TABLE 操作
_parse_nested_key(key: str)
内部方法 解析嵌套键:点号(.)总是表示嵌套访问,即使根键不存在也会创建嵌套结构
- key (
键名,如): "user.settings.theme" 返回值 ((根键名,): 路径列表)
_get_nested_value(obj: Any, key_path: list[str])
内部方法 从嵌套对象中获取值
_set_nested_value(obj: Any, key_path: list[str], value: Any)
内部方法 在嵌套对象中设置值(中间层一律预创建为字典)
_delete_nested_value(obj: Any, key_path: list[str])
内部方法 从嵌套对象中删除值,返回 (更新后的对象, 是否删除成功)
_loads(raw: Any)
内部方法 JSON 反序列化(失败时返回原始文本)
_kv_table_sql()
内部方法 返回 (已引用表名, 已引用键列, 已引用值列)
async _init_kv_table()
内部方法 创建默认 KV 表
async _recover_missing_table()
内部方法 缺表自动恢复
_shadow_overlay()
内部方法 影子存储覆盖层(方向十一):当前 owner 属影子时返回其覆盖层, 否则 None。KV 三件套(aget / aset / adelete)据此实现"写隔离、 读透传、删为墓碑";ORM 读写不在覆盖层语义内(文档声明)。
async aget(key: str, default: Any = None)
异步获取存储项的值
支持嵌套键访问,如 "user.settings.theme" 会从存储的嵌套对象中获取值
- key (
存储项键名,支持嵌套路径): - default: 默认值(当键不存在时返回) - conn (
内部参数:事务连接路由(勿手动传入)): 返回值 (存储项的值): 示例:
>>> timeout = await storage.aget("network.timeout", 30)
async aset(key: str, value: Any)
异步设置存储项的值
支持嵌套键设置,如 "user.settings.theme" 会更新存储的嵌套对象中的对应字段
- **key** (`存储项键名,支持嵌套路径`): - **value**: 存储项的值
- **conn** (`内部参数:事务连接路由(勿手动传入)`): **返回值** (`操作是否成功`):
示例:
>>> await storage.aset("user.settings.theme", "dark")
async adelete(key: str)
异步删除存储项
支持嵌套键删除,如 "user.settings.theme" 会删除嵌套对象中的对应字段
- **key** (`存储项键名,支持嵌套路径`): - **conn**: 内部参数:事务连接路由(勿手动传入)
**返回值** (`操作是否成功`):
示例:
>>> await storage.adelete("user.settings.theme")
async aget_all_keys()
异步获取所有存储项的键名
- conn (
内部参数:事务连接路由(勿手动传入)): 返回值 (键名列表): 示例:
>>> all_keys = await storage.aget_all_keys()
async aclear()
异步清空所有存储项
- conn (
内部参数:事务连接路由(勿手动传入)): 返回值 (操作是否成功): 示例:
>>> await storage.aclear()
async aget_multi(keys: list[str])
异步批量获取多个存储项的值(单次 IN 查询,仅返回存在的键)
- keys (
键名列表): 返回值 (键值对字典): 示例:
>>> settings = await storage.aget_multi(["app.name", "app.version"])
get_multi(keys: list[str])
批量获取多个存储项的值(同步兼容,桥接到 :meth:aget_multi)
- keys (
键名列表): 返回值 (键值对字典): 示例:
>>> settings = storage.get_multi(["app.name", "app.version"])
async aset_multi(items: dict[str, Any])
异步批量设置多个存储项
与旧版语义一致:键按字面量直写,不做点号嵌套解析
(需要嵌套行为请逐键调用 :meth:aset)。
- items (
键值对字典): 返回值 (操作是否成功): 示例:
>>> await storage.aset_multi({"app.name": "MyApp", "app.debug": True})
set_multi(items: dict[str, Any])
批量设置多个存储项(同步兼容,桥接到 :meth:aset_multi)
- items (
键值对字典): 返回值 (操作是否成功): 示例:
>>> storage.set_multi({"app.name": "MyApp", "app.debug": True})
async adelete_multi(keys: list[str])
异步批量删除多个存储项(单次批量 DELETE)
- keys (
键名列表): 返回值 (操作是否成功): 示例:
>>> await storage.adelete_multi(["temp.key1", "temp.key2"])
delete_multi(keys: list[str])
批量删除多个存储项(同步兼容,桥接到 :meth:adelete_multi)
- keys (
键名列表): 返回值 (操作是否成功): 示例:
>>> storage.delete_multi(["temp.key1", "temp.key2"])
Table(table_name: str)
获取指定表的查询构建器
- table_name (
表名): 返回值 (SQLQueryBuilder): 实例
示例:
>>> rows = await storage.Table("users").Select("name", "age").Where("age > ?", 18).aExecute()
async aCreateTable(table_name: str, columns: dict[str, str])
异步创建表
列类型支持 SQLite 风格定义(如 "INTEGER PRIMARY KEY AUTOINCREMENT"), 由方言自动翻译为目标后端等价写法。
- table_name (
表名): - columns: 列名到类型的映射 返回值 (操作是否成功): 示例:
>>> await storage.aCreateTable("users", {
... "id": "INTEGER PRIMARY KEY AUTOINCREMENT",
... "name": "TEXT NOT NULL"
... })
async aDropTable(table_name: str)
异步删除表
- table_name (
表名): 返回值 (操作是否成功): 示例:
>>> await storage.aDropTable("users")
async aHasTable(table_name: str)
异步检查表是否存在
- table_name (
表名): 返回值 (是否存在): 示例:
>>> if await storage.aHasTable("users"):
async aGetTableColumns(table_name: str)
异步列举表的现有列名(ORM 自动迁移用)
- table_name (
表名): 返回值: 列名列表;查询失败时返回空列表
AlterTable(table_name: str)
获取 ALTER TABLE 构建器
- table_name (
表名): 返回值 (AlterTableBuilder): 实例
示例:
>>> storage.AlterTable("users").AddColumn("email", "TEXT").Execute()
async aclose()
异步关闭当前事件循环上绑定的连接资源(连接池/共享连接)
事务专用连接不受影响(由事务自行管理)。
示例:
>>> await storage.aclose()