Skip to main content

Crate sz_orm_core

Crate sz_orm_core 

Source
Expand description

§SZ-ORM — 鲜视达 ORM

Rust 异步 ORM 工作空间(原型阶段),兼容 ThinkORM 风格。

§架构概览

SZ-ORM 工作空间由 39 个成员 组成(37 个 sz-orm-* lib + cli + examples):

§核心引擎 (sz-orm-core)

模块功能
modelModel trait — 定义表名、主键、时间戳、软删除、关联关系
queryQueryBuilder<M> — 链式 API,支持 SELECT/INSERT/UPDATE/DELETE/聚合/分页/JOIN
dialect多数据库方言 — MySQL (反引号)、PostgreSQL (双引号)、SQLite、Oracle 23ai
pool异步连接池 — 可配置大小、超时、空闲回收、健康检查、最大生命周期
transactionACID 事务 — 隔离级别、保存点、TransactionManager 多事务管理
migration文件迁移系统 — up/down/rollback/reset/refresh,含 SchemaBuilder
cache多级缓存 — MemoryCacheMultiLevelCache,支持 TTL
value统一值类型 — 20 种变体 (整数/浮点/字符串/字节/UUID/日期/JSON/数组)
db_type数据库类型枚举 — MySQL、PostgreSQL、SQLite、Oracle、Redis、MongoDB 等 11 种
error错误类型体系 — DbError(20 变体)、PoolErrorCacheErrorTxError

§数据库适配器

  • sz-orm-sqlx — sqlx 适配器,连接真实 MySQL/PostgreSQL/SQLite/Oracle
  • sz-orm-sql-validator — SQL 校验与注入检测

§扩展生态包 (18 个)

包名功能
sz-orm-crypto加密原语 (AES-256-GCM, PBKDF2, HMAC-SHA256)
sz-orm-authJWT 鉴权 (HS256)
sz-orm-schedulerCron 定时任务调度
sz-orm-mqttMQTT 客户端 (rumqttc)
sz-orm-websocketWebSocket 服务端 (tokio-tungstenite)
sz-orm-queue消息队列 (RabbitMQ/lapin, Kafka, NATS, ActiveMQ, RocketMQ, Pulsar)
sz-orm-storage对象存储 (S3/阿里云/腾讯云/华为云/七牛/又拍云/本地)
sz-orm-aiAI 集成 (Embedding, RAG, Vector)
sz-orm-grpcgRPC 服务/客户端
sz-orm-graphqlGraphQL 查询支持
sz-orm-esElasticsearch 集成
sz-orm-tracing分布式追踪
sz-orm-logger日志系统
sz-orm-swaggerAPI 文档生成
sz-orm-masking数据脱敏
sz-orm-health健康检查
sz-orm-audit审计日志
sz-orm-batch批量操作

§高级特性包 (6 个)

包名功能
sz-orm-dtx分布式事务
sz-orm-rw读写分离
sz-orm-sharding分库分表
sz-orm-limit限流控制
sz-orm-config配置管理
sz-orm-mig迁移管理增强

§平台支持

  • sz-orm-wasm — WebAssembly 编译目标
  • sz-orm-lc — 本地/边缘计算
  • sz-orm-back — 备份与恢复

§快速入门

use sz_orm_core::*;

// 1. 定义模型
#[derive(Clone)]
struct User {
    id: i64,
    name: String,
    email: String,
}

impl Model for User {
    type PrimaryKey = i64;
    fn table_name() -> &'static str { "users" }
    fn pk(&self) -> Self::PrimaryKey { self.id }
    fn set_pk(&mut self, pk: Self::PrimaryKey) { self.id = pk; }
}

// 2. 构建查询
let dialect = get_dialect(DbType::MySQL).unwrap();
let sql = QueryBuilder::<User>::new(dialect)
    .table("users")
    .select(vec!["id", "name", "email"])
    .where_cond("status = 'active'")
    .order_by("created_at")
    .order_desc("id")
    .limit(10)
    .build_select();

