Skip to main content

north_consul/
kv.rs

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}