1use std::time::{Duration, SystemTime};
11
12use ruvio_client::{
13 Client, FieldExpireCondition, SetMode, SortedSetAddOptions, StreamClaimOptions,
14 StreamPendingFilter, StreamReadOptions,
15};
16
17fn main() -> Result<(), ruvio_client::Error> {
18 let host = std::env::args()
19 .nth(1)
20 .unwrap_or_else(|| "127.0.0.1".to_owned());
21 let port = std::env::args()
22 .nth(2)
23 .and_then(|text| text.parse().ok())
24 .unwrap_or(6379);
25 let mut db = Client::connect(&host, port)?;
26
27 db.on_connection_lost(|notice| {
28 println!(
29 "connection lost {}:{} subscriber={}",
30 notice.host, notice.port, notice.subscriber
31 );
32 });
33
34 db.on_connection_restored(|notice| {
35 println!(
36 "connection restored {}:{} subscriber={}",
37 notice.host, notice.port, notice.subscriber
38 );
39 });
40
41 db.set_reconnect_delay(|attempt| {
42 if attempt >= 8 {
43 None
44 } else {
45 Some(Duration::from_millis(100 * 2u64.pow(attempt - 1)))
46 }
47 });
48
49 println!(
50 "connected cluster={} shards={} slot={}",
51 db.is_cluster(),
52 db.shard_count(),
53 ruvio_client::hash_slot("{tour}")
54 );
55 println!("reconnect armed; the command that sees a drop fails, the next one opens the socket");
56
57 if !db.is_cluster() {
58 db.select(15)?;
59 db.flushdb()?;
60 println!("selected database {}", db.database());
61 }
62
63 strings(&mut db)?;
64 owned(&mut db)?;
65 lists(&mut db)?;
66 sets(&mut db)?;
67 hashes(&mut db)?;
68 sorted_sets(&mut db)?;
69 bloom(&mut db)?;
70 streams(&mut db)?;
71 transactions(&mut db)?;
72 server(&mut db)?;
73 pubsub(&host, port)?;
74
75 println!("tour finished");
76
77 Ok(())
78}
79
80fn strings(db: &mut Client) -> Result<(), ruvio_client::Error> {
81 let key = "{tour}:s";
82
83 db.set(key, "hello")?;
84 println!("get {:?}", db.get_string(key)?);
85 db.set_with(key, "nx", None, SetMode::IfNotExists)?;
86 db.set_expires_at(
87 key,
88 "later",
89 SystemTime::now() + Duration::from_secs(120),
90 SetMode::Always,
91 )?;
92 db.set_keep_ttl(key, "kept", SetMode::Always)?;
93 println!(
94 "set_and_get {:?}",
95 db.set_and_get(key, "now", None, SetMode::Always)?
96 );
97 db.mset(&[("{tour}:a", "1"), ("{tour}:b", "2")])?;
98 println!(
99 "mget {:?}",
100 db.mget(&["{tour}:a", "{tour}:missing", "{tour}:b"])?
101 );
102 println!("getset {:?}", db.getset(key, "10")?);
103 println!("append {}", db.append(key, "0")?);
104 println!("strlen {}", db.strlen(key)?);
105 println!("incr {}", db.incr("{tour}:n")?);
106 println!("incr_by {}", db.incr_by("{tour}:n", 5)?);
107 println!("decr {}", db.decr("{tour}:n")?);
108 println!("decr_by {}", db.decr_by("{tour}:n", 2)?);
109 println!("exists {}", db.exists(&["{tour}:a", "{tour}:b"])?);
110 println!("exists_key {}", db.exists_key(key)?);
111 println!("type {}", db.key_type(key)?);
112 db.rename("{tour}:b", "{tour}:c")?;
113 println!("unlink {}", db.unlink(&["{tour}:c"])?);
114 println!("unlink_key {}", db.unlink_key("{tour}:gone")?);
115 db.expire(key, Duration::from_secs(30))?;
116 db.pexpire("{tour}:a", Duration::from_millis(30_000))?;
117 db.expire_at("{tour}:a", SystemTime::now() + Duration::from_secs(60))?;
118 println!("ttl {} pttl {}", db.ttl(key)?, db.pttl(key)?);
119 println!("persist {}", db.persist(key)?);
120 let page = db.scan(0, Some("{tour}:*"), Some(20))?;
121
122 println!("scan cursor={} keys={}", page.cursor, page.keys.len());
123 println!("dbsize {}", db.dbsize()?);
124 println!("del {}", db.del(&["{tour}:a"])?);
125 println!("del_key {}", db.del_key("{tour}:n")?);
126
127 Ok(())
128}
129
130fn owned(db: &mut Client) -> Result<(), ruvio_client::Error> {
131 let first = db.lease("{tour}:lock", 30_000)?;
132
133 println!(
134 "lease {:?} held {:?} release {}",
135 first,
136 db.lease("{tour}:lock", 30_000)?,
137 db.release("{tour}:lock", first.unwrap_or(0) as u64)?
138 );
139 let permit = db.semaphore("{tour}:workers", 1, 30_000)?;
140
141 println!(
142 "semaphore {:?} full {:?} release {}",
143 permit,
144 db.semaphore("{tour}:workers", 1, 30_000)?,
145 db.release("{tour}:workers", permit.unwrap_or(0) as u64)?
146 );
147 println!(
148 "incr_by_max {:?} over {:?}",
149 db.incr_by_max("{tour}:seats", 3, 4)?,
150 db.incr_by_max("{tour}:seats", 2, 4)?
151 );
152 let created = db.setv("{tour}:doc", "12", None)?;
153 let next = db.setv("{tour}:doc", "11", Some(created as u64))?;
154
155 println!("setv {created} then {next}");
156 println!(
157 "setv mismatch {}",
158 db.setv("{tour}:doc", "10", Some(created as u64))
159 .err()
160 .map(|error| error.to_string())
161 .unwrap_or_else(|| "accepted".to_owned())
162 );
163 let cursor = db.changes_start()?;
164
165 db.set("{tour}:note", "paid")?;
166 let notes = db.changes(cursor as u64)?;
167
168 println!(
169 "changes cursor={cursor} notes={} {}",
170 notes.len(),
171 notes.first().map(|(_, text)| text.as_str()).unwrap_or("")
172 );
173 println!(
174 "limit {:?} {:?}",
175 db.limit("{tour}:login", 2, 60_000)?,
176 db.limit("{tour}:login", 2, 60_000)?
177 );
178 println!(
179 "once {} {}",
180 db.once("{tour}:pay", 60_000)?,
181 db.once("{tour}:pay", 60_000)?
182 );
183 db.set("{tour}:seats", "5")?;
184 println!(
185 "take {:?} {:?}",
186 db.take("{tour}:seats", 2)?,
187 db.take("{tour}:seats", 9)?
188 );
189 println!("getex {:?}", db.getex("{tour}:seats", 60_000)?);
190
191 Ok(())
192}
193
194fn lists(db: &mut Client) -> Result<(), ruvio_client::Error> {
195 let key = "{tour}:l";
196
197 println!("rpush {}", db.rpush(key, &["a", "b", "c", "d"])?);
198 println!("lpush {}", db.lpush(key, &["z"])?);
199 println!("lpop {:?}", db.lpop(key)?);
200 println!("rpop_count {:?}", db.rpop_count(key, 2)?);
201 println!("llen {}", db.llen(key)?);
202 println!("lindex {:?}", db.lindex(key, 0)?);
203 println!("lrange {:?}", db.lrange(key, 0, -1)?);
204 db.ltrim(key, 0, 0)?;
205 println!("lpop_count {:?}", db.lpop_count(key, 1)?);
206 db.rpush(key, &["only"])?;
207 println!("blpop {:?}", db.blpop_key(key, Duration::from_millis(200))?);
208 db.rpush(key, &["right"])?;
209 println!("brpop {:?}", db.brpop(&[key], Duration::from_millis(200))?);
210 println!(
211 "brpop empty {:?}",
212 db.brpop_key(key, Duration::from_millis(1))?
213 );
214
215 Ok(())
216}
217
218fn sets(db: &mut Client) -> Result<(), ruvio_client::Error> {
219 println!("sadd {}", db.sadd("{tour}:s1", &["a", "b", "c"])?);
220 db.sadd("{tour}:s2", &["b", "c", "d"])?;
221 println!("sismember {}", db.sismember("{tour}:s1", "a")?);
222 println!("scard {}", db.scard("{tour}:s1")?);
223 println!("smembers {:?}", db.smembers("{tour}:s1")?);
224 println!("sinter {:?}", db.sinter(&["{tour}:s1", "{tour}:s2"])?);
225 println!("sunion {:?}", db.sunion(&["{tour}:s1", "{tour}:s2"])?);
226 println!("sdiff {:?}", db.sdiff(&["{tour}:s1", "{tour}:s2"])?);
227 println!(
228 "sinterstore {}",
229 db.sinterstore("{tour}:s3", &["{tour}:s1", "{tour}:s2"])?
230 );
231 println!(
232 "sunionstore {}",
233 db.sunionstore("{tour}:s4", &["{tour}:s1", "{tour}:s2"])?
234 );
235 println!(
236 "sdiffstore {}",
237 db.sdiffstore("{tour}:s5", &["{tour}:s1", "{tour}:s2"])?
238 );
239 println!("smove {}", db.smove("{tour}:s1", "{tour}:s3", "a")?);
240 println!("srem {}", db.srem("{tour}:s3", &["a"])?);
241 println!("srandmember {:?}", db.srandmember("{tour}:s2")?);
242 println!(
243 "srandmember_count {:?}",
244 db.srandmember_count("{tour}:s2", 2)?
245 );
246 println!("spop_count {:?}", db.spop_count("{tour}:s5", 5)?);
247 println!("spop {:?}", db.spop("{tour}:s4")?);
248
249 Ok(())
250}
251
252fn hashes(db: &mut Client) -> Result<(), ruvio_client::Error> {
253 let key = "{tour}:h";
254
255 println!("hset {}", db.hset(key, &[("f1", "1"), ("f2", "two")])?);
256 println!("hget {:?}", db.hget(key, "f2")?);
257 println!("hincr_by {}", db.hincr_by(key, "f1", 4)?);
258 println!("hsetnx {}", db.hsetnx(key, "f1", "9")?);
259 println!("hstrlen {}", db.hstrlen(key, "f2")?);
260 println!("hlen {}", db.hlen(key)?);
261 println!("hexists {}", db.hexists(key, "f1")?);
262 println!("hkeys {:?}", db.hkeys(key)?);
263 println!("hvals {:?}", db.hvals(key)?);
264 println!("hgetall {:?}", db.hgetall(key)?);
265 println!("hmget {:?}", db.hmget(key, &["f2", "nope"])?);
266 let page = db.hscan(key, 0, Some("f*"), Some(10))?;
267
268 println!("hscan cursor={} fields={}", page.cursor, page.fields.len());
269 println!(
270 "hexpire {:?}",
271 db.hexpire(key, Duration::from_secs(60), &["f1", "nope"])?
272 );
273 println!(
274 "hexpire_with GT {:?}",
275 db.hexpire_with(
276 key,
277 Duration::from_secs(120),
278 FieldExpireCondition::IfGreater,
279 &["f1"]
280 )?
281 );
282 println!(
283 "hpexpire {:?}",
284 db.hpexpire(key, Duration::from_millis(90_500), &["f2"])?
285 );
286 println!("httl {:?}", db.httl(key, &["f1", "f2"])?);
287 println!("hpttl {:?}", db.hpttl(key, &["f1"])?);
288 println!(
289 "hexpire_at {:?}",
290 db.hexpire_at(key, SystemTime::now() + Duration::from_secs(300), &["f2"])?
291 );
292 println!("hexpiretime {:?}", db.hexpiretime(key, &["f2"])?);
293 println!("hpexpiretime {:?}", db.hpexpiretime(key, &["f2"])?);
294 println!("hpersist {:?}", db.hpersist(key, &["f1", "f2"])?);
295 println!("hdel {}", db.hdel(key, &["f1"])?);
296
297 Ok(())
298}
299
300fn sorted_sets(db: &mut Client) -> Result<(), ruvio_client::Error> {
301 let key = "{tour}:z";
302
303 println!("zadd {}", db.zadd(key, &[("ada", 100.0), ("bob", 50.0)])?);
304 println!(
305 "geoadd {}",
306 db.geo_add(key, &[(13.361389, 38.115556, "Palermo")])?
307 );
308 println!(
309 "geosearch {:?}",
310 db.geo_search(key, 13.36, 38.11, 20.0, "km")?
311 );
312 let greater = SortedSetAddOptions {
313 greater_than: true,
314 changed: true,
315 ..SortedSetAddOptions::default()
316 };
317
318 println!(
319 "zadd_with {}",
320 db.zadd_with(key, greater, &[("bob", 10.0)])?
321 );
322 println!("zincr_by {}", db.zincr_by(key, "bob", 5.5)?);
323 println!(
324 "zadd_incr {:?}",
325 db.zadd_incr(key, "cy", 1.0, SortedSetAddOptions::default())?
326 );
327 println!("zrange {:?}", db.zrange(key, 0, -1)?);
328 println!("zrevrange {:?}", db.zrevrange(key, 0, 0)?);
329 println!("zrange_with_scores {:?}", db.zrange_with_scores(key, 0, 0)?);
330 println!(
331 "zrevrange_with_scores {:?}",
332 db.zrevrange_with_scores(key, 0, 0)?
333 );
334 println!(
335 "zrange_by_score {:?}",
336 db.zrange_by_score_with_scores(key, "(1", "+inf")?
337 );
338 println!("zscore {:?}", db.zscore(key, "ada")?);
339 println!("zrank {:?}", db.zrank(key, "ada")?);
340 println!("zrevrank {:?}", db.zrevrank(key, "ada")?);
341 println!("zrem {}", db.zrem(key, &["cy"])?);
342 println!("zcard {}", db.zcard(key)?);
343
344 let queue = "{tour}:zq";
345
346 db.zadd(
347 queue,
348 &[("a", 1.0), ("b", 2.0), ("c", 3.0), ("d", 4.0), ("e", 5.0)],
349 )?;
350 println!("zpopmin {:?}", db.zpopmin(queue, 1)?);
351 println!("zpopmax {:?}", db.zpopmax(queue, 1)?);
352 println!(
353 "bzpopmin {:?}",
354 db.bzpopmin(&[queue], Duration::from_secs(1))?
355 );
356 println!(
357 "bzpopmax {:?}",
358 db.bzpopmax(&[queue], Duration::from_secs(1))?
359 );
360 db.zadd(queue, &[("a", 1.0), ("b", 2.0), ("c", 3.0), ("x", 0.0)])?;
361 println!("zremrangebyrank {}", db.zremrangebyrank(queue, 0, 0)?);
362 println!(
363 "zremrangebyscore {}",
364 db.zremrangebyscore(queue, "(1", "2")?
365 );
366 println!("zremrangebylex {}", db.zremrangebylex(queue, "-", "+")?);
367
368 Ok(())
369}
370
371fn bloom(db: &mut Client) -> Result<(), ruvio_client::Error> {
372 db.bf_reserve("{tour}:bf", 0.01, 1000)?;
373 println!("bf_add {}", db.bf_add("{tour}:bf", "x")?);
374 println!(
375 "bf_exists {} {}",
376 db.bf_exists("{tour}:bf", "x")?,
377 db.bf_exists("{tour}:bf", "y")?
378 );
379
380 Ok(())
381}
382
383fn streams(db: &mut Client) -> Result<(), ruvio_client::Error> {
384 let key = "{tour}:stream";
385 let first = db.xadd(key, "*", &[("f", "1")])?;
386
387 db.xadd(key, "*", &[("f", "2")])?;
388 println!(
389 "xadd_maxlen {}",
390 db.xadd_maxlen(key, 10, "*", &[("f", "3")])?
391 );
392 println!(
393 "xadd_delay {}",
394 db.xadd_delay(key, 60_000, "*", &[("f", "later")])?
395 );
396 println!("xlen {}", db.xlen(key)?);
397 println!("xrange {}", db.xrange(key, "-", "+", None)?[0].id);
398 println!(
399 "xrevrange {}",
400 db.xrevrange(key, "+", "-", Some(1))?[0].fields[0].1
401 );
402 let read = db.xread(
403 &StreamReadOptions {
404 count: Some(2),
405 ..StreamReadOptions::default()
406 },
407 &[(key, "0")],
408 )?;
409
410 println!("xread {}", read[0].entries.len());
411 db.xgroup_create(key, "g", "0", false)?;
412 db.xgroup_create("{tour}:fresh", "g", "$", true)?;
413 let delivered = db.xreadgroup("g", "c1", &StreamReadOptions::default(), &[(key, ">")])?;
414
415 println!("xreadgroup {}", delivered[0].entries.len());
416 let pending = db.xpending(key, "g", "-", "+", 10, &StreamPendingFilter::default())?;
417
418 println!(
419 "xpending {} summary-array {}",
420 pending.len(),
421 db.xpending_summary(key, "g")?.as_array().is_some()
422 );
423 let claimed = db.xclaim(
424 key,
425 "g",
426 "c2",
427 Duration::ZERO,
428 &[&first],
429 &StreamClaimOptions::default(),
430 )?;
431
432 println!("xclaim {}", claimed[0].id);
433 println!(
434 "xclaim_ids {:?}",
435 db.xclaim_ids(
436 key,
437 "g",
438 "c2",
439 Duration::ZERO,
440 &[&first],
441 &StreamClaimOptions::default()
442 )?
443 );
444 println!("xack {}", db.xack(key, "g", &[&first])?);
445 println!(
446 "xgroup_delconsumer {}",
447 db.xgroup_delconsumer(key, "g", "c1")?
448 );
449 db.xgroup_setid(key, "g", "0")?;
450 println!("xgroup_destroy {}", db.xgroup_destroy(key, "g")?);
451 println!("xdel {}", db.xdel(key, &[&first])?);
452 println!("xtrim {}", db.xtrim_maxlen(key, 1)?);
453
454 Ok(())
455}
456
457fn transactions(db: &mut Client) -> Result<(), ruvio_client::Error> {
458 db.watch_key("{tour}:w")?;
459 db.watch(&["{tour}:w"])?;
460 db.multi()?;
461 db.execute(&["SET", "{tour}:w", "1"])?;
462 db.execute(&["INCR", "{tour}:w"])?;
463 println!(
464 "exec {} replies",
465 db.exec()?.map(|replies| replies.len()).unwrap_or(0)
466 );
467 db.unwatch()?;
468 db.multi()?;
469 db.discard()?;
470 println!(
471 "eval {:?}",
472 db.eval("return ARGV[1]", &[], &["hi"])?.as_string()?
473 );
474 let sha = db.script_load("return ruvio.call('GET', KEYS[1])")?;
475
476 println!("script_exists {:?}", db.script_exists(&[&sha, "0000"])?);
477 println!(
478 "evalsha {:?}",
479 db.evalsha(&sha, &["{tour}:w"], &[])?.as_string()?
480 );
481 db.function_load("tour_echo", "return ARGV[1]")?;
482 println!(
483 "fcall {:?}",
484 db.fcall("tour_echo", &[], &["x"])?.as_string()?
485 );
486 println!(
487 "function_list array {}",
488 db.function_list()?.as_array().is_some()
489 );
490 println!("function_delete {}", db.function_delete("tour_echo")?);
491 db.script_flush()?;
492 match db.script_kill() {
493 Ok(()) => println!("script_kill ok"),
494 Err(error) => println!("script_kill {error}"),
495 }
496
497 if !db.is_cluster() {
498 db.set("{tour}:move", "v")?;
499 println!("move {}", db.move_key("{tour}:move", 14)?);
500 db.swapdb(14, 15)?;
501 println!("swapdb back");
502 db.swapdb(14, 15)?;
503 }
504
505 Ok(())
506}
507
508fn server(db: &mut Client) -> Result<(), ruvio_client::Error> {
509 println!("ping {}", db.ping()?);
510 println!("ping_message {}", db.ping_message("tour")?);
511 println!("info has ruvio {}", db.info(None)?.contains("ruvio"));
512 println!("info server {}", db.info(Some("server"))?.len());
513 println!("config_get {}", db.config_get("databases")?.len());
514 println!("whoami {}", db.acl_whoami()?);
515 println!("acl_users {:?}", db.acl_users()?);
516 println!("acl_list {}", db.acl_list()?.len());
517 println!("acl_cat {}", db.acl_cat()?.len());
518 println!("acl_getuser null {}", db.acl_getuser("default")?.is_null());
519 db.acl_setuser("tour-reader", &["on", "~{tour}:*", "+get"])?;
520 db.client_setname("ruvio-tour")?;
521 println!("client_getname {:?}", db.client_getname()?);
522 db.client_tracking(false)?;
523 println!("cluster_keyslot {}", db.cluster_keyslot("{tour}")?);
524
525 if db.is_cluster() {
526 println!(
527 "cluster_nodes {}",
528 db.cluster_nodes()?.chars().take(40).collect::<String>()
529 );
530 println!(
531 "cluster_slots array {}",
532 db.cluster_slots()?.as_array().is_some()
533 );
534 println!(
535 "cluster_shards array {}",
536 db.cluster_shards()?.as_array().is_some()
537 );
538 println!("cluster_info {}", db.cluster_info()?.len());
539 println!("cluster_myid {}", db.cluster_myid()?);
540 db.readonly()?;
541 db.readwrite()?;
542 }
543
544 db.save()?;
545 db.bgsave()?;
546 println!("hello array {}", db.hello(2)?.as_array().is_some());
547 db.execute_many(&[&["PING"], &["PING"]])?;
548
549 Ok(())
550}
551
552fn pubsub(host: &str, port: u16) -> Result<(), ruvio_client::Error> {
553 let mut commands = Client::connect(host, port)?;
554 let mut listener = commands.subscriber()?;
555
556 listener.subscribe(&["{tour}:news"])?;
557 listener.psubscribe(&["{tour}:n*"])?;
558 listener.set_read_timeout(Some(Duration::from_secs(2)))?;
559 println!("publish {}", commands.publish("{tour}:news", "hello-tour")?);
560
561 let exact = listener.next_message()?;
562 let pattern = listener.next_message()?;
563
564 println!(
565 "message {} {:?}",
566 exact.channel,
567 String::from_utf8_lossy(&exact.payload)
568 );
569 println!("pmessage {:?}", pattern.pattern);
570 listener.unsubscribe(&["{tour}:news"])?;
571 listener.punsubscribe(&["{tour}:n*"])?;
572 listener.ssubscribe(&["{tour}:shard"])?;
573 println!("spublish {}", commands.spublish("{tour}:shard", "sharded")?);
574 println!("smessage {}", listener.next_message()?.kind);
575 listener.sunsubscribe(&["{tour}:shard"])?;
576 println!("ping after subscribe {}", commands.ping()?);
577
578 Ok(())
579}