use std::collections::HashMap;
use std::io::{BufReader, Write};
use std::net::{TcpStream, ToSocketAddrs};
use std::time::Duration;
use crate::client::ClientOptions;
use crate::codec::{encode_strings, read_value};
use crate::error::Error;
use crate::value::RespValue;
pub const SLOT_COUNT: u16 = 16_384;
const MAX_REDIRECTS: usize = 3;
const KEYLESS: &[&str] = &[
"ACL",
"AUTH",
"BGSAVE",
"CLIENT",
"CLUSTER",
"COMMAND",
"CONFIG",
"DBSIZE",
"DISCARD",
"ECHO",
"EXEC",
"FLUSHALL",
"FLUSHDB",
"FUNCTION",
"HELLO",
"INFO",
"KEYS",
"LASTSAVE",
"MULTI",
"PING",
"PSUBSCRIBE",
"PUBLISH",
"PUNSUBSCRIBE",
"QUIT",
"RANDOMKEY",
"READONLY",
"READWRITE",
"RESET",
"SAVE",
"SCAN",
"SCRIPT",
"SELECT",
"SLOWLOG",
"SUBSCRIBE",
"SWAPDB",
"TIME",
"UNSUBSCRIBE",
"UNWATCH",
"WAIT",
];
const EVERY_ARGUMENT: &[&str] = &[
"DEL",
"EXISTS",
"MGET",
"PFCOUNT",
"PFMERGE",
"SDIFF",
"SDIFFSTORE",
"SINTER",
"SINTERSTORE",
"SSUBSCRIBE",
"SUNION",
"SUNIONSTORE",
"SUNSUBSCRIBE",
"TOUCH",
"UNLINK",
"WATCH",
];
const TWO_KEYS: &[&str] = &[
"BLMOVE",
"BRPOPLPUSH",
"COPY",
"LMOVE",
"RENAME",
"RENAMENX",
"RPOPLPUSH",
"SMOVE",
];
const SPLIT: &[&str] = &["DEL", "EXISTS", "MGET", "MSET", "TOUCH", "UNLINK"];
const TRANSACTION: &[&str] = &["DISCARD", "EXEC", "MULTI", "UNWATCH", "WATCH"];
pub fn hash_slot(key: &str) -> u16 {
crc16(hash_tag(key.as_bytes())) % SLOT_COUNT
}
fn hash_tag(key: &[u8]) -> &[u8] {
let Some(open) = key.iter().position(|byte| *byte == b'{') else {
return key;
};
let tagged = &key[open + 1..];
let Some(close) = tagged.iter().position(|byte| *byte == b'}') else {
return key;
};
if close == 0 { key } else { &tagged[..close] }
}
fn crc16(bytes: &[u8]) -> u16 {
let mut crc = 0u16;
for byte in bytes {
crc ^= u16::from(*byte) << 8;
for _ in 0..8 {
crc = if crc & 0x8000 != 0 {
(crc << 1) ^ 0x1021
} else {
crc << 1
};
}
}
crc
}
fn command_name(arguments: &[&str]) -> String {
arguments
.first()
.map(|name| name.to_ascii_uppercase())
.unwrap_or_default()
}
fn listed(list: &[&str], name: &str) -> bool {
list.contains(&name)
}
pub(crate) fn command_keys<'a>(name: &str, arguments: &[&'a str]) -> Vec<&'a str> {
if arguments.len() < 2 || listed(KEYLESS, name) {
return Vec::new();
}
if listed(EVERY_ARGUMENT, name) {
return arguments[1..].to_vec();
}
if listed(TWO_KEYS, name) {
return arguments[1..].iter().take(2).copied().collect();
}
match name {
"MSET" | "MSETNX" => arguments[1..].iter().step_by(2).copied().collect(),
"BLPOP" | "BRPOP" | "BZPOPMIN" | "BZPOPMAX" => {
arguments[1..arguments.len().saturating_sub(1)].to_vec()
}
"EVAL" | "EVALSHA" | "EVAL_RO" | "EVALSHA_RO" | "FCALL" | "FCALL_RO" => {
counted(arguments, 2)
}
"ZUNION" | "ZINTER" | "ZDIFF" => counted(arguments, 1),
"ZUNIONSTORE" | "ZINTERSTORE" | "ZDIFFSTORE" => {
let mut keys = vec![arguments[1]];
keys.extend(counted(arguments, 2));
keys
}
"XREAD" | "XREADGROUP" => streams(arguments),
"OBJECT" | "MEMORY" if arguments.len() > 2 => vec![arguments[2]],
_ => vec![arguments[1]],
}
}
fn counted<'a>(arguments: &[&'a str], count_index: usize) -> Vec<&'a str> {
let Some(count) = arguments
.get(count_index)
.and_then(|value| value.parse().ok())
else {
return Vec::new();
};
arguments
.get(count_index + 1..)
.unwrap_or(&[])
.iter()
.take(count)
.copied()
.collect()
}
fn streams<'a>(arguments: &[&'a str]) -> Vec<&'a str> {
let Some(index) = arguments
.iter()
.position(|argument| argument.eq_ignore_ascii_case("STREAMS"))
else {
return Vec::new();
};
let rest = arguments.len().saturating_sub(index + 1);
arguments[index + 1..index + 1 + rest / 2].to_vec()
}
pub(crate) struct Connection {
stream: BufReader<TcpStream>,
host: String,
port: u16,
}
impl Connection {
pub(crate) fn open(host: &str, port: u16, options: &ClientOptions) -> Result<Self, Error> {
let mut last_error = None;
for address in (host, port).to_socket_addrs()? {
match TcpStream::connect_timeout(&address, options.connect_timeout) {
Ok(stream) => {
stream.set_nodelay(true)?;
let mut connection = Self {
stream: BufReader::new(stream),
host: host.to_owned(),
port,
};
if let Some(password) = options.password.clone() {
let reply = if let Some(username) = options.username.clone() {
connection.send(&["AUTH", &username, &password])?
} else {
connection.send(&["AUTH", &password])?
};
throw_error(reply)?;
}
return Ok(connection);
}
Err(error) => last_error = Some(error),
}
}
Err(Error::Io(last_error.unwrap_or_else(|| {
std::io::Error::new(std::io::ErrorKind::NotFound, "no addresses for host")
})))
}
pub(crate) fn send(&mut self, arguments: &[&str]) -> Result<RespValue, Error> {
let mut values = self.send_many(&[arguments])?;
values
.pop()
.ok_or_else(|| Error::Protocol("missing reply".to_owned()))
}
pub(crate) fn send_owned(&mut self, arguments: &[String]) -> Result<RespValue, Error> {
let borrowed: Vec<&str> = arguments.iter().map(String::as_str).collect();
self.send(&borrowed)
}
pub(crate) fn send_replies(
&mut self,
arguments: &[&str],
replies: usize,
) -> Result<Vec<RespValue>, Error> {
let payload = encode_strings(arguments)?;
self.stream.get_mut().write_all(&payload)?;
self.stream.get_mut().flush()?;
let mut values = Vec::with_capacity(replies);
for _ in 0..replies {
values.push(read_value(&mut self.stream)?);
}
Ok(values)
}
pub(crate) fn send_many(&mut self, commands: &[&[&str]]) -> Result<Vec<RespValue>, Error> {
let mut payload = Vec::new();
for command in commands {
payload.extend(encode_strings(command)?);
}
self.stream.get_mut().write_all(&payload)?;
self.stream.get_mut().flush()?;
let mut values = Vec::with_capacity(commands.len());
for _ in 0..commands.len() {
values.push(read_value(&mut self.stream)?);
}
Ok(values)
}
pub(crate) fn send_many_owned(
&mut self,
commands: &[Vec<String>],
) -> Result<Vec<RespValue>, Error> {
let borrowed: Vec<Vec<&str>> = commands
.iter()
.map(|command| command.iter().map(String::as_str).collect())
.collect();
let refs: Vec<&[&str]> = borrowed.iter().map(Vec::as_slice).collect();
self.send_many(&refs)
}
pub(crate) fn read(&mut self) -> Result<RespValue, Error> {
read_value(&mut self.stream)
}
pub(crate) fn set_read_timeout(&mut self, timeout: Option<Duration>) -> Result<(), Error> {
self.stream.get_ref().set_read_timeout(timeout)?;
Ok(())
}
}
pub(crate) fn throw_error(reply: RespValue) -> Result<RespValue, Error> {
match reply {
RespValue::Error(message) => Err(Error::Server(message)),
value => Ok(value),
}
}
pub(crate) enum Discovery {
Standalone(Connection),
Cluster(Box<ClusterRouter>),
}
#[derive(Clone)]
struct Topology {
owners: Vec<Option<usize>>,
nodes: Vec<(String, u16)>,
}
impl Topology {
fn parse(reply: &RespValue) -> Option<Self> {
let RespValue::Array(ranges) = reply else {
return None;
};
if ranges.is_empty() {
return None;
}
let mut owners = vec![None; SLOT_COUNT as usize];
let mut nodes = Vec::new();
let mut ordered = ranges.clone();
ordered.sort_by_key(|range| match range {
RespValue::Array(fields) => fields.first().and_then(RespValue::as_integer).unwrap_or(0),
_ => 0,
});
for range in &ordered {
let RespValue::Array(fields) = range else {
return None;
};
if fields.len() < 3 {
return None;
}
let RespValue::Array(primary) = &fields[2] else {
return None;
};
if primary.len() < 2 {
return None;
}
let host = primary[0].as_string().ok().flatten()?;
let port = u16::try_from(primary[1].as_integer()?).ok()?;
let node = match nodes
.iter()
.position(|(existing, existing_port)| existing == &host && *existing_port == port)
{
Some(index) => index,
None => {
nodes.push((host, port));
nodes.len() - 1
}
};
let first = fields[0].as_integer()?.clamp(0, i64::from(SLOT_COUNT) - 1) as usize;
let last = fields[1].as_integer()?.clamp(0, i64::from(SLOT_COUNT) - 1) as usize;
owners[first..=last].fill(Some(node));
}
Some(Self { owners, nodes })
}
fn owner(&self, slot: u16) -> Option<&(String, u16)> {
self.owners
.get(usize::from(slot))
.and_then(|node| node.and_then(|index| self.nodes.get(index)))
}
}
pub(crate) struct ClusterRouter {
options: ClientOptions,
seed: String,
connections: HashMap<String, Connection>,
topology: Topology,
queued: Vec<Vec<String>>,
pinned: Option<String>,
multi_pending: bool,
in_multi: bool,
subscriber: Option<String>,
}
impl ClusterRouter {
pub(crate) fn discover(
mut seed: Connection,
options: ClientOptions,
) -> Result<Discovery, Error> {
let reply = seed.send(&["CLUSTER", "SLOTS"])?;
let Some(topology) = Topology::parse(&reply) else {
return Ok(Discovery::Standalone(seed));
};
let key = endpoint(seed.host.as_str(), seed.port);
Ok(Discovery::Cluster(Box::new(Self {
options,
seed: key.clone(),
connections: HashMap::from([(key, seed)]),
topology,
queued: Vec::new(),
pinned: None,
multi_pending: false,
in_multi: false,
subscriber: None,
})))
}
pub(crate) fn node_count(&self) -> usize {
self.topology.nodes.len()
}
pub(crate) fn execute(&mut self, arguments: &[&str]) -> Result<RespValue, Error> {
let name = command_name(arguments);
if self.pinned.is_some() || self.multi_pending {
return self.transaction(&name, arguments);
}
match name.as_str() {
"MULTI" => {
self.multi_pending = true;
return Ok(RespValue::Simple("OK".to_owned()));
}
"WATCH" => return self.watch(arguments),
"DBSIZE" | "FLUSHDB" | "FLUSHALL" => return self.every_node(&name, arguments),
"SCRIPT" | "FUNCTION" if spreads_to_every_node(arguments) => {
return self.every_node(&name, arguments);
}
"SCAN" => return self.scan(arguments),
"PUBLISH" => {
let key = self.for_slot(0)?;
return self.send_routed(&key, arguments);
}
_ => {}
}
let keys = command_keys(&name, arguments);
if listed(SPLIT, &name) && spans_slots(&keys) {
return self.split(&name, arguments, &keys);
}
let target = self.target(&keys)?;
self.send_routed(&target, arguments)
}
pub(crate) fn execute_many(&mut self, commands: &[&[&str]]) -> Result<Vec<RespValue>, Error> {
if self.pinned.is_some()
|| self.multi_pending
|| commands
.iter()
.any(|command| listed(TRANSACTION, &command_name(command)))
{
let key = self.first_keyed_target(commands)?;
let owned: Vec<Vec<String>> = commands
.iter()
.map(|command| command.iter().map(|item| (*item).to_owned()).collect())
.collect();
return self.connection(&key)?.send_many_owned(&owned);
}
let mut groups: HashMap<String, Vec<usize>> = HashMap::new();
for (index, command) in commands.iter().enumerate() {
let name = command_name(command);
let keys = command_keys(&name, command);
let target = self.target(&keys)?;
groups.entry(target).or_default().push(index);
}
let mut replies = vec![RespValue::Null; commands.len()];
for (target, positions) in groups {
let batch: Vec<Vec<String>> = positions
.iter()
.map(|index| {
commands[*index]
.iter()
.map(|item| (*item).to_owned())
.collect()
})
.collect();
let values = self.connection(&target)?.send_many_owned(&batch)?;
for (offset, value) in values.into_iter().enumerate() {
replies[positions[offset]] = value;
}
}
for (index, reply) in replies.iter_mut().enumerate() {
if let Some((host, port)) = moved(reply) {
self.refresh()?;
let key = self.ensure(&host, port)?;
*reply = self.send_routed(&key, commands[index])?;
}
}
Ok(replies)
}
pub(crate) fn run_replies(
&mut self,
arguments: &[&str],
replies: usize,
) -> Result<Vec<RespValue>, Error> {
let name = command_name(arguments);
let slot = if matches!(name.as_str(), "SSUBSCRIBE" | "SUNSUBSCRIBE") && arguments.len() > 1
{
hash_slot(arguments[1])
} else {
0
};
let key = self.for_slot(slot)?;
self.subscriber = Some(key.clone());
let values = self.connection(&key)?.send_replies(arguments, replies)?;
if values.len() < replies {
return Err(Error::Protocol("missing subscribe confirmation".to_owned()));
}
Ok(values)
}
pub(crate) fn read_message(&mut self) -> Result<RespValue, Error> {
let key = self.subscriber.clone().unwrap_or_else(|| self.seed.clone());
self.connection(&key)?.read()
}
pub(crate) fn set_read_timeout(&mut self, timeout: Option<Duration>) -> Result<(), Error> {
for connection in self.connections.values_mut() {
connection.set_read_timeout(timeout)?;
}
Ok(())
}
fn transaction(&mut self, name: &str, arguments: &[&str]) -> Result<RespValue, Error> {
if let Some(pinned) = self.pinned.clone() {
let reply = self.connection(&pinned)?.send(arguments)?;
if name == "EXEC" || name == "DISCARD" || (name == "UNWATCH" && !self.in_multi) {
self.pinned = None;
self.in_multi = false;
} else if name == "MULTI" && !matches!(reply, RespValue::Error(_)) {
self.in_multi = true;
}
return Ok(reply);
}
match name {
"MULTI" => {
return Ok(RespValue::Error(
"ERR MULTI calls can not be nested".to_owned(),
));
}
"WATCH" => {
return Ok(RespValue::Error(
"ERR WATCH inside MULTI is not allowed".to_owned(),
));
}
"DISCARD" => {
self.reset_transaction();
return Ok(RespValue::Simple("OK".to_owned()));
}
"EXEC" => {
let key = self.seed.clone();
let mut commands = vec![vec!["MULTI".to_owned()]];
commands.extend(self.queued.iter().cloned());
commands.push(arguments.iter().map(|item| (*item).to_owned()).collect());
let replies = self.connection(&key)?.send_many_owned(&commands)?;
self.reset_transaction();
return Ok(replies.into_iter().next_back().unwrap_or(RespValue::Null));
}
_ => {}
}
let keys = command_keys(name, arguments);
if keys.is_empty() {
self.queued
.push(arguments.iter().map(|item| (*item).to_owned()).collect());
return Ok(RespValue::Simple("QUEUED".to_owned()));
}
let owner = self.for_slot(hash_slot(keys[0]))?;
let mut commands = vec![vec!["MULTI".to_owned()]];
commands.extend(self.queued.iter().cloned());
commands.push(arguments.iter().map(|item| (*item).to_owned()).collect());
let replies = self.connection(&owner)?.send_many_owned(&commands)?;
self.queued.clear();
self.multi_pending = false;
if matches!(replies.first(), Some(RespValue::Error(_))) {
return Ok(replies.into_iter().next().unwrap());
}
self.pinned = Some(owner);
self.in_multi = true;
Ok(replies.into_iter().next_back().unwrap_or(RespValue::Null))
}
fn reset_transaction(&mut self) {
self.queued.clear();
self.multi_pending = false;
self.in_multi = false;
self.pinned = None;
}
fn watch(&mut self, arguments: &[&str]) -> Result<RespValue, Error> {
let keys = command_keys("WATCH", arguments);
let owner = self.target(&keys)?;
let reply = self.connection(&owner)?.send(arguments)?;
if !matches!(reply, RespValue::Error(_)) {
self.pinned = Some(owner);
}
Ok(reply)
}
fn every_node(&mut self, name: &str, arguments: &[&str]) -> Result<RespValue, Error> {
let nodes: Vec<(String, u16)> = self.topology.nodes.clone();
let mut replies = Vec::new();
for (host, port) in nodes {
let key = self.ensure(&host, port)?;
let reply = self.connection(&key)?.send(arguments)?;
if let RespValue::Error(_) = reply {
return Ok(reply);
}
replies.push(reply);
}
if name == "DBSIZE" {
return Ok(RespValue::Integer(
replies.iter().filter_map(RespValue::as_integer).sum(),
));
}
Ok(replies.into_iter().next().unwrap_or(RespValue::Null))
}
fn scan(&mut self, arguments: &[&str]) -> Result<RespValue, Error> {
let nodes = self.topology.nodes.len() as u64;
let cursor = arguments
.get(1)
.and_then(|value| value.parse::<u64>().ok())
.ok_or_else(|| Error::Protocol("ERR invalid cursor".to_owned()))?;
let node = (cursor % nodes) as usize;
let mut forwarded: Vec<String> = arguments.iter().map(|item| (*item).to_owned()).collect();
forwarded[1] = (cursor / nodes).to_string();
let (host, port) = self.topology.nodes[node].clone();
let key = self.ensure(&host, port)?;
let reply = self.connection(&key)?.send_owned(&forwarded)?;
let RespValue::Array(items) = &reply else {
return Ok(reply);
};
if items.len() < 2 {
return Ok(reply);
}
let server_next = items[0]
.as_string()?
.unwrap_or_else(|| "0".to_owned())
.parse::<u64>()
.unwrap_or(0);
let client_next = if server_next != 0 {
server_next
.checked_mul(nodes)
.and_then(|value| value.checked_add(node as u64))
.ok_or_else(|| Error::Protocol("scan cursor overflow".to_owned()))?
} else if (node as u64) + 1 < nodes {
node as u64 + 1
} else {
0
};
Ok(RespValue::Array(vec![
RespValue::Bulk(client_next.to_string().into_bytes()),
items[1].clone(),
]))
}
fn split(&mut self, name: &str, arguments: &[&str], keys: &[&str]) -> Result<RespValue, Error> {
let pairs = name == "MSET";
let mut by_slot: HashMap<u16, Vec<usize>> = HashMap::new();
for (index, key) in keys.iter().enumerate() {
by_slot.entry(hash_slot(key)).or_default().push(index);
}
let mut parts = Vec::new();
for (slot, positions) in &by_slot {
let mut command = vec![arguments[0].to_owned()];
for position in positions {
if pairs {
command.push(arguments[1 + position * 2].to_owned());
command.push(arguments[2 + position * 2].to_owned());
} else {
command.push(keys[*position].to_owned());
}
}
parts.push((*slot, positions.clone(), command));
}
let mut replies = Vec::new();
for (slot, _, command) in &parts {
let key = self.for_slot(*slot)?;
let reply = self.connection(&key)?.send_owned(command)?;
if let Some((host, port)) = moved(&reply) {
self.refresh()?;
let redirected = self.ensure(&host, port)?;
replies.push(self.send_routed(&redirected, &borrowed(command))?);
} else if matches!(reply, RespValue::Error(_)) {
return Ok(reply);
} else {
replies.push(reply);
}
}
match name {
"MSET" => Ok(RespValue::Simple("OK".to_owned())),
"MGET" => {
let mut values = vec![RespValue::Null; keys.len()];
for (index, (_, positions, _)) in parts.iter().enumerate() {
let Some(items) = replies[index].as_array() else {
continue;
};
for (offset, item) in items.iter().enumerate() {
values[positions[offset]] = item.clone();
}
}
Ok(RespValue::Array(values))
}
_ => Ok(RespValue::Integer(
replies.iter().filter_map(RespValue::as_integer).sum(),
)),
}
}
fn send_routed(&mut self, key: &str, arguments: &[&str]) -> Result<RespValue, Error> {
let mut current = key.to_owned();
for redirect in 0..=MAX_REDIRECTS {
let reply = self.connection(¤t)?.send(arguments)?;
if redirect == MAX_REDIRECTS {
return Ok(reply);
}
let Some((host, port)) = moved(&reply) else {
return Ok(reply);
};
self.refresh()?;
current = self.ensure(&host, port)?;
}
unreachable!("redirect loop ends")
}
fn refresh(&mut self) -> Result<(), Error> {
let keys: Vec<String> = self.connections.keys().cloned().collect();
for key in keys {
let reply = self.connection(&key)?.send(&["CLUSTER", "SLOTS"])?;
if let Some(topology) = Topology::parse(&reply) {
self.topology = topology;
return Ok(());
}
}
Ok(())
}
fn target(&mut self, keys: &[&str]) -> Result<String, Error> {
if let Some(key) = keys.first() {
return self.for_slot(hash_slot(key));
}
Ok(self.seed.clone())
}
fn first_keyed_target(&mut self, commands: &[&[&str]]) -> Result<String, Error> {
for command in commands {
let keys = command_keys(&command_name(command), command);
if let Some(key) = keys.first() {
return self.for_slot(hash_slot(key));
}
}
Ok(self.seed.clone())
}
fn for_slot(&mut self, slot: u16) -> Result<String, Error> {
if let Some((host, port)) = self.topology.owner(slot).cloned() {
return self.ensure(&host, port);
}
Ok(self.seed.clone())
}
fn ensure(&mut self, host: &str, port: u16) -> Result<String, Error> {
let key = endpoint(host, port);
if !self.connections.contains_key(&key) {
let opened = Connection::open(host, port, &self.options)?;
self.connections.insert(key.clone(), opened);
}
Ok(key)
}
fn connection(&mut self, key: &str) -> Result<&mut Connection, Error> {
self.connections
.get_mut(key)
.ok_or_else(|| Error::Protocol(format!("no connection for {key}")))
}
}
fn borrowed(command: &[String]) -> Vec<&str> {
command.iter().map(String::as_str).collect()
}
fn spreads_to_every_node(arguments: &[&str]) -> bool {
arguments.get(1).is_some_and(|subcommand| {
matches!(
subcommand.to_ascii_uppercase().as_str(),
"LOAD" | "FLUSH" | "DELETE"
)
})
}
fn spans_slots(keys: &[&str]) -> bool {
let Some(first) = keys.first() else {
return false;
};
let slot = hash_slot(first);
keys.iter().any(|key| hash_slot(key) != slot)
}
fn moved(reply: &RespValue) -> Option<(String, u16)> {
let RespValue::Error(text) = reply else {
return None;
};
let rest = text.strip_prefix("MOVED ")?;
let address = rest.split_whitespace().nth(1)?;
let (host, port) = address.rsplit_once(':')?;
Some((host.to_owned(), port.parse().ok()?))
}
fn endpoint(host: &str, port: u16) -> String {
format!("{host}:{port}")
}
#[cfg(test)]
mod tests {
use super::{command_keys, hash_slot};
#[test]
fn hash_slot_matches_redis_and_honors_a_hash_tag() {
assert_eq!(hash_slot("123456789"), 12_739);
assert_eq!(hash_slot("somekey"), 11_058);
assert_eq!(hash_slot("{user1000}.following"), 3_443);
assert_eq!(hash_slot("user:{42}:name"), hash_slot("cart:{42}"));
}
#[test]
fn command_keys_follow_the_redis_positions() {
assert_eq!(command_keys("MGET", &["MGET", "a", "b"]), ["a", "b"]);
assert_eq!(
command_keys("MSET", &["MSET", "a", "1", "c", "2"]),
["a", "c"]
);
assert_eq!(
command_keys("EVAL", &["EVAL", "return 1", "2", "k1", "k2", "arg"]),
["k1", "k2"]
);
assert!(command_keys("PING", &["PING"]).is_empty());
assert!(command_keys("PUBLISH", &["PUBLISH", "chan", "hi"]).is_empty());
assert_eq!(
command_keys("SPUBLISH", &["SPUBLISH", "chan", "hi"]),
["chan"]
);
}
}