1use async_trait::async_trait;
2use std::collections::HashMap;
3
4use crate::errors::Error;
5use crate::errors::Result;
6use crate::request::{delete, get, get_vec, put, put_body};
7use crate::{Client, QueryMeta, QueryOptions, WriteMeta, WriteOptions};
8
9#[derive(Clone, Default, Eq, PartialEq, Serialize, Deserialize, Debug)]
10#[serde(default)]
11#[allow(clippy::upper_case_acronyms)]
12pub struct KVPair {
13 pub Key: String,
14 pub CreateIndex: Option<u64>,
15 pub ModifyIndex: Option<u64>,
16 pub LockIndex: Option<u64>,
17 pub Flags: Option<u64>,
18 pub Value: String,
19 pub Session: Option<String>,
20}
21
22#[allow(clippy::upper_case_acronyms)]
23#[async_trait]
24pub trait KV {
25 async fn acquire(&self, _: &KVPair, _: Option<&WriteOptions>) -> Result<(bool, WriteMeta)>;
26 async fn delete(&self, _: &str, _: Option<&WriteOptions>) -> Result<(bool, WriteMeta)>;
27 async fn get(&self, _: &str, _: Option<&QueryOptions>) -> Result<(Option<KVPair>, QueryMeta)>;
28 async fn list(&self, _: &str, _: Option<&QueryOptions>) -> Result<(Vec<KVPair>, QueryMeta)>;
29 async fn put(&self, _: &KVPair, _: Option<&WriteOptions>) -> Result<(bool, WriteMeta)>;
30 async fn release(&self, _: &KVPair, _: Option<&WriteOptions>) -> Result<(bool, WriteMeta)>;
31}
32
33#[async_trait]
34impl KV for Client {
35 async fn acquire(&self, pair: &KVPair, o: Option<&WriteOptions>) -> Result<(bool, WriteMeta)> {
36 let mut params = HashMap::new();
37 if let Some(i) = pair.Flags {
38 if i != 0 {
39 params.insert(String::from("flags"), i.to_string());
40 }
41 }
42 if let Some(ref session) = pair.Session {
43 params.insert(String::from("acquire"), session.to_owned());
44 let path = format!("/v1/kv/{}", pair.Key);
45 put(&path, Some(&pair.Value), &self.config, params, o).await
46 } else {
47 Err(Error::from("Session flag is required to acquire lock"))
48 }
49 }
50
51 async fn delete(&self, key: &str, options: Option<&WriteOptions>) -> Result<(bool, WriteMeta)> {
52 let path = format!("/v1/kv/{}", key);
53 delete(&path, &self.config, HashMap::new(), options).await
54 }
55
56 async fn get(
57 &self,
58 key: &str,
59 options: Option<&QueryOptions>,
60 ) -> Result<(Option<KVPair>, QueryMeta)> {
61 let path = format!("/v1/kv/{}", key);
62 let x: Result<(Vec<KVPair>, QueryMeta)> =
63 get(&path, &self.config, HashMap::new(), options).await;
64 x.map(|r| (r.0.first().cloned(), r.1))
65 }
66
67 async fn list(
68 &self,
69 prefix: &str,
70 o: Option<&QueryOptions>,
71 ) -> Result<(Vec<KVPair>, QueryMeta)> {
72 let mut params = HashMap::new();
73 params.insert(String::from("recurse"), String::from(""));
74 let path = format!("/v1/kv/{}", prefix);
75 get_vec(&path, &self.config, params, o).await
76 }
77
78 async fn put(&self, pair: &KVPair, o: Option<&WriteOptions>) -> Result<(bool, WriteMeta)> {
79 let mut params = HashMap::new();
80 if let Some(i) = pair.Flags {
81 if i != 0 {
82 params.insert(String::from("flags"), i.to_string());
83 }
84 }
85 let path = format!("/v1/kv/{}", pair.Key);
86 put_body(&path, Some(pair.Value.clone()), &self.config, params, o).await
87 }
88
89 async fn release(&self, pair: &KVPair, o: Option<&WriteOptions>) -> Result<(bool, WriteMeta)> {
90 let mut params = HashMap::new();
91 if let Some(i) = pair.Flags {
92 if i != 0 {
93 params.insert(String::from("flags"), i.to_string());
94 }
95 }
96 if let Some(ref session) = pair.Session {
97 params.insert(String::from("release"), session.to_owned());
98 let path = format!("/v1/kv/{}", pair.Key);
99 put(&path, Some(&pair.Value), &self.config, params, o).await
100 } else {
101 Err(Error::from("Session flag is required to release a lock"))
102 }
103 }
104}