ruvio-client 0.2.9

RESP2 client for the Ruvio key-value server
Documentation

ruvio-client

Blocking RESP2 client for Ruvio. Connecting does not send CLIENT SETINFO or HELLO. Give the client one address; it sends CLUSTER SLOTS and, on a sharded server, opens one socket per shard and routes by hash slot. MOVED refreshes the map. MGET / MSET / DEL split across shards. Methods take &mut self and are named after the command they send. They cover every command Ruvio supports. execute and execute_many send any argument list, as redis-cli would.

use std::time::Duration;
use ruvio_client::Client;

let mut db = Client::connect("127.0.0.1", 6379)?;
db.set("session:ada", "hello")?;
let value = db.get_string("session:ada")?;
db.incr("visits")?;
db.expire("session:ada", Duration::from_secs(60))?;
db.hset("user:1", &[("name", "Ada"), ("plan", "pro")])?;
db.zadd("leaderboard", &[("player1", 100.0), ("player2", 80.0)])?;
let top = db.zrevrange_with_scores("leaderboard", 0, 9)?;

Connecting

hash_slot(key) is the same 0–16383 slot Ruvio uses, including {hash tags}. Set ClientOptions::discover_cluster to false to stay on the seed socket.

A standalone client is bound to one logical database. Database 0 sends nothing after discovery. Any other index sends SELECT before the connect call returns. Cluster mode has only database 0.

use std::net::{IpAddr, Ipv4Addr, SocketAddr};
use ruvio_client::{Client, ClientOptions};

let mut cache = Client::connect_database("127.0.0.1", 6379, 2)?;
let mut sessions = Client::connect_ip(IpAddr::V4(Ipv4Addr::LOCALHOST), 6379, 3)?;
let mut main = Client::connect_addr(SocketAddr::from(([10, 0, 0, 5], 6379)), 0)?;
let mut secured = Client::connect_with(ClientOptions {
    host: "10.0.0.5".into(),
    username: Some("app".into()),
    password: Some("secret".into()),
    database: 4,
    ..ClientOptions::default()
})?;

select switches the same connection, and into_database does the same and returns the client. flushdb, flushall, swapdb, and move_key send FLUSHDB, FLUSHALL, SWAPDB, and MOVE.

A password is sent as AUTH only when ClientOptions::password is set. A server error reply leaves the connection usable. A broken read or write does not.

Commands

idempotency_begin(key, fingerprint, owner, ttl_ms) sends native IDEM key BEGIN ... and returns the state array with optional binary result. idempotency_complete(key, fingerprint, owner, result) sends native COMPLETE with a text result and returns 1 or 0. Only an acquired reservation allows execution; a pending response must not be treated as permission to retry. Neither completion nor replay renews retention. This is not exactly-once for external effects; see the native contract.

Family Methods
Ruvio jobs lease, release, semaphore, setv, limit, once, take, getex, incr_by_max, changes_start, changes
Strings get, get_string, set, set_with, set_expires_at, set_keep_ttl, set_and_get, mget, mset, getset, append, strlen, incr, incr_by, decr, decr_by
Keys del, del_key, unlink, unlink_key, exists, exists_key, key_type, rename, scan, dbsize, expire, pexpire, expire_at, ttl, pttl, persist
Lists lpush, rpush, lpop, lpop_count, rpop, rpop_count, llen, lindex, lrange, ltrim, blpop, blpop_key, brpop, brpop_key
Sets sadd, srem, sismember, scard, smembers, sinter, sunion, sdiff, sinterstore, sunionstore, sdiffstore, smove, spop, spop_count, srandmember, srandmember_count
Hashes hset, hget, hdel, hlen, hgetall, hmget, hexists, hkeys, hvals, hincr_by, hsetnx, hstrlen, hscan, hexpire, hexpire_with, hpexpire, hexpire_at, httl, hpttl, hexpiretime, hpexpiretime, hpersist
Sorted sets zadd, zadd_with, zadd_incr, zincr_by, zrange, zrevrange, zrange_with_scores, zrevrange_with_scores, zrange_by_score_with_scores, zrem, zcard, zscore, zrank, zrevrank, zpopmin, zpopmax, bzpopmin, bzpopmax, zremrangebyrank, zremrangebyscore, zremrangebylex
Bloom filters bf_reserve, bf_add, bf_exists
Streams xadd, xadd_maxlen, xlen, xrange, xrevrange, xdel, xtrim_maxlen, xread, xgroup_create, xreadgroup, xgroup_destroy, xgroup_setid, xgroup_delconsumer, xack, xpending_summary, xpending, xclaim, xclaim_ids
Transactions multi, exec, discard, watch, watch_key, unwatch
Scripts and functions eval, evalsha, script_load, script_exists, script_flush, script_kill, fcall, function_load, function_list, function_delete
Pub/Sub publish, spublish, subscribe, unsubscribe, psubscribe, punsubscribe, ssubscribe, sunsubscribe, next_message, read_message, subscriber, listen
Server ping, ping_message, info, config_get, save, bgsave, select, flushdb, flushall, swapdb, move_key, hello, client_getname, client_setname, client_tracking, cluster_slots, cluster_nodes
ACL acl_whoami, acl_users, acl_list, acl_getuser, acl_cat, acl_setuser