// 3. 执行前校验
QueryBuilder::<User>::new(get_dialect(DbType::MySQL).unwrap())
    .table("users")
    .select(vec!["id", "name"])
    .validate()?; // 校验 SQL 语法、注入、括号平衡

// 4. 其他操作
let mut data = std::collections::HashMap::new();
data.insert("name".to_string(), Value::String("Alice".to_string()));
data.insert("age".to_string(), Value::I64(25));

let insert_sql = QueryBuilder::<User>::new(dialect)
    .table("users")
    .build_insert(&data);

let update_sql = QueryBuilder::<User>::new(get_dialect(DbType::MySQL).unwrap())
    .table("users")
    .where_cond("id = 1")
    .build_update(&data);

let delete_sql = QueryBuilder::<User>::new(get_dialect(DbType::MySQL).unwrap())
    .table("users")
    .where_cond("id = 1")
    .build_delete();

§支持的数据库

数据库方言实现真实连接引用方式
MySQLMySqlDialect (` 反引号)sz-orm-sqlx
PostgreSQLPostgreSqlDialect (" 双引号)sz-orm-sqlx
SQLite 3.35+SqliteDialect (" 双引号)sz-orm-sqlx
Oracle 23aiOracleDialect (类型自动映射)sz-orm-sqlx

通过 get_dialect(DbType::MySQL) 获取方言实例。每个方言处理:

  • 标识符引用风格
  • 字符串转义规则
  • 分页语法 (LIMIT/OFFSET vs OFFSET/FETCH)
  • JSON 提取函数 (JSON_EXTRACT vs #>> vs json_extract vs JSON_VALUE)
  • 全文搜索 (MATCH AGAINST vs to_tsvector vs CONTAINS)
  • 布尔转整数 (IF/CASE)
  • 自增关键字 (AUTO_INCREMENT/GENERATED BY DEFAULT AS IDENTITY)

§核心功能详解

§QueryBuilder API

所有查询方法返回 Self,支持链式调用:

// 基础查询
QueryBuilder::<M>::new(dialect)
    .table("users")
    .select(vec!["id", "name"])
    .where_cond("status = 'active'")    // AND
    .or_where("role = 'admin'")         // OR
    .where_in("id", vec![Value::I64(1), Value::I64(2)])
    .where_between("age", Value::I64(18), Value::I64(30))
    .where_null("deleted_at")
    .order_by("created_at")
    .order_desc("id")
    .group_by("status")
    .having("COUNT(*) > 5")
    .limit(20)
    .offset(40)
    .page(3, 20)                       // page=3, page_size=20
    .join_inner("posts", "users.id", "posts.user_id")
    .join_left("profiles", "users.id", "profiles.user_id")
    .build_select();

// 聚合函数
builder.build_count();    // SELECT COUNT(*)
builder.build_exists();   // SELECT EXISTS(...)
builder.build_max("score");
builder.build_min("price");
builder.build_sum("amount");
builder.build_avg("value");

§SQL 校验

// 编译时 + 运行时双重校验
builder.validate()?;              // 校验 SELECT
builder.validate_insert(&data)?;  // 校验 INSERT(含空数据检测)
builder.validate_update(&data)?;  // 校验 UPDATE(含空数据检测)
builder.validate_delete()?;       // 校验 DELETE

// 校验内容包括:SQL 语法、注入检测、括号平衡、
// 表名/列名合法性、JOIN 列名校验

§Model Trait

pub trait Model: Send + Sync + Sized + 'static {
    type PrimaryKey: Send + Sync + Debug + Display + Clone + Default;

    fn table_name() -> &'static str;          // 表名(必需)
    fn pk_name() -> &'static str { "id" }     // 主键列名
    fn pk(&self) -> Self::PrimaryKey;         // 获取主键值
    fn set_pk(&mut self, pk: Self::PrimaryKey); // 设置主键值
    fn foreign_key(relation: &str) -> String; // 外键命名 "user_id"
    fn timestamp_fields() -> Option<TimestampFields>; // 自动时间戳
    fn soft_delete_field() -> Option<&'static str>;   // 软删除字段
}

// ModelExt 扩展
pub trait ModelExt: Model {
    fn columns() -> Vec<&'static str>;     // 所有列
    fn fillable() -> Vec<&'static str>;    // 可填充列
    fn guarded() -> Vec<&'static str>;     // 保护列(默认含主键)
    fn hidden() -> Vec<&'static str>;      // 隐藏列(不序列化)
    fn relations() -> HashMap<&str, Relation>; // 关联关系
    fn fill(&mut self, data: HashMap<String, Value>); // 批量赋值
    fn to_json(&self) -> serde_json::Value; // 序列化
}

// 四种关联关系
// BelongsTo   — 多对一(Order → User)
// HasMany     — 一对多(User → Orders)
// HasOne      — 一对一(User → Profile)
// BelongsToMany — 多对多(User ↔ Role,通过中间表)

§连接池

// 通过 Builder 配置
let config = PoolConfigBuilder::new()
    .max_size(100)       // 最大连接数
    .min_idle(10)        // 最小空闲连接
    .acquire_timeout(30) // 获取超时(秒)
    .idle_timeout(600)   // 空闲超时(秒)
    .max_lifetime(1800)  // 最大生命周期(秒)
    .build()?;

let pool = Pool::new(config, factory)?;
let conn = pool.acquire().await?;  // 获取连接(带超时)
pool.release(conn).await;         // 归还连接
pool.status().await;               // PoolStatus { idle, active, max, min }
pool.reap_idle().await;           // 回收空闲连接
pool.close_all().await;           // 关闭所有连接

§事务

// 事务选项
let opts = TransactOptions::default()
    .with_isolation(IsolationLevel::Serializable)
    .read_only()
    .with_timeout(Duration::from_secs(30));

let mut tx = Transaction::new(conn, opts);
tx.execute("INSERT INTO users VALUES (1)").await?;
tx.query("SELECT * FROM users").await?;

// 保存点(嵌套事务)
let sp = tx.savepoint().await?;         // SAVEPOINT sp_N
tx.rollback_to_savepoint(&sp).await?;   // ROLLBACK TO SAVEPOINT sp_N
tx.release_savepoint(&sp).await?;       // RELEASE SAVEPOINT sp_N

tx.commit().await?;
// tx.rollback().await?;

// TransactionManager:管理多个命名事务
let mgr = TransactionManager::new();
mgr.begin("tx1", conn, opts).await?;
mgr.commit("tx1").await?;
mgr.list().await;        // ["tx1"]
mgr.state("tx1").await;  // Some(TransactionState::Committed)

§迁移系统

// 文件命名:<version>_<name>_up.sql / <version>_<name>_down.sql
// 示例:001_create_users_up.sql, 001_create_users_down.sql

let resolver = FileMigrationResolver::new(PathBuf::from("./migrations"));
let migrations = resolver.resolve(DbType::MySQL)?;

let mut migrator = Migrator::new(MigrationContext::default())
    .add_migrations(migrations);

migrator.migrate().await?;                     // 执行所有待迁移
migrator.up(Some("003")).await?;               // 执行到指定版本
migrator.down(Some("001")).await?;             // 回滚到指定版本
migrator.rollback("002").await?;               // 回滚单个迁移
migrator.reset().await?;                       // 全部回滚 + 重新执行
migrator.refresh().await?;                     // 同 reset
migrator.progress();                            // MigrationProgress { total, applied, pending }

// SchemaBuilder:程序化建表
let sql = SchemaBuilder::new("users")
    .add_column(ColumnDef::new("id", "INT").not_null().auto_increment())
    .add_column(ColumnDef::new("name", "VARCHAR").length(255).not_null())
    .add_index(IndexDef::new("idx_name", vec!["name"]).unique())
    .add_foreign_key(
        ForeignKeyDef::new("fk_role", "role_id", "roles", "id")
            .on_delete("CASCADE")
    )
    .build(DbType::MySQL);

§值类型 (Value)

// 20 种变体,覆盖所有数据库类型
Value::Null | Bool(bool) | I8..I64 | U8..U64 | F32 | F64
| String(String) | Bytes(Vec<u8>) | Uuid(String) | Date(String)
| DateTime(String) | Time(String) | Json(String) | Array(Vec<Value>)

// 类型转换
value.as_str()    // Option<&str>
value.as_i64()    // Option<i64>(支持 F32/F64/Bool/String→i64 转换)
value.as_f64()    // Option<f64>
value.as_bool()   // Option<bool>(支持 "true"/"1"/"yes"/"on" 等)
value.as_bytes()  // Option<&[u8]>
value.to_param()  // Cow<str> — SQL 参数格式

// From 实现
let v: Value = 42i64.into();
let v: Value = "hello".into();
let v: Value = vec![1u8, 2u8].into();

§错误处理

统一错误类型体系,每种错误携带唯一错误码:

// DbError — 20 种变体,错误码 DB001-DB020
DbError::QueryError("...")
DbError::ConnectionRefused("...")
DbError::ConnectionTimeout("...")
DbError::NotFound("...")
DbError::ConstraintViolation("...")
// ... 等

// PoolError — 6 种变体,错误码 PL001-PL006
PoolError::Exhausted | Timeout | AlreadyAcquired | InvalidConfig | ...

// CacheError — 6 种变体,错误码 CH001-CH006
// TxError — 6 种变体(NotStarted, CommitFailed, SavepointError 等)

// 便捷方法
DbError::query("test failed")        // 创建查询错误
DbError::connection("timeout")       // 创建连接错误
DbError::not_found("user #42")       // 创建未找到错误
err.is_retryable()                   // 是否可重试
err.error_code()                     // "DB001"

§验证方法

SZ-ORM 通过 7 线验证体系 保证质量:

验证方法描述测试文件
TDD核心模块 115+ 单元测试core.rs
集成真实 MySQL/PG/SQLite/Oracle 端到端integration_mysql.rs, integration_pg.rs, integration_sqlite.rs
Jepsen29 并发正确性测试 + 10 真实 DB Jepsenjepsen.rs, real_db_jepsen.rs
Fuzz11 边界/边缘案例发现fuzz.rs
Stress77 性能基准测试stress.rs, core_bench.rs
Chaos16 故障鲁棒性测试chaos.rs
Formal14 形式化验证不变量formal.rs

总计:2,271 测试,0 失败(1,635 #[test] + 636 #[tokio::test];部分需真实服务)

§类型别名与常量

// 类型别名
pub type Shared<T> = Arc<T>;
pub type Boxed<T> = Box<T>;
pub type DbResult<T> = Result<T, DbError>;
pub type PoolResult<T> = Result<T, PoolError>;
pub type CacheResult<T> = Result<T, CacheError>;
pub type TxResult<T> = Result<T, TxError>;

// 默认常量
pub const DEFAULT_BATCH_SIZE: usize = 1000;
pub const DEFAULT_ACQUIRE_TIMEOUT: u64 = 30;   // 秒
pub const DEFAULT_IDLE_TIMEOUT: u64 = 600;      // 秒
pub const DEFAULT_MAX_LIFETIME: u64 = 1800;     // 秒
pub const DEFAULT_MIN_IDLE: u32 = 5;
pub const DEFAULT_MAX_SIZE: u32 = 100;

§导出清单

use sz_orm_core::*; 将导入以下模块的全部公共符号:

  • async_trait (重导出)、bytes::Byteschrono::{DateTime, Utc}serde::{Deserialize, Serialize}
  • cache::*Cache, MemoryCache, MultiLevelCache, CacheStats
  • db_type::*DbType 枚举 (11 种数据库)
  • dialect::*Dialect, MySqlDialect, PostgreSqlDialect, SqliteDialect, OracleDialect, get_dialect()
  • error::*DbError, PoolError, CacheError, TxError
  • migration::*Migration, Migrator, SchemaBuilder, ColumnDef, IndexDef, ForeignKeyDef
  • model::*Model, ModelExt, Relation, BelongsTo, HasMany, HasOne, BelongsToMany
  • pool::*Pool, PoolConfig, PoolConfigBuilder, Connection, ConnectionFactory, PoolStatus
  • query::*QueryBuilder<M> (链式 SQL 构造器)
  • transaction::*Transaction, TransactionManager, TransactOptions, IsolationLevel
  • value::*Value 枚举 (20 种变体)

Re-exports§

pub use dialect::*;
pub use migration::*;

Modules§

access_control
行级和字段级权限控制
accessors
Accessors / Mutators + Attribute Casting
behaviors
行为系统(Behaviors)— 可插拔代码复用单元
data_permission
数据权限拦截器(Data Permission Interceptor)
dialect
不同数据库的方言抽象
dirty_attributes
脏字段追踪(Dirty Attributes)+ @DynamicInsert / @DynamicUpdate
dynamic_filter
动态 Filter(Hibernate @Filter / @FilterDef 风格)
dynamic_sql
XML/py_sql 动态 SQL 构造器(rbatis 风格)
entity_graph
Entity Graph + @BatchSize 批量抓取
find_with_related
find_with_related 关联查询流畅 API
guard
防全表 UPDATE/DELETE 攻击守卫(Safe SQL Guard)
hooks
钩子系统(Hooks)— 软删除 + 多租户
hydration_plugin
Hydration Modes + Plugin 拦截器链
i18n
国际化(i18n)支持
join_dsl
JoinDSL — 类型安全的 JOIN 语法(Diesel 风格)
json_query
JSON 字段查询增强
l2_cache
L2 二级缓存(Level-2 Cache)
lambda
Lambda 类型安全 Wrapper
migration
Migration system
observer
Observer + Event Subscriber — 模型生命周期观察者模式
optimistic_lock
乐观锁(Optimistic Locking)
phinx_migration
Phinx 风格 migration 链式 API
queryable
derive(Queryable) — 从 SELECT 结果自动派生结构体(Diesel 风格)
quick_query
快捷查询(Db::name 风格)
repository
Repository Pattern 仓储模式
result_map
ResultMap 高级映射 + Native Query + ResultSetMapping
retry
通用错误重试器
schema_gen
Diesel 风格 schema.rs 自动生成
shadow
双轨影子流量校验模块
sql_safety
SQL 安全工具:标识符与外键动作校验
type_handler
TypeHandler SPI — 自定义类型处理器注册
typed
强类型 AST 支持模块
typed_ast
强类型 AST 表达式层(Diesel 风格探索)

Macros§

define_columns
为 Model 定义一组类型安全的字段标记
schema
Compile-time SQL schema generator.
sql_string
Compile-time SQL validation macro.
typed_query
Diesel 风格强类型 AST 宏(与 sql_string! / query! 并存)。

Structs§

BelongsTo
多对一关系配置
BelongsToMany
多对多关系配置
Bytes
Re-export common types A cheaply cloneable and sliceable chunk of contiguous memory.
CacheStats
DateTime
ISO 8601 combined date and time with time zone.
ErrorContext
#6 修复:错误上下文链节点
HasMany
一对多关系配置
HasOne
一对一关系配置
MemoryCache
MorphMany
多态一对多配置(父模型侧)
MorphTo
多态反向配置(子模型侧)
MultiLevelCache
NegativeCache
带负缓存的缓存包装器
Pool
连接池核心实现
PoolConfig
PoolConfigBuilder
PoolStatus
PooledConnection
连接池中的连接条目,记录创建时间和最后使用时间
QueryBuilder
用于构造 SQL 查询的查询构造器
QueryBuilderWrapper
查询构造器包装类型,用于挂载作用域
TimestampFields
时间戳字段配置
TlsConfig
TLS 配置
TransactOptions
Transaction
事务对象,封装一个数据库事务
TransactionManager
事务管理器,管理多个事务
Utc
The UTC time zone. This is the most efficient time zone when you don’t need the local time. It is also used as an offset (which is also a dummy type).
WithRelation
关系预加载构造器

Enums§

AutoCommit
CacheError
缓存特有错误
CacheLookup
缓存查找结果
ColType
列类型枚举(v1.1.0 新增)
DbError
数据库错误类型
DbType
支持的数据库类型
IsolationLevel
PoolError
连接池特有错误
PoolEvent
连接池事件
PropagationBehavior
事务传播行为
Relation
模型间的关系描述
RelationError
关系操作错误类型
TlsVersion
TLS 版本
TransactionState
事务状态
TxError
事务特有错误
Value
数据库值类型

Constants§

DEFAULT_ACQUIRE_TIMEOUT
Default connection timeout in seconds
DEFAULT_BATCH_SIZE
Default batch size for bulk operations
DEFAULT_IDLE_TIMEOUT
Default idle timeout in seconds
DEFAULT_MAX_LIFETIME
Default max lifetime in seconds
DEFAULT_MAX_NESTING_DEPTH
H-8 默认最大嵌套深度
DEFAULT_MAX_SIZE
Default maximum pool size
DEFAULT_MIN_IDLE
Default minimum idle connections

Traits§

ActiveRecord
支持关系加载的模型 trait(ActiveRecord 模式)
Cache
Connection
数据库连接 trait
ConnectionFactory
连接工厂 trait,用于创建新连接
Deserialize
A data structure that can be deserialized from any data format supported by Serde.
Model
所有 ORM 模型必须实现的核心 trait
ModelExt
模型扩展 trait,提供额外功能
QueryBuilderExt
RelationAccess
ModelExt 的关系访问扩展方法
RelationLoader
可存储已加载关系数据的模型 trait
Scope
查询结果过滤作用域
Serialize
A data structure that can be serialized into any data format supported by Serde.

Functions§

is_deadlock_error
M-8 修复:检测错误字符串是否表示死锁
read_through
Read-through(同步版):缓存未命中时通过 loader 回源加载,写入缓存后返回
read_through_async
Read-through(异步版):缓存未命中时通过异步 loader 回源加载,写入缓存后返回
retry_on_deadlock
M-8 修复:在死锁时自动重试事务
rows_to_values
将查询结果行转换为 Vec<HashMap<String, Value>> 以便存入关系字段
set_error_hook
设置全局错误上报 hook
trigger_error_hook
触发错误 hook(在 DbError 创建/返回时调用)
value_to_json
将 Value 转换为 serde_json::Value(递归处理 Array)
write_around
Write-around(写旁路):仅写后端存储,同时失效缓存中的旧值
write_through
Write-through(同步版):同时写入缓存和后端存储(通过 writer 回调)
write_through_async
Write-through(异步版):同时写入缓存和异步后端存储

Type Aliases§

Boxed
Alias for Box<T>
CacheResult
Result type for cache operations
DbResult
Alias for Result<T, DbError>
PoolEventCallback
连接池事件回调
PoolResult
Result type for pool operations
QueryRows
查询结果行类型别名:避免 Connection::query 签名触发 clippy::type_complexity
QueryStreamItem
流式查询结果项类型别名:避免 Connection::query_stream 签名触发 clippy::type_complexity
QueryValues
位置式查询结果类型
Shared
Alias for Arc<T>
TxResult
Result type for transaction operations

Attribute Macros§

async_trait
Re-export async traits

Derive Macros§

Deserialize
Serialize