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(())

Reply via email to