mod common;
use glide::{AsyncCommands, ConnectionManagementCommands, CustomCommand, Route};
#[tokio::test]
async fn cluster_set_get_routed_by_key() {
let cluster = cluster_or_skip!();
let client = match cluster.client().await {
Some(c) => c,
None => {
eprintln!("SKIP: cluster client connect failed");
return;
}
};
for i in 0..50 {
let k = format!("clusterkey:{i}");
let _: () = client.set(&k, "v").await.unwrap();
let got: Option<String> = client.get(&k).await.unwrap();
assert_eq!(got.as_deref(), Some("v"));
}
}
#[tokio::test]
async fn cluster_ping_all_primaries() {
let cluster = cluster_or_skip!();
let client = match cluster.client().await {
Some(c) => c,
None => {
eprintln!("SKIP: cluster client connect failed");
return;
}
};
let reply = client
.custom_command_with_route(&["PING"], Route::AllPrimaries)
.await
.unwrap();
assert!(!matches!(reply, glide::Value::Nil));
}
#[tokio::test]
async fn cluster_info_reports_ok() {
let cluster = cluster_or_skip!();
let client = match cluster.client().await {
Some(c) => c,
None => {
eprintln!("SKIP: cluster client connect failed");
return;
}
};
let mut info = String::new();
for _ in 0..10 {
let reply = client
.custom_command_with_route(&["CLUSTER", "INFO"], Route::RandomNode)
.await
.unwrap();
info = glide::value::to_string(reply).unwrap_or_default();
if info.contains("cluster_state:ok") {
break;
}
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
}
assert!(
info.contains("cluster_state:ok"),
"cluster never reported ok: {info}"
);
}
#[tokio::test]
async fn cluster_del_and_exists() {
let cluster = cluster_or_skip!();
let client = match cluster.client().await {
Some(c) => c,
None => {
eprintln!("SKIP: cluster client connect failed");
return;
}
};
let k = "cluster:delkey";
let _: () = client.set(k, "v").await.unwrap();
let exists: i64 = client.exists(k).await.unwrap();
assert_eq!(exists, 1);
let deleted: i64 = client.del(k).await.unwrap();
assert_eq!(deleted, 1);
let exists: i64 = client.exists(k).await.unwrap();
assert_eq!(exists, 0);
}
#[tokio::test]
async fn cluster_incr() {
let cluster = cluster_or_skip!();
let client = match cluster.client().await {
Some(c) => c,
None => {
eprintln!("SKIP: cluster client connect failed");
return;
}
};
let k = "cluster:counter";
let v: i64 = client.incr(k, 1i64).await.unwrap();
assert_eq!(v, 1);
let v: i64 = client.incr(k, 4i64).await.unwrap();
assert_eq!(v, 5);
}
#[tokio::test]
async fn cluster_hashtag_same_slot() {
let cluster = cluster_or_skip!();
let client = match cluster.client().await {
Some(c) => c,
None => {
eprintln!("SKIP: cluster client connect failed");
return;
}
};
let _: () = client.set("{tag}:a", "1").await.unwrap();
let _: () = client.set("{tag}:b", "2").await.unwrap();
let got: Vec<Option<String>> = client.mget(&["{tag}:a", "{tag}:b"]).await.unwrap();
assert_eq!(got[0].as_deref(), Some("1"));
assert_eq!(got[1].as_deref(), Some("2"));
}
#[tokio::test]
async fn cluster_ping_resp2_and_resp3() {
let cluster = cluster_or_skip!();
for proto in [glide::ProtocolVersion::RESP2, glide::ProtocolVersion::RESP3] {
let client = match cluster.client_with_protocol(proto).await {
Some(c) => c,
None => {
eprintln!("SKIP: cluster client connect failed for {proto:?}");
return;
}
};
assert_eq!(client.ping().await.unwrap(), "PONG");
}
}