-
-
Notifications
You must be signed in to change notification settings - Fork 0
github-actions[bot] edited this page Jul 27, 2026
·
4 revisions
定义面向业务的命令而不是通用CRUD。命令将领域语言带入你的API,并启用命令查询职责分离(CQRS)模式。
#[derive(Entity)] #[entity(table = "users", commands)] #[command(Register)] #[command(UpdateEmail: email)] #[command(Deactivate, requires_id)] pub struct User { #[id] pub id: Uuid, #[field(create, update, response)] pub email: String, #[field(create, response)] pub name: String, #[field(response)] pub active: bool, }
/// User上Register操作的命令载荷。 #[derive(Debug, Clone)] pub struct RegisterUser { pub email: String, pub name: String, } /// User上UpdateEmail操作的命令载荷。 #[derive(Debug, Clone)] pub struct UpdateEmailUser { pub id: Uuid, pub email: String, } /// User上Deactivate操作的命令载荷。 #[derive(Debug, Clone)] pub struct DeactivateUser { pub id: Uuid, }
/// User实体的命令枚举。 #[derive(Debug, Clone)] pub enum UserCommand { Register(RegisterUser), UpdateEmail(UpdateEmailUser), Deactivate(DeactivateUser), } impl EntityCommand for UserCommand { fn kind(&self) -> CommandKind { match self { UserCommand::Register(_) => CommandKind::Create, UserCommand::UpdateEmail(_) => CommandKind::Update, UserCommand::Deactivate(_) => CommandKind::Custom, } } fn name(&self) -> &'static str { match self { UserCommand::Register(_) => "Register", UserCommand::UpdateEmail(_) => "UpdateEmail", UserCommand::Deactivate(_) => "Deactivate", } } }
/// User命令执行的结果枚举。 #[derive(Debug, Clone)] pub enum UserCommandResult { Register(User), UpdateEmail(User), Deactivate, }
/// 处理User命令的异步trait。 #[async_trait] pub trait UserCommandHandler: Send + Sync { type Error: std::error::Error + Send + Sync; type Context: Send + Sync; /// 将命令分派到适当的处理器。 async fn handle(&self, cmd: UserCommand, ctx: &Self::Context) -> Result<UserCommandResult, Self::Error>; /// 处理Register命令。 async fn handle_register(&self, cmd: RegisterUser, ctx: &Self::Context) -> Result<User, Self::Error>; /// 处理UpdateEmail命令。 async fn handle_update_email(&self, cmd: UpdateEmailUser, ctx: &Self::Context) -> Result<User, Self::Error>; /// 处理Deactivate命令。 async fn handle_deactivate(&self, cmd: DeactivateUser, ctx: &Self::Context) -> Result<(), Self::Error>; }
使用所有 #[field(create)] 字段:
#[command(Register)] // 生成:RegisterUser { email, name }
仅使用指定字段(自动添加 requires_id):
#[command(UpdateEmail: email)] // 生成:UpdateEmailUser { id, email } #[command(UpdateProfile: name, bio, avatar)] // 生成:UpdateProfileUser { id, name, bio, avatar }
声明了 sets(...) 的命令直接写入指定的列,这些列无需标记 #[field(update)],因此既不会出现在公开的补丁 DTO 中,也不会进入 upsert 的 SET 列表:
#[command(VerifyPassport, payload(passport_provider), sets( passport_verified = "true", passport_verified_at = "NOW()" ))] // 生成:VerifyPassportUser { id, passport_provider } // pool.verify_passport(command) -> User
一条 UPDATE 只写入固定表达式和 payload 列,别的都不动。表达式与 #[column(default = "...")] 一样原样进入语句;列名在编译期对照实体校验。
启用 transactions 后,同一操作也出现在事务适配器上,可与其他写入一同提交;此处行不存在返回 Ok(None),与适配器其他方法一致。
let mut tx = pool.begin().await?; let verified = UserTransactionRepo::new(&mut tx) .verify_passport(VerifyPassportUser { id, passport_provider: Some("gov".into()) }) .await?; tx.commit().await?;
只添加ID字段:
#[command(Deactivate, requires_id)] // 生成:DeactivateUser { id } #[command(Delete, requires_id, kind = "delete")] // 生成:DeleteUser { id },返回 ()
使用外部结构体:
pub struct TransferPayload { pub from_account: Uuid, pub to_account: Uuid, pub amount: i64, } #[command(Transfer, payload = "TransferPayload")] // 直接使用TransferPayload
使用自定义结果类型:
pub struct TransferResult { pub transaction_id: Uuid, pub success: bool, } #[command(Transfer, payload = "TransferPayload", result = "TransferResult")] // 返回TransferResult而不是实体
控制使用哪些字段:
#[command(Create, source = "create")] // 使用#[field(create)]字段(默认) #[command(Modify, source = "update")] // 使用#[field(update)]字段(可选) #[command(Ping, source = "none")] // 无载荷字段
影响结果类型推断:
#[command(Create, kind = "create")] // 返回实体(默认) #[command(Update, kind = "update")] // 返回实体 #[command(Remove, kind = "delete")] // 返回 () #[command(Process, kind = "custom")] // 从source推断
use async_trait::async_trait; struct UserHandler { pool: PgPool, email_service: EmailService, } struct RequestContext { user_id: Option<Uuid>, correlation_id: Uuid, } #[async_trait] impl UserCommandHandler for UserHandler { type Error = AppError; type Context = RequestContext; async fn handle(&self, cmd: UserCommand, ctx: &Self::Context) -> Result<UserCommandResult, Self::Error> { match cmd { UserCommand::Register(c) => { let user = self.handle_register(c, ctx).await?; Ok(UserCommandResult::Register(user)) } UserCommand::UpdateEmail(c) => { let user = self.handle_update_email(c, ctx).await?; Ok(UserCommandResult::UpdateEmail(user)) } UserCommand::Deactivate(c) => { self.handle_deactivate(c, ctx).await?; Ok(UserCommandResult::Deactivate) } } } async fn handle_register(&self, cmd: RegisterUser, ctx: &Self::Context) -> Result<User, Self::Error> { // 验证 if cmd.email.is_empty() { return Err(AppError::Validation("邮箱必填".into())); } // 创建用户 let user = User { id: Uuid::now_v7(), email: cmd.email.to_lowercase(), name: cmd.name, active: true, }; // 持久化 sqlx::query( "INSERT INTO users (id, email, name, active) VALUES (1,ドル 2,ドル 3,ドル 4ドル)" ) .bind(user.id) .bind(&user.email) .bind(&user.name) .bind(user.active) .execute(&self.pool) .await?; // 副作用 self.email_service.send_welcome(&user.email).await?; Ok(user) } async fn handle_update_email(&self, cmd: UpdateEmailUser, ctx: &Self::Context) -> Result<User, Self::Error> { // 授权检查 if ctx.user_id != Some(cmd.id) { return Err(AppError::Forbidden("无法更新其他用户的邮箱".into())); } // 更新 let user: User = sqlx::query_as( "UPDATE users SET email = 1ドル WHERE id = 2ドル RETURNING *" ) .bind(&cmd.email.to_lowercase()) .bind(cmd.id) .fetch_one(&self.pool) .await?; // 发送验证 self.email_service.send_verification(&user.email).await?; Ok(user) } async fn handle_deactivate(&self, cmd: DeactivateUser, ctx: &Self::Context) -> Result<(), Self::Error> { sqlx::query("UPDATE users SET active = false WHERE id = 1ドル") .bind(cmd.id) .execute(&self.pool) .await?; Ok(()) } }
async fn register_user( handler: &impl UserCommandHandler, email: String, name: String, ) -> Result<User, AppError> { let cmd = RegisterUser { email, name }; let ctx = RequestContext { user_id: None, correlation_id: Uuid::new_v4(), }; match handler.handle(UserCommand::Register(cmd), &ctx).await? { UserCommandResult::Register(user) => Ok(user), _ => unreachable!(), } } // 或直接调用特定handler async fn update_email( handler: &impl UserCommandHandler, user_id: Uuid, new_email: String, ctx: &RequestContext, ) -> Result<User, AppError> { let cmd = UpdateEmailUser { id: user_id, email: new_email, }; handler.handle_update_email(cmd, ctx).await }
所有命令枚举都实现 EntityCommand trait:
use entity_derive::{EntityCommand, CommandKind}; let cmd = UserCommand::Register(register_data); // 获取命令元数据 assert_eq!(cmd.name(), "Register"); assert!(matches!(cmd.kind(), CommandKind::Create)); // 模式匹配 match cmd.kind() { CommandKind::Create => println!("创建实体"), CommandKind::Update => println!("更新实体"), CommandKind::Delete => println!("删除实体"), CommandKind::Custom => println!("自定义操作"), }
当同时启用 commands 和 hooks 时:
#[derive(Entity)] #[entity(table = "orders", commands, hooks)] #[command(Place)] #[command(Cancel, requires_id)] pub struct Order { /* ... */ }
生成的钩子:
#[async_trait] pub trait OrderHooks: Send + Sync { type Error: std::error::Error + Send + Sync; // 标准CRUD钩子... // 命令特定钩子 async fn before_command(&self, cmd: &OrderCommand) -> Result<(), Self::Error>; async fn after_command(&self, cmd: &OrderCommand, result: &OrderCommandResult) -> Result<(), Self::Error>; }
-
领域语言 — 使用业务术语:
RegisterUser而不是CreateUser - 单一职责 — 一个命令 = 一个业务操作
- 明确意图 — 命令名应描述操作
- 在handler中验证 — 将验证逻辑保留在命令handler中
- 尽可能幂等 — 设计命令以便安全重试
- 使用上下文 — 通过上下文传递请求元数据(用户、关联ID)
命令是CQRS的一半。与投影结合用于查询端:
#[derive(Entity)] #[entity(table = "orders", commands)] #[projection(Summary: id, status, total_cents, created_at)] #[projection(Details: id, status, items, shipping_address, total_cents)] #[command(Place)] #[command(Ship, requires_id)] #[command(Cancel, requires_id)] pub struct Order { /* ... */ } // 命令(写入端) let result = handler.handle(OrderCommand::Place(place_order), &ctx).await?; // 查询(读取端) let summary = repo.find_by_id_summary(order_id).await?; let details = repo.find_by_id_details(order_id).await?;
🇬🇧 English | 🇷🇺 Русский | 🇰🇷 한국어 | 🇪🇸 Español | 🇨🇳 中文
Getting Started
Features
Advanced
Начало работы
Возможности
Продвинутое
시작하기
기능
고급
Comenzando
Características
Avanzado
入门
功能
高级