You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

为Rust Diesel-async异步闭包事务确定正确特征边界

调整Diesel-async中DbContext的txw方法以支持异步闭包直接调用

需求说明

希望让自定义DbContext的txw方法支持直接传入异步闭包,无需手动装箱Future,目标调用方式如下:

async fn txw_custom(db: DbContext) -> DbResult<()> {
    db.txw(async |conn| { Ok(conn.batch_execute("SELECT 1").await?) }).await
}

当前实现依赖BoxFuture,调用时必须手动用Box::pin包装异步闭包,代码冗余。

原实现代码

use diesel_async::{
    AsyncConnection, AsyncMigrationHarness, AsyncPgConnection,
    pooled_connection::{
        AsyncDieselConnectionManager, ManagerConfig, PoolableConnection, RecyclingMethod,
        bb8::{Pool, PooledConnection, RunError},
    },
};
use futures_util::{FutureExt, future::BoxFuture, try_join};


#[derive(Debug, thiserror::Error)]
pub enum DbError {
    #[error("failed to run migrations because {0}")]
    MigrationError(String),

    #[error("something went wrong when getting connection from the pool because {0}")]
    PoolConnection(#[from] bb8::RunError),

    #[error("failed to execute query because {0}")]
    FailedQuery(#[from] diesel::result::Error),

    #[error("failed to {0}")]
    Other(String),
}

pub type DbResult<T> = std::result::Result<T, DbError>;


#[derive(Clone)]
pub struct DbContext {
    primary_pool: Pool<AsyncPgConnection>,
    replica_pool: Pool<AsyncPgConnection>,
}

impl DbContext {
    pub async fn primary(&self) -> DbResult<PooledConnection<'_, AsyncPgConnection>> {
        self.primary_pool
            .get()
            .await
            .map_err(DbError::PoolConnection)
    }

    pub async fn replica(&self) -> DbResult<PooledConnection<'_, AsyncPgConnection>> {
        self.replica_pool
            .get()
            .await
            .map_err(DbError::PoolConnection)
    }

    /// Run a transaction on a primary connection
    ///
    /// How to use this function:
    /// ```rust
    /// use db::{DbContext, DbResult};
    /// use diesel_async::SimpleAsyncConnection;
    /// async fn txw_custom(db: DbContext) -> DbResult<()> {
    ///     db.txw(|conn| Box::pin(async move { Ok(conn.batch_execute("SELECT 1").await?) })).await
    /// }
    /// ```
    pub async fn txw<F, T>(&self, f: F) -> DbResult<T>
    where
        F: FnOnce(&mut AsyncPgConnection) -> BoxFuture<DbResult<T>> + Send,
        T: Send,
    {
        self.primary()
            .await?
            .transaction(async |c| f(c).await)
            .await
    }

}

修改方案

核心是调整txw方法的泛型约束,不再要求返回BoxFuture,而是直接关联闭包返回的Future类型,同时添加必要的Send约束以适配异步事务的要求。

步骤1:引入Future trait

在顶部添加Future的导入:

use futures_util::Future;

步骤2:修改txw方法的泛型约束

替换原txw方法为以下实现:

/// Run a transaction on a primary connection
///
/// How to use this function:
/// ```rust
/// use db::{DbContext, DbResult};
/// use diesel_async::SimpleAsyncConnection;
/// async fn txw_custom(db: DbContext) -> DbResult<()> {
///     db.txw(async |conn| { Ok(conn.batch_execute("SELECT 1").await?) }).await
/// }
/// ```
pub async fn txw<F, Fut, T>(&self, f: F) -> DbResult<T>
where
    F: FnOnce(&mut AsyncPgConnection) -> Fut + Send,
    Fut: Future<Output = DbResult<T>> + Send,
    T: Send,
{
    self.primary()
        .await?
        .transaction(async |c| f(c).await)
        .await
}

说明

  • 新增泛型参数Fut,代表异步闭包返回的具体Future类型,彻底避免了装箱操作
  • 约束Fut实现Future<Output = DbResult<T>>和Send,确保Future可以在异步事务的多线程上下文安全执行
  • 闭包F保留FnOnce(&mut AsyncPgConnection) -> Fut + Send约束,保证可以被安全传递到异步执行环境

修改完成后,即可直接使用目标方式调用txw方法,无需手动装箱Future。

内容的提问来源于stack exchange,提问作者Vraj Shah

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.17 18:57:51