use super::{GlideClient, GlideClusterClient};
use glide_core::client::Client as CoreClient;
use redis::{Cmd, FromRedisValue, RedisFuture, Value};
async fn dispatch_pipeline(
core: &mut CoreClient,
pipeline: &redis::Pipeline,
retry: Option<redis::PipelineRetryStrategy>,
) -> redis::RedisResult<Vec<Value>> {
if pipeline.is_atomic() {
let value = core.send_transaction(pipeline, None, None, true).await?;
Ok(vec![value])
} else {
let value = core
.send_pipeline(pipeline, None, true, None, retry.unwrap_or_default())
.await?;
match value {
Value::Array(items) => Ok(items),
other => Err(redis::RedisError::from((
redis::ErrorKind::ResponseError,
"unexpected non-array pipeline reply from glide-core",
format!("{other:?}"),
))),
}
}
}
struct PipelineConn {
core: CoreClient,
db: i64,
}
impl redis::aio::ConnectionLike for PipelineConn {
fn req_packed_command<'a>(&'a mut self, _cmd: &'a Cmd) -> RedisFuture<'a, Value> {
Box::pin(async {
Err(redis::RedisError::from((
redis::ErrorKind::ClientError,
"single-command dispatch is not supported on the pipeline adapter; \
use the unified command API",
)))
})
}
fn req_packed_commands<'a>(
&'a mut self,
cmd: &'a redis::Pipeline,
_offset: usize,
_count: usize,
pipeline_retry_strategy: Option<redis::PipelineRetryStrategy>,
) -> RedisFuture<'a, Vec<Value>> {
Box::pin(dispatch_pipeline(
&mut self.core,
cmd,
pipeline_retry_strategy,
))
}
fn get_db(&self) -> i64 {
self.db
}
fn is_closed(&self) -> bool {
false
}
}
mod sealed {
pub trait Sealed {}
impl Sealed for super::GlideClient {}
impl Sealed for super::GlideClusterClient {}
}
pub trait GlidePipelineTarget: sealed::Sealed + Sync {
#[doc(hidden)]
fn core_handle(&self) -> CoreClient;
#[doc(hidden)]
fn db_index(&self) -> i64;
}
impl GlidePipelineTarget for GlideClient {
fn core_handle(&self) -> CoreClient {
self.inner.clone()
}
fn db_index(&self) -> i64 {
self.db()
}
}
impl GlidePipelineTarget for GlideClusterClient {
fn core_handle(&self) -> CoreClient {
self.inner.clone()
}
fn db_index(&self) -> i64 {
0 }
}
pub trait PipelineExt {
fn query_glide<'a, C: GlidePipelineTarget, T: FromRedisValue + Send + 'a>(
&'a self,
con: &'a C,
) -> RedisFuture<'a, T>;
}
impl PipelineExt for redis::Pipeline {
fn query_glide<'a, C: GlidePipelineTarget, T: FromRedisValue + Send + 'a>(
&'a self,
con: &'a C,
) -> RedisFuture<'a, T> {
let mut conn = PipelineConn {
core: con.core_handle(),
db: con.db_index(),
};
Box::pin(async move { self.query_async(&mut conn).await })
}
}