This is an automated email from the ASF dual-hosted git repository. xuanwo pushed a commit to branch reactor-redis in repository https://gitbox.apache.org/repos/asf/opendal.git
commit bf5719fcaae8ba02c09d1529f01ac51406f37762 Author: Xuanwo <[email protected]> AuthorDate: Tue Mar 11 22:45:48 2025 +0800 refactor(services/redis): Implement ConnectionLike for RedisConnection Signed-off-by: Xuanwo <[email protected]> --- core/Cargo.lock | 141 +++++++++++++++++++++++++++++++------ core/Cargo.toml | 132 +++++++++++++++++----------------- core/src/services/redis/backend.rs | 66 +++++++++++++++-- core/src/services/redis/core.rs | 93 +++++------------------- 4 files changed, 263 insertions(+), 169 deletions(-) diff --git a/core/Cargo.lock b/core/Cargo.lock index 9aac4604f..572a8aadd 100644 --- a/core/Cargo.lock +++ b/core/Cargo.lock @@ -68,7 +68,7 @@ dependencies = [ "getrandom 0.2.15", "once_cell", "version_check", - "zerocopy", + "zerocopy 0.7.35", ] [[package]] @@ -1129,9 +1129,9 @@ dependencies = [ [[package]] name = "backon" -version = "1.3.0" +version = "1.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ba5289ec98f68f28dd809fd601059e6aa908bb8f6108620930828283d4ee23d7" +checksum = "49fef586913a57ff189f25c9b3d034356a5bf6b3fa9a7f067588fe1698ba1f5d" dependencies = [ "fastrand", "gloo-timers", @@ -2044,6 +2044,16 @@ dependencies = [ "libc", ] +[[package]] +name = "core-foundation" +version = "0.10.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b55271e5c8c478ad3f38ad24ef34923091e0548492a266d19b3c0b4d82574c63" +dependencies = [ + "core-foundation-sys", + "libc", +] + [[package]] name = "core-foundation-sys" version = "0.8.7" @@ -3251,6 +3261,18 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "getrandom" +version = "0.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "43a49c392881ce6d5c3b8cb70f98717b7c07aabbdff06687b9030dbfbe2725f8" +dependencies = [ + "cfg-if", + "libc", + "wasi 0.13.3+wasi-0.2.2", + "windows-targets 0.52.6", +] + [[package]] name = "ghac" version = "0.2.0" @@ -4936,7 +4958,7 @@ dependencies = [ "openssl-probe", "openssl-sys", "schannel", - "security-framework", + "security-framework 2.11.1", "security-framework-sys", "tempfile", ] @@ -6057,7 +6079,7 @@ version = "0.2.20" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "77957b295656769bb8ad2b6a6b09d897d94f05c41b069aede1fcdaa675eaea04" dependencies = [ - "zerocopy", + "zerocopy 0.7.35", ] [[package]] @@ -6499,6 +6521,17 @@ dependencies = [ "rand_core 0.6.4", ] +[[package]] +name = "rand" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3779b94aeb87e8bd4e834cee3650289ee9e0d5677f976ecdb6d219e5f4f6cd94" +dependencies = [ + "rand_chacha 0.9.0", + "rand_core 0.9.3", + "zerocopy 0.8.23", +] + [[package]] name = "rand_chacha" version = "0.2.2" @@ -6519,6 +6552,16 @@ dependencies = [ "rand_core 0.6.4", ] +[[package]] +name = "rand_chacha" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d3022b5f1df60f26e1ffddd6c66e8aa15de382ae63b3a0c1bfc0e4d3e3f325cb" +dependencies = [ + "ppv-lite86", + "rand_core 0.9.3", +] + [[package]] name = "rand_core" version = "0.5.1" @@ -6537,6 +6580,15 @@ dependencies = [ "getrandom 0.2.15", ] +[[package]] +name = "rand_core" +version = "0.9.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "99d9a13982dcf210057a8a78572b2217b667c3beacbf3a0d8b454f6f82837d38" +dependencies = [ + "getrandom 0.3.1", +] + [[package]] name = "rand_hc" version = "0.2.0" @@ -6598,30 +6650,27 @@ dependencies = [ [[package]] name = "redis" -version = "0.27.6" +version = "0.29.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "09d8f99a4090c89cc489a94833c901ead69bfbf3877b4867d5482e321ee875bc" +checksum = "8034fb926579ff49d3fe58d288d5dcb580bf11e9bccd33224b45adebf0fd0c23" dependencies = [ "arc-swap", - "async-trait", "backon", "bytes", "combine", "crc16", - "futures", + "futures-channel", + "futures-sink", "futures-util", - "itertools 0.13.0", "itoa", "log", "native-tls", "num-bigint", "percent-encoding", "pin-project-lite", - "rand 0.8.5", + "rand 0.9.0", "rustls 0.23.20", - "rustls-native-certs 0.7.3", - "rustls-pemfile 2.2.0", - "rustls-pki-types", + "rustls-native-certs 0.8.1", "ryu", "sha1_smol", "socket2", @@ -7176,20 +7225,19 @@ dependencies = [ "openssl-probe", "rustls-pemfile 1.0.4", "schannel", - "security-framework", + "security-framework 2.11.1", ] [[package]] name = "rustls-native-certs" -version = "0.7.3" +version = "0.8.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e5bfb394eeed242e909609f56089eecfe5fda225042e8b171791b9c95f5931e5" +checksum = "7fcff2dd52b58a8d98a70243663a0d234c4e2b79235637849d15913394a247d3" dependencies = [ "openssl-probe", - "rustls-pemfile 2.2.0", "rustls-pki-types", "schannel", - "security-framework", + "security-framework 3.2.0", ] [[package]] @@ -7346,7 +7394,20 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "897b2245f0b511c87893af39b033e5ca9cce68824c4d7e7630b5a1d339658d02" dependencies = [ "bitflags 2.6.0", - "core-foundation", + "core-foundation 0.9.4", + "core-foundation-sys", + "libc", + "security-framework-sys", +] + +[[package]] +name = "security-framework" +version = "3.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "271720403f46ca04f7ba6f55d438f8bd878d6b8ca0a1046e8228c4145bcbb316" +dependencies = [ + "bitflags 2.6.0", + "core-foundation 0.10.0", "core-foundation-sys", "libc", "security-framework-sys", @@ -9208,6 +9269,15 @@ version = "0.11.0+wasi-snapshot-preview1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9c8d87e72b64a3b4db28d11ce29237c246188f4f51057d65a7eab63b7987e423" +[[package]] +name = "wasi" +version = "0.13.3+wasi-0.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "26816d2e1a4a36a2940b96c5296ce403917633dff8f3440e9b236ed6f6bacad2" +dependencies = [ + "wit-bindgen-rt", +] + [[package]] name = "wasite" version = "0.1.0" @@ -9709,6 +9779,15 @@ dependencies = [ "windows-sys 0.48.0", ] +[[package]] +name = "wit-bindgen-rt" +version = "0.33.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3268f3d866458b787f390cf61f4bbb563b922d091359f9608842999eaee3943c" +dependencies = [ + "bitflags 2.6.0", +] + [[package]] name = "write16" version = "1.0.0" @@ -9804,7 +9883,16 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1b9b4fd18abc82b8136838da5d50bae7bdea537c574d8dc1a34ed098d6c166f0" dependencies = [ "byteorder", - "zerocopy-derive", + "zerocopy-derive 0.7.35", +] + +[[package]] +name = "zerocopy" +version = "0.8.23" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fd97444d05a4328b90e75e503a34bad781f14e28a823ad3557f0750df1ebcbc6" +dependencies = [ + "zerocopy-derive 0.8.23", ] [[package]] @@ -9818,6 +9906,17 @@ dependencies = [ "syn 2.0.95", ] +[[package]] +name = "zerocopy-derive" +version = "0.8.23" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6352c01d0edd5db859a63e2605f4ea3183ddbd15e2c4a9e7d32184df75e4f154" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.95", +] + [[package]] name = "zerofrom" version = "0.1.5" diff --git a/core/Cargo.toml b/core/Cargo.toml index fe60c1ac0..f14aa7bc3 100644 --- a/core/Cargo.toml +++ b/core/Cargo.toml @@ -56,16 +56,16 @@ default = ["reqwest/rustls-tls", "executors-tokio", "services-memory"] # # You should never enable this feature unless you are developing opendal. tests = [ - "dep:rand", - "dep:sha2", - "dep:dotenvy", - "layers-blocking", - "services-azblob", - "services-fs", - "services-http", - "services-memory", - "internal-tokio-rt", - "services-s3", + "dep:rand", + "dep:sha2", + "dep:dotenvy", + "layers-blocking", + "services-azblob", + "services-fs", + "services-http", + "services-memory", + "internal-tokio-rt", + "services-s3", ] # Enable path cache. @@ -109,29 +109,29 @@ services-aliyun-drive = [] services-alluxio = [] services-atomicserver = ["dep:atomic_lib"] services-azblob = [ - "dep:sha2", - "dep:reqsign", - "reqsign?/services-azblob", - "reqsign?/reqwest_request", + "dep:sha2", + "dep:reqsign", + "reqsign?/services-azblob", + "reqsign?/reqwest_request", ] services-azdls = [ - "dep:reqsign", - "reqsign?/services-azblob", - "reqsign?/reqwest_request", + "dep:reqsign", + "reqsign?/services-azblob", + "reqsign?/reqwest_request", ] services-azfile = [ - "dep:reqsign", - "reqsign?/services-azblob", - "reqsign?/reqwest_request", + "dep:reqsign", + "reqsign?/services-azblob", + "reqsign?/reqwest_request", ] services-b2 = [] services-cacache = ["dep:cacache"] services-cloudflare-kv = [] services-compfs = ["dep:compio"] services-cos = [ - "dep:reqsign", - "reqsign?/services-tencent", - "reqsign?/reqwest_request", + "dep:reqsign", + "reqsign?/services-tencent", + "reqsign?/reqwest_request", ] services-d1 = [] services-dashmap = ["dep:dashmap"] @@ -142,9 +142,9 @@ services-foundationdb = ["dep:foundationdb"] services-fs = ["tokio/fs", "internal-tokio-rt"] services-ftp = ["dep:suppaftp", "dep:bb8", "dep:async-tls"] services-gcs = [ - "dep:reqsign", - "reqsign?/services-google", - "reqsign?/reqwest_request", + "dep:reqsign", + "reqsign?/services-google", + "reqsign?/reqwest_request", ] services-gdrive = ["internal-path-cache"] services-ghac = ["dep:ghac", "dep:prost", "services-azblob"] @@ -168,15 +168,15 @@ services-monoiofs = ["dep:monoio", "dep:flume"] services-mysql = ["dep:sqlx", "sqlx?/mysql"] services-nebula-graph = ["dep:rust-nebula", "dep:bb8", "dep:snowflaked"] services-obs = [ - "dep:reqsign", - "reqsign?/services-huaweicloud", - "reqsign?/reqwest_request", + "dep:reqsign", + "reqsign?/services-huaweicloud", + "reqsign?/reqwest_request", ] services-onedrive = [] services-oss = [ - "dep:reqsign", - "reqsign?/services-aliyun", - "reqsign?/reqwest_request", + "dep:reqsign", + "reqsign?/services-aliyun", + "reqsign?/reqwest_request", ] services-pcloud = [] services-persy = ["dep:persy", "internal-tokio-rt"] @@ -186,10 +186,10 @@ services-redis = ["dep:redis", "dep:bb8", "redis?/tokio-rustls-comp"] services-redis-native-tls = ["services-redis", "redis?/tokio-native-tls-comp"] services-rocksdb = ["dep:rocksdb", "internal-tokio-rt"] services-s3 = [ - "dep:reqsign", - "reqsign?/services-aws", - "reqsign?/reqwest_request", - "dep:crc32c", + "dep:reqsign", + "reqsign?/services-aws", + "reqsign?/reqwest_request", + "dep:crc32c", ] services-seafile = [] services-sftp = ["dep:openssh", "dep:openssh-sftp-client", "dep:bb8"] @@ -234,12 +234,12 @@ backon = { version = "1.2", features = ["tokio-sleep"] } base64 = "0.22" bytes = "1.6" chrono = { version = "0.4.28", default-features = false, features = [ - "clock", - "std", + "clock", + "std", ] } futures = { version = "0.3", default-features = false, features = [ - "std", - "async-await", + "std", + "async-await", ] } http = "1.1" log = "0.4" @@ -250,7 +250,7 @@ once_cell = "1" percent-encoding = "2" quick-xml = { version = "0.36", features = ["serialize", "overlapped-lists"] } reqwest = { version = "0.12.2", features = [ - "stream", + "stream", ], default-features = false } serde = { version = "1", features = ["derive"] } serde_json = "1" @@ -270,7 +270,7 @@ prost = { version = "0.13", optional = true } sha1 = { version = "0.10.6", optional = true } sha2 = { version = "0.10", optional = true } sqlx = { version = "0.8.0", features = [ - "runtime-tokio-rustls", + "runtime-tokio-rustls", ], optional = true } # For http based services. @@ -283,8 +283,8 @@ ouroboros = { version = "0.18.4", optional = true } atomic_lib = { version = "0.39.0", optional = true } # for services-cacache cacache = { version = "13.0", default-features = false, features = [ - "tokio-runtime", - "mmap", + "tokio-runtime", + "mmap", ], optional = true } # for services-dashmap dashmap = { version = "6", optional = true } @@ -292,8 +292,8 @@ dashmap = { version = "6", optional = true } etcd-client = { version = "0.14", optional = true, features = ["tls"] } # for services-foundationdb foundationdb = { version = "0.9.0", features = [ - "embedded-fdb-include", - "fdb-7_3", + "embedded-fdb-include", + "fdb-7_3", ], optional = true } # for services-hdfs hdrs = { version = "0.3.2", optional = true, features = ["async_file"] } @@ -309,18 +309,18 @@ mongodb-internal-macros = { version = "3.2.2", optional = true } # for services-sftp openssh = { version = "0.11.0", optional = true } openssh-sftp-client = { version = "0.15.2", optional = true, features = [ - "openssh", - "tracing", + "openssh", + "tracing", ] } # for services-persy persy = { version = "1.4.6", optional = true } # for services-redb redb = { version = "2", optional = true } # for services-redis -redis = { version = "0.27", features = [ - "cluster-async", - "tokio-comp", - "connection-manager", +redis = { version = "0.29", features = [ + "cluster-async", + "tokio-comp", + "connection-manager", ], optional = true } # for services-rocksdb rocksdb = { version = "0.21.0", default-features = false, optional = true } @@ -328,9 +328,9 @@ rocksdb = { version = "0.21.0", default-features = false, optional = true } sled = { version = "0.34.7", optional = true } # for services-ftp suppaftp = { version = "6.0.3", default-features = false, features = [ - "async-secure", - "rustls", - "async-rustls", + "async-secure", + "rustls", + "async-rustls", ], optional = true } # for services-tikv tikv-client = { version = "0.3.0", optional = true, default-features = false } @@ -340,10 +340,10 @@ hdfs-native = { version = "0.10", optional = true } surrealdb = { version = "2", optional = true, features = ["protocol-http"] } # for services-compfs compio = { version = "0.12.0", optional = true, features = [ - "runtime", - "bytes", - "polling", - "dispatcher", + "runtime", + "bytes", + "polling", + "dispatcher", ] } # for services-s3 crc32c = { version = "0.6.6", optional = true } @@ -353,10 +353,10 @@ snowflaked = { version = "1", optional = true, features = ["sync"] } # for services-monoiofs flume = { version = "0.11", optional = true } monoio = { version = "0.2.4", optional = true, features = [ - "sync", - "mkdirat", - "unlinkat", - "renameat", + "sync", + "mkdirat", + "unlinkat", + "renameat", ] } # Layers @@ -395,7 +395,7 @@ fastrace = { version = "0.7", features = ["enable"] } fastrace-jaeger = "0.7" libtest-mimic = "0.8" opentelemetry = { version = "0.28", default-features = false, features = [ - "trace", + "trace", ] } opentelemetry-otlp = { version = "0.28", features = ["grpc-tonic"] } opentelemetry_sdk = { version = "0.28", features = ["rt-tokio"] } @@ -406,6 +406,6 @@ size = "0.4" tokio = { version = "1.27", features = ["fs", "macros", "rt-multi-thread"] } tracing-opentelemetry = "0.29.0" tracing-subscriber = { version = "0.3", features = [ - "env-filter", - "tracing-log", + "env-filter", + "tracing-log", ] } diff --git a/core/src/services/redis/backend.rs b/core/src/services/redis/backend.rs index c9eb20426..ddbf28976 100644 --- a/core/src/services/redis/backend.rs +++ b/core/src/services/redis/backend.rs @@ -16,26 +16,27 @@ // under the License. use bb8::RunError; +use bytes::Bytes; use http::Uri; use redis::cluster::ClusterClient; use redis::cluster::ClusterClientBuilder; -use redis::Client; use redis::ConnectionAddr; use redis::ConnectionInfo; use redis::ProtocolVersion; use redis::RedisConnectionInfo; +use redis::{AsyncCommands, Client}; use std::fmt::Debug; use std::fmt::Formatter; use std::path::PathBuf; use std::time::Duration; use tokio::sync::OnceCell; +use super::core::*; use crate::raw::adapters::kv; use crate::raw::*; use crate::services::RedisConfig; use crate::*; -use super::core::*; const DEFAULT_REDIS_ENDPOINT: &str = "tcp://127.0.0.1:6379"; const DEFAULT_REDIS_PORT: u16 = 6379; @@ -345,26 +346,77 @@ impl kv::Adapter for Adapter { async fn get(&self, key: &str) -> Result<Option<Buffer>> { let mut conn = self.conn().await?; - let result = conn.get(key).await?; - Ok(result) + let result: Option<Bytes> = conn.get(key).await.map_err(format_redis_error)?; + Ok(result.map(Buffer::from)) } async fn set(&self, key: &str, value: Buffer) -> Result<()> { let mut conn = self.conn().await?; let value = value.to_vec(); - conn.set(key, value, self.default_ttl).await?; + if let Some(dur) = self.default_ttl { + let _: () = conn + .set_ex(key, value, dur.as_secs()) + .await + .map_err(format_redis_error)?; + } else { + let _: () = conn.set(key, value).await.map_err(format_redis_error)?; + } Ok(()) } async fn delete(&self, key: &str) -> Result<()> { let mut conn = self.conn().await?; - conn.delete(key).await?; + let _: () = conn.del(key).await.map_err(format_redis_error)?; Ok(()) } async fn append(&self, key: &str, value: &[u8]) -> Result<()> { let mut conn = self.conn().await?; - conn.append(key, value).await?; + let _: () = conn.append(key, value).await.map_err(format_redis_error)?; Ok(()) } } + +// impl Access for Adapter { +// type Reader = (); +// type Writer = (); +// type Lister = (); +// type Deleter = (); +// type BlockingReader = (); +// type BlockingWriter = (); +// type BlockingLister = (); +// type BlockingDeleter = (); +// +// fn info(&self) -> Arc<AccessorInfo> { +// todo!() +// } +// +// async fn stat(&self, path: &str, args: OpStat) -> Result<RpStat> { +// let mut conn = self.conn().await?; +// let size = conn.strlen().await?; +// } +// +// async fn read(&self, path: &str, args: OpRead) -> Result<(RpRead, Self::Reader)> { +// todo!() +// } +// +// async fn write(&self, path: &str, args: OpWrite) -> Result<(RpWrite, Self::Writer)> { +// todo!() +// } +// +// async fn delete(&self) -> Result<(RpDelete, Self::Deleter)> { +// todo!() +// } +// +// async fn list(&self, path: &str, args: OpList) -> Result<(RpList, Self::Lister)> { +// todo!() +// } +// +// async fn copy(&self, from: &str, to: &str, args: OpCopy) -> Result<RpCopy> { +// todo!() +// } +// +// async fn rename(&self, from: &str, to: &str, args: OpRename) -> Result<RpRename> { +// todo!() +// } +// } diff --git a/core/src/services/redis/core.rs b/core/src/services/redis/core.rs index 041ed8716..f5eab5cc7 100644 --- a/core/src/services/redis/core.rs +++ b/core/src/services/redis/core.rs @@ -15,95 +15,49 @@ // specific language governing permissions and limitations // under the License. -use crate::Buffer; use crate::Error; use crate::ErrorKind; use redis::aio::ConnectionLike; use redis::aio::ConnectionManager; - use redis::cluster::ClusterClient; use redis::cluster_async::ClusterConnection; -use redis::from_redis_value; use redis::AsyncCommands; use redis::Client; use redis::RedisError; - -use std::time::Duration; +use redis::{Cmd, Pipeline, RedisFuture, Value}; #[derive(Clone)] pub enum RedisConnection { Normal(ConnectionManager), Cluster(ClusterConnection), } -impl RedisConnection { - pub async fn get(&mut self, key: &str) -> crate::Result<Option<Buffer>> { - let result: Option<bytes::Bytes> = match self { - RedisConnection::Normal(ref mut conn) => { - conn.get(key).await.map_err(format_redis_error) - } - RedisConnection::Cluster(ref mut conn) => { - conn.get(key).await.map_err(format_redis_error) - } - }?; - Ok(result.map(Buffer::from)) - } - pub async fn set( - &mut self, - key: &str, - value: Vec<u8>, - ttl: Option<Duration>, - ) -> crate::Result<()> { - let value = value.to_vec(); - if let Some(ttl) = ttl { - match self { - RedisConnection::Normal(ref mut conn) => conn - .set_ex(key, value, ttl.as_secs()) - .await - .map_err(format_redis_error)?, - RedisConnection::Cluster(ref mut conn) => conn - .set_ex(key, value, ttl.as_secs()) - .await - .map_err(format_redis_error)?, - } - } else { - match self { - RedisConnection::Normal(ref mut conn) => { - conn.set(key, value).await.map_err(format_redis_error)? - } - RedisConnection::Cluster(ref mut conn) => { - conn.set(key, value).await.map_err(format_redis_error)? - } - } +impl ConnectionLike for RedisConnection { + fn req_packed_command<'a>(&'a mut self, cmd: &'a Cmd) -> RedisFuture<'a, Value> { + match self { + RedisConnection::Normal(conn) => conn.req_packed_command(cmd), + RedisConnection::Cluster(conn) => conn.req_packed_command(cmd), } - - Ok(()) } - pub async fn delete(&mut self, key: &str) -> crate::Result<()> { + fn req_packed_commands<'a>( + &'a mut self, + cmd: &'a Pipeline, + offset: usize, + count: usize, + ) -> RedisFuture<'a, Vec<Value>> { match self { - RedisConnection::Normal(ref mut conn) => { - let _: () = conn.del(key).await.map_err(format_redis_error)?; - } - RedisConnection::Cluster(ref mut conn) => { - let _: () = conn.del(key).await.map_err(format_redis_error)?; - } + RedisConnection::Normal(conn) => conn.req_packed_commands(cmd, offset, count), + RedisConnection::Cluster(conn) => conn.req_packed_commands(cmd, offset, count), } - - Ok(()) } - pub async fn append(&mut self, key: &str, value: &[u8]) -> crate::Result<()> { + fn get_db(&self) -> i64 { match self { - RedisConnection::Normal(ref mut conn) => { - () = conn.append(key, value).await.map_err(format_redis_error)?; - } - RedisConnection::Cluster(ref mut conn) => { - () = conn.append(key, value).await.map_err(format_redis_error)?; - } + RedisConnection::Normal(conn) => conn.get_db(), + RedisConnection::Cluster(conn) => conn.get_db(), } - Ok(()) } } @@ -136,18 +90,7 @@ impl bb8::ManageConnection for RedisConnectionManager { } async fn is_valid(&self, conn: &mut Self::Connection) -> Result<(), Self::Error> { - let pong_value = match conn { - RedisConnection::Normal(ref mut conn) => conn - .send_packed_command(&redis::cmd("PING")) - .await - .map_err(format_redis_error)?, - - RedisConnection::Cluster(ref mut conn) => conn - .req_packed_command(&redis::cmd("PING")) - .await - .map_err(format_redis_error)?, - }; - let pong: String = from_redis_value(&pong_value).map_err(format_redis_error)?; + let pong: String = conn.ping().await.map_err(format_redis_error)?; if pong == "PONG" { Ok(())