SortedSetAddOptions, StreamReadOptions, StreamPendingFilter, and StreamClaimOptions carry the optional ZADD, XREAD/XREADGROUP, XPENDING, and XCLAIM arguments. eval, evalsha, and fcall take keys and arguments separately and send keys.len() as the key count.

Inside MULTI the server answers QUEUED instead of the typed reply, so queue commands with execute and read the results from exec:

db.watch(&["balance"])?;
db.multi()?;
db.execute(&["INCRBY", "balance", "10"])?;
db.execute(&["GET", "balance"])?;
let replies = db.exec()?; // None when a watched key changed

subscribe opens a second socket, the same way Ruvio.Client does. Keep publish, GET, and PING on the same client. listen (callback) or next_message reads that subscriber socket. subscriber() still returns the extra socket if you want it separately. set_read_timeout bounds the wait.

A dropped socket is opened again before the next command, including a shard socket. The command that discovers the drop returns an error. on_connection_lost and on_connection_restored report the change. Set ClientOptions::reconnect to false and call reconnect yourself, or set_reconnect_delay to own the wait. Returning None stops the retry.

db.subscribe(&["news"])?;
assert_eq!(db.ping()?, "PONG");
db.listen(|message| println!("{}", String::from_utf8_lossy(&message.payload)))?;

Tests

cargo test runs the unit tests. RUVIO_TEST_ADDR=127.0.0.1:16390 cargo test also runs tests/live.rs against a server on that address. Those tests use databases 0 to 7 and flush them.

0.2.5

  • geo_add, geo_dist, geo_hash, geo_pos, geo_search, geo_search_store, geo_radius, and geo_radius_by_member.

0.2.4

  • A dropped command or shard socket is opened again before the next command. The command that sees the drop returns an error. Subscriptions on a dedicated Pub/Sub socket are sent again after it reopens.
  • on_connection_lost and on_connection_restored take a ConnectionNotice (host, port, subscriber).
  • ClientOptions::reconnect (default true), max_reconnect_attempts, reconnect_base_delay, reconnect_max_delay, set_reconnect_delay, and reconnect. None from the delay stops the retry. max_reconnect_attempts is not applied while a custom delay is set.

0.2.3

  • subscribe / psubscribe / ssubscribe open a dedicated Pub/Sub socket so command traffic stays on the original connection. subscriber() still returns that extra socket if you want it separately. listen calls a callback for each delivery.

0.2.2

  • del_key, unlink_key, exists_key, watch_key, blpop_key, and brpop_key take one key. The slice methods remain for many keys.

0.2.1

  • A sharded server is discovered from one seed address. Commands are routed by hash slot, MOVED is followed, and MGET / MSET / DEL split across shards.
  • cluster_shards, cluster_info, cluster_myid, cluster_keyslot, readonly, readwrite.
  • hash_slot, is_cluster, shard_count, and ClientOptions::discover_cluster.

0.2.0

  • Typed methods for every command Ruvio supports.
  • ClientOptions::database, connect_database, connect_ip, connect_addr, select, into_database, and database.
  • set_expires_at (PXAT), set_keep_ttl (KEEPTTL), set_and_get (GET), and ping_message.
  • expire sends PEXPIRE when the duration is not whole seconds. It used to round up to the next second.
  • RespValue::as_integer, as_array, and is_null.
  • The minimum Rust version is 1.87.