use crate::client::Client;
use redis::{Arg, Cmd};
use std::borrow::Borrow;
use telemetrylib::GlideSpan;
enum MaskingPattern {
ShowAll,
MaskAll,
ShowFirst(usize),
InterleavedKeyValue,
ShowFirstThenInterleave(usize),
}
fn masking_pattern(cmd_name: &str) -> MaskingPattern {
match cmd_name.to_ascii_uppercase().as_str() {
"AUTH" | "ECHO" | "HELLO" => MaskingPattern::MaskAll,
"APPEND" | "GETSET" | "LPUSH" | "LPUSHX" | "PFADD" | "PUBLISH" | "RPUSH" | "RPUSHX"
| "SADD" | "SET" | "SETNX" | "SPUBLISH" | "XADD" | "ZADD" => {
MaskingPattern::ShowFirst(1)
}
"SETEX" | "PSETEX" | "HSETNX" | "LSET" | "LPOS" => MaskingPattern::ShowFirst(2),
"LINSERT" => MaskingPattern::ShowFirst(3),
"MSET" | "MSETNX" => MaskingPattern::InterleavedKeyValue,
"HSET" | "HMSET" => MaskingPattern::ShowFirstThenInterleave(1),
"DECR" | "DECRBY" | "GET" | "GETBIT" | "GETDEL" | "GETEX" | "GETRANGE"
| "INCR" | "INCRBY" | "INCRBYFLOAT" | "MGET" | "STRLEN" | "SUBSTR"
| "BITCOUNT" | "BITFIELD" | "BITFIELD_RO" | "BITOP" | "BITPOS"
| "COPY" | "DEL" | "DUMP" | "EXISTS" | "EXPIRE" | "EXPIREAT" | "EXPIRETIME"
| "KEYS" | "MOVE" | "OBJECT" | "PERSIST" | "PEXPIRE" | "PEXPIREAT"
| "PEXPIRETIME" | "PTTL" | "RANDOMKEY" | "RENAME" | "RENAMENX" | "SCAN"
| "SORT" | "SORT_RO" | "TOUCH" | "TTL" | "TYPE" | "UNLINK" | "WAIT" | "WAITAOF"
| "WATCH" | "UNWATCH"
| "HDEL" | "HEXISTS" | "HGET" | "HGETALL" | "HINCRBY" | "HINCRBYFLOAT"
| "HKEYS" | "HLEN" | "HMGET" | "HRANDFIELD" | "HSCAN" | "HSTRLEN" | "HVALS"
| "LINDEX" | "LLEN" | "LMOVE" | "LMPOP" | "LPOP" | "LRANGE" | "LREM" | "LTRIM"
| "RPOP" | "RPOPLPUSH"
| "SCARD" | "SDIFF" | "SDIFFSTORE" | "SINTER" | "SINTERCARD" | "SINTERSTORE"
| "SISMEMBER" | "SMEMBERS" | "SMISMEMBER" | "SMOVE" | "SPOP" | "SRANDMEMBER"
| "SREM" | "SSCAN" | "SUNION" | "SUNIONSTORE"
| "ZCARD" | "ZCOUNT" | "ZDIFF" | "ZDIFFSTORE" | "ZINCRBY" | "ZINTER"
| "ZINTERCARD" | "ZINTERSTORE" | "ZLEXCOUNT" | "ZMPOP" | "ZMSCORE"
| "ZPOPMAX" | "ZPOPMIN" | "ZRANDMEMBER" | "ZRANGE" | "ZRANGEBYLEX"
| "ZRANGEBYSCORE" | "ZRANGESTORE" | "ZRANK" | "ZREM" | "ZREMRANGEBYLEX"
| "ZREMRANGEBYRANK" | "ZREMRANGEBYSCORE" | "ZREVRANGE" | "ZREVRANGEBYLEX"
| "ZREVRANGEBYSCORE" | "ZREVRANK" | "ZSCAN" | "ZSCORE" | "ZUNION" | "ZUNIONSTORE"
| "GEODIST" | "GEOHASH" | "GEOPOS" | "GEORADIUS" | "GEORADIUS_RO"
| "GEORADIUSBYMEMBER" | "GEORADIUSBYMEMBER_RO" | "GEOSEARCH" | "GEOSEARCHSTORE"
| "PFCOUNT" | "PFMERGE"
| "XACK" | "XCLAIM" | "XDEL" | "XGROUP" | "XINFO" | "XLEN" | "XPENDING"
| "XRANGE" | "XREAD" | "XREADGROUP" | "XREVRANGE" | "XTRIM"
| "COMMAND" | "DBSIZE" | "FLUSHALL" | "FLUSHDB" | "INFO" | "LOLWUT"
| "PING" | "RESET" | "SELECT" | "SLOWLOG" | "SWAPDB" | "TIME"
| "PUBSUB"
| "DISCARD" | "EXEC" | "MULTI"
| "SCRIPT" | "FUNCTION" => MaskingPattern::ShowAll,
_ => MaskingPattern::MaskAll,
}
}
pub fn serialize_query_text(cmd: &Cmd) -> Option<String> {
let mut args = cmd.args_iter().filter_map(|arg| match arg {
Arg::Simple(b) => Some(String::from_utf8_lossy(b).into_owned()),
_ => None,
});
let cmd_name = args.next()?;
let remaining: Vec<String> = args.collect();
if remaining.is_empty() {
return Some(cmd_name);
}
let pattern = masking_pattern(&cmd_name);
let mut parts = vec![cmd_name];
match pattern {
MaskingPattern::ShowAll => parts.extend(remaining),
MaskingPattern::MaskAll => {
parts.extend(remaining.iter().map(|_| "?".to_string()));
}
MaskingPattern::ShowFirst(n) => {
let n = n.min(remaining.len());
parts.extend_from_slice(&remaining[..n]);
parts.extend((0..remaining.len() - n).map(|_| "?".to_string()));
}
MaskingPattern::InterleavedKeyValue => {
for (i, arg) in remaining.iter().enumerate() {
if i % 2 == 0 {
parts.push(arg.clone());
} else {
parts.push("?".to_string());
}
}
}
MaskingPattern::ShowFirstThenInterleave(n) => {
let n = n.min(remaining.len());
parts.extend_from_slice(&remaining[..n]);
for (i, arg) in remaining[n..].iter().enumerate() {
if i % 2 == 0 {
parts.push(arg.clone()); } else {
parts.push("?".to_string()); }
}
}
}
Some(parts.join(" "))
}
const DB_SYSTEM_NAME: &str = "redis";
pub fn set_db_connection_attributes(span: &GlideSpan, client: &Client) {
span.set_attribute("db.system.name", DB_SYSTEM_NAME);
span.set_attribute("server.address", client.server_address().to_string());
span.set_attribute_i64("server.port", client.server_port() as i64);
span.set_attribute("db.namespace", client.db_namespace().to_string());
}
pub fn set_db_attributes(span: &GlideSpan, cmd: &Cmd, client: &Client) {
set_db_connection_attributes(span, client);
if let Some(Arg::Simple(name_bytes)) = cmd.args_iter().next() {
let cmd_name = String::from_utf8_lossy(name_bytes).into_owned();
span.set_attribute("db.operation.name", cmd_name);
}
if let Some(query_text) = serialize_query_text(cmd) {
span.set_attribute("db.query.text", query_text);
}
}
fn serialize_script_query_text(hash: &str, keys: &[&[u8]], args: &[&[u8]]) -> String {
let mut parts: Vec<String> = Vec::with_capacity(3 + keys.len() + args.len());
parts.push("EVALSHA".to_string());
parts.push(hash.to_string());
parts.push(keys.len().to_string());
for key in keys {
parts.push(String::from_utf8_lossy(key).into_owned());
}
for _ in args {
parts.push("?".to_string());
}
parts.join(" ")
}
pub fn set_db_script_attributes(
span: &GlideSpan,
hash: &str,
keys: &[&[u8]],
args: &[&[u8]],
client: &Client,
) {
set_db_connection_attributes(span, client);
span.set_attribute("db.operation.name", "EVALSHA");
span.set_attribute(
"db.query.text",
serialize_script_query_text(hash, keys, args),
);
}
pub fn set_db_batch_attributes<T: Borrow<Cmd>>(span: &GlideSpan, cmds: &[T], client: &Client) {
set_db_connection_attributes(span, client);
let mut query_texts: Vec<String> = Vec::with_capacity(cmds.len());
let mut cmd_names: Vec<String> = Vec::with_capacity(cmds.len());
for cmd in cmds {
let cmd: &Cmd = cmd.borrow();
if let Some(text) = serialize_query_text(cmd) {
query_texts.push(text);
}
if let Some(Arg::Simple(name_bytes)) = cmd.args_iter().next() {
cmd_names.push(String::from_utf8_lossy(name_bytes).into_owned());
}
}
if !query_texts.is_empty() {
span.set_attribute("db.query.text", query_texts.join("\n"));
}
let op_name = if !cmd_names.is_empty() && cmd_names.iter().all(|n| n == &cmd_names[0]) {
format!("PIPELINE {}", cmd_names[0])
} else {
"PIPELINE".to_string()
};
span.set_attribute("db.operation.name", op_name);
}
#[cfg(test)]
mod tests {
use super::*;
fn make_cmd(name: &str, args: &[&str]) -> Cmd {
let mut cmd = redis::cmd(name);
for arg in args {
cmd.arg(*arg);
}
cmd
}
const QUERY_TEXT_CASES: &[(&str, &[&str], &str)] = &[
("PING", &[], "PING"),
("AUTH", &["s3cret!"], "AUTH ?"),
("AUTH", &["admin", "s3cret!"], "AUTH ? ?"),
("ECHO", &["sensitive-payload"], "ECHO ?"),
(
"SET",
&["user:1001", r#"{"name":"Alice"}"#],
"SET user:1001 ?",
),
(
"SET",
&["session:abc", "token123", "EX", "3600", "NX"],
"SET session:abc ? ? ? ?",
),
(
"SETNX",
&["lock:resource", "owner-id"],
"SETNX lock:resource ?",
),
(
"SETEX",
&["cache:page:home", "300", "<html>..."],
"SETEX cache:page:home 300 ?",
),
(
"PSETEX",
&["ratelimit:user:42", "1500", "1"],
"PSETEX ratelimit:user:42 1500 ?",
),
(
"APPEND",
&["audit:log", "2026-02-24 action=login"],
"APPEND audit:log ?",
),
(
"GETSET",
&["counter:visits", "0"],
"GETSET counter:visits ?",
),
(
"LPUSH",
&["queue:jobs", "job1", "job2", "job3"],
"LPUSH queue:jobs ? ? ?",
),
(
"LPUSHX",
&["queue:existing", "new-item"],
"LPUSHX queue:existing ?",
),
(
"RPUSH",
&["notifications:user:5", r#"{"type":"mention"}"#],
"RPUSH notifications:user:5 ?",
),
(
"RPUSHX",
&["pending:emails", "msg-body"],
"RPUSHX pending:emails ?",
),
(
"SADD",
&["tags:article:99", "rust", "async", "valkey"],
"SADD tags:article:99 ? ? ?",
),
(
"PFADD",
&["unique:visitors", "user1", "user2", "user3"],
"PFADD unique:visitors ? ? ?",
),
(
"PUBLISH",
&["chat:room:1", "Hello everyone!"],
"PUBLISH chat:room:1 ?",
),
(
"SPUBLISH",
&["orders:shard1", r#"{"orderId":42}"#],
"SPUBLISH orders:shard1 ?",
),
(
"ZADD",
&["leaderboard", "9500", "player:42"],
"ZADD leaderboard ? ?",
),
(
"XADD",
&["stream:events", "*", "action", "login", "user", "alice"],
"XADD stream:events ? ? ? ? ?",
),
(
"HSET",
&["user:1001", "email", "alice@example.com"],
"HSET user:1001 email ?",
),
(
"HSET",
&["user:1001", "email", "a@b.com", "name", "Alice"],
"HSET user:1001 email ? name ?",
),
(
"HMSET",
&["product:500", "price", "29.99", "stock", "150"],
"HMSET product:500 price ? stock ?",
),
(
"HSETNX",
&["user:1001", "created_at", "2026-01-01"],
"HSETNX user:1001 created_at ?",
),
(
"LSET",
&["queue:jobs", "0", "updated-payload"],
"LSET queue:jobs 0 ?",
),
(
"LINSERT",
&["playlist", "BEFORE", "track:5", "track:new"],
"LINSERT playlist BEFORE track:5 ?",
),
(
"LPOS",
&["queue:jobs", "target-item", "RANK", "1", "COUNT", "2"],
"LPOS queue:jobs target-item ? ? ? ?",
),
("MSET", &["config:retries", "5"], "MSET config:retries ?"),
(
"MSET",
&[
"user:1:name",
"Alice",
"user:2:name",
"Bob",
"user:3:name",
"Carol",
],
"MSET user:1:name ? user:2:name ? user:3:name ?",
),
(
"MSETNX",
&["lock:a", "owner1", "lock:b", "owner2"],
"MSETNX lock:a ? lock:b ?",
),
("GET", &["session:abc"], "GET session:abc"),
(
"DEL",
&["temp:1", "temp:2", "temp:3"],
"DEL temp:1 temp:2 temp:3",
),
("HGET", &["user:1001", "email"], "HGET user:1001 email"),
("KEYS", &["user:*"], "KEYS user:*"),
(
"SINTERCARD",
&["2", "set:a", "set:b", "LIMIT", "10"],
"SINTERCARD 2 set:a set:b LIMIT 10",
),
(
"SCAN",
&["0", "MATCH", "user:*", "COUNT", "100"],
"SCAN 0 MATCH user:* COUNT 100",
),
("CUSTOMCMD", &["arg1", "arg2"], "CUSTOMCMD ? ?"),
("CONFIG", &["SET", "requirepass", "s3cret"], "CONFIG ? ? ?"),
("ACL", &["SETUSER", "admin", ">password"], "ACL ? ? ?"),
(
"MIGRATE",
&["host", "6379", "key", "0", "5000", "AUTH", "s3cret"],
"MIGRATE ? ? ? ? ? ? ?",
),
("HELLO", &["3", "AUTH", "user", "pass"], "HELLO ? ? ? ?"),
];
#[test]
fn test_serialize_query_text() {
for (cmd_name, args, expected) in QUERY_TEXT_CASES {
let cmd = make_cmd(cmd_name, args);
let result = serialize_query_text(&cmd).unwrap();
assert_eq!(&result, expected, "query: {cmd_name} {}", args.join(" "));
}
}
#[test]
fn test_masking_pattern_case_insensitive() {
assert!(matches!(
masking_pattern("set"),
MaskingPattern::ShowFirst(1)
));
}
#[test]
fn test_batch_query_text_multiple_commands() {
let cmds = vec![
make_cmd("SET", &["key1", "val1"]),
make_cmd("GET", &["key2"]),
make_cmd("HSET", &["hash1", "field", "value"]),
];
let mut query_texts: Vec<String> = Vec::new();
let mut cmd_names: Vec<String> = Vec::new();
for cmd in &cmds {
if let Some(text) = serialize_query_text(cmd) {
query_texts.push(text);
}
if let Some(Arg::Simple(name_bytes)) = cmd.args_iter().next() {
cmd_names.push(String::from_utf8_lossy(name_bytes).into_owned());
}
}
let joined = query_texts.join("\n");
assert_eq!(joined, "SET key1 ?\nGET key2\nHSET hash1 field ?");
assert!(!cmd_names.iter().all(|n| n == &cmd_names[0]));
}
#[test]
fn test_batch_query_text_same_commands() {
let cmds = vec![
make_cmd("GET", &["key1"]),
make_cmd("GET", &["key2"]),
make_cmd("GET", &["key3"]),
];
let mut query_texts: Vec<String> = Vec::new();
let mut cmd_names: Vec<String> = Vec::new();
for cmd in &cmds {
if let Some(text) = serialize_query_text(cmd) {
query_texts.push(text);
}
if let Some(Arg::Simple(name_bytes)) = cmd.args_iter().next() {
cmd_names.push(String::from_utf8_lossy(name_bytes).into_owned());
}
}
let joined = query_texts.join("\n");
assert_eq!(joined, "GET key1\nGET key2\nGET key3");
assert!(cmd_names.iter().all(|n| n == &cmd_names[0]));
let op_name = format!("PIPELINE {}", cmd_names[0]);
assert_eq!(op_name, "PIPELINE GET");
}
#[test]
fn test_batch_query_text_with_masking() {
let cmds = vec![
make_cmd("AUTH", &["password123"]), make_cmd("SET", &["key", "secret-val"]), make_cmd("MSET", &["k1", "v1", "k2", "v2"]), ];
let mut query_texts: Vec<String> = Vec::new();
for cmd in &cmds {
if let Some(text) = serialize_query_text(cmd) {
query_texts.push(text);
}
}
let joined = query_texts.join("\n");
assert_eq!(joined, "AUTH ?\nSET key ?\nMSET k1 ? k2 ?");
}
#[test]
fn test_serialize_script_query_text_with_keys_and_args() {
let result = serialize_script_query_text(
"abc123def456",
&[b"user:1", b"user:2"],
&[b"secret-val1", b"secret-val2"],
);
assert_eq!(result, "EVALSHA abc123def456 2 user:1 user:2 ? ?");
}
#[test]
fn test_serialize_script_query_text_keys_only() {
let result = serialize_script_query_text("sha1hash", &[b"mykey"], &[]);
assert_eq!(result, "EVALSHA sha1hash 1 mykey");
}
#[test]
fn test_serialize_script_query_text_no_keys_no_args() {
let result = serialize_script_query_text("emptyscript", &[], &[]);
assert_eq!(result, "EVALSHA emptyscript 0");
}
#[test]
fn test_serialize_script_query_text_args_only() {
let result = serialize_script_query_text("argsonly", &[], &[b"arg1", b"arg2", b"arg3"]);
assert_eq!(result, "EVALSHA argsonly 0 ? ? ?");
}
}