Skip to main content

secrets_engine_kv/
lib.rs

1use std::collections::BTreeMap;
2
3use async_trait::async_trait;
4use chrono::{DateTime, Utc};
5use secrets_core::engine::{
6    CredentialShape, EngineDoc, EngineError, EngineResult, PathDoc, SecretsEngine, TtlDoc,
7};
8use secrets_core::storage::{StorageBackend, StorageEntry};
9use serde::{Deserialize, Serialize};
10use serde_json::json;
11
12const METADATA_PREFIX: &str = "kv-metadata/";
13const DATA_PREFIX: &str = "kv-data/";
14
15#[derive(Debug, Clone, Serialize, Deserialize)]
16struct VersionMetadata {
17    created_at: DateTime<Utc>,
18    deleted: bool,
19}
20
21#[derive(Debug, Clone, Serialize, Deserialize, Default)]
22struct SecretMetadata {
23    current_version: u32,
24    versions: BTreeMap<u32, VersionMetadata>,
25}
26
27/// Versioned static KV engine (kv-v2-style): every write creates a new
28/// version, `delete` soft-deletes the current version without erasing
29/// history, `read` always returns the latest non-deleted version.
30#[derive(Default)]
31pub struct KvEngine;
32
33impl KvEngine {
34    pub fn new() -> Self {
35        Self
36    }
37
38    async fn load_metadata(
39        storage: &dyn StorageBackend,
40        path: &str,
41    ) -> EngineResult<Option<SecretMetadata>> {
42        let key = format!("{METADATA_PREFIX}{path}");
43        let Some(entry) = storage.get(&key).await? else {
44            return Ok(None);
45        };
46        let metadata: SecretMetadata =
47            serde_json::from_slice(&entry.value).map_err(|e| EngineError::Other(e.to_string()))?;
48        Ok(Some(metadata))
49    }
50
51    async fn save_metadata(
52        storage: &dyn StorageBackend,
53        path: &str,
54        metadata: &SecretMetadata,
55    ) -> EngineResult<()> {
56        let key = format!("{METADATA_PREFIX}{path}");
57        let value = serde_json::to_vec(metadata).map_err(|e| EngineError::Other(e.to_string()))?;
58        storage
59            .put(
60                &key,
61                StorageEntry {
62                    value,
63                    expires_at: None,
64                },
65            )
66            .await?;
67        Ok(())
68    }
69
70    fn data_key(path: &str, version: u32) -> String {
71        format!("{DATA_PREFIX}{path}/v{version}")
72    }
73}
74
75#[async_trait]
76impl SecretsEngine for KvEngine {
77    fn doc(&self) -> EngineDoc {
78        EngineDoc {
79            provider: "built-in".to_string(),
80            mechanism: "versioned key/value storage behind the AES-256-GCM barrier".to_string(),
81            shape: CredentialShape::StaticCustody,
82            revocable: false,
83            revoke_effect: "not applicable — nothing is leased. Delete the secret \
84                            with DELETE, and rotate it at whatever system issued it: \
85                            deleting our copy does not invalidate the credential."
86                .to_string(),
87            ttl: TtlDoc {
88                min_seconds: None,
89                max_seconds: None,
90                fixed: false,
91                note: "stored values do not expire. Writes create a new version \
92                       rather than overwriting, and DELETE soft-deletes the current one."
93                    .to_string(),
94            },
95            scoping: "by path prefix, via policy. A consumer granted read on \
96                      secret/data/thirdparty/github/ci-bot cannot reach a sibling path."
97                .to_string(),
98            root_credential: "none — but whatever is stored here IS a long-lived \
99                              credential, so treat this mount as custody of last resort."
100                .to_string(),
101            paths: vec![
102                PathDoc::new(
103                    "secret/data/{path}",
104                    &["GET", "POST", "DELETE"],
105                    "read / create / delete",
106                    "read, write or soft-delete a secret. Writing bumps the version.",
107                ),
108                PathDoc::new(
109                    "secret/metadata/{prefix}",
110                    &["GET"],
111                    "list",
112                    "list the secret paths under a prefix",
113                ),
114            ],
115            docs_url: Some("docs/delegation/README.md#five-shapes-of-delegation".to_string()),
116            caveats: vec![
117                "Writes require the `create` capability, not `update` — `update` is \
118                 never checked for this mount."
119                    .to_string(),
120                "Deletion is a soft delete: earlier versions remain in storage."
121                    .to_string(),
122            ],
123        }
124    }
125
126    async fn read(&self, storage: &dyn StorageBackend, path: &str) -> EngineResult<serde_json::Value> {
127        let metadata = Self::load_metadata(storage, path)
128            .await?
129            .ok_or(EngineError::NotFound)?;
130        let version_meta = metadata
131            .versions
132            .get(&metadata.current_version)
133            .ok_or(EngineError::NotFound)?;
134        if version_meta.deleted {
135            return Err(EngineError::NotFound);
136        }
137        let key = Self::data_key(path, metadata.current_version);
138        let entry = storage.get(&key).await?.ok_or(EngineError::NotFound)?;
139        let value: serde_json::Value =
140            serde_json::from_slice(&entry.value).map_err(|e| EngineError::Other(e.to_string()))?;
141        Ok(json!({
142            "data": value,
143            "metadata": {
144                "version": metadata.current_version,
145                "created_at": version_meta.created_at,
146            },
147        }))
148    }
149
150    async fn write(
151        &self,
152        storage: &dyn StorageBackend,
153        path: &str,
154        data: serde_json::Value,
155    ) -> EngineResult<()> {
156        let mut metadata = Self::load_metadata(storage, path).await?.unwrap_or_default();
157        let new_version = metadata.current_version + 1;
158        let key = Self::data_key(path, new_version);
159        let value = serde_json::to_vec(&data).map_err(|e| EngineError::Other(e.to_string()))?;
160        storage
161            .put(
162                &key,
163                StorageEntry {
164                    value,
165                    expires_at: None,
166                },
167            )
168            .await?;
169        metadata.current_version = new_version;
170        metadata.versions.insert(
171            new_version,
172            VersionMetadata {
173                created_at: Utc::now(),
174                deleted: false,
175            },
176        );
177        Self::save_metadata(storage, path, &metadata).await
178    }
179
180    async fn delete(&self, storage: &dyn StorageBackend, path: &str) -> EngineResult<()> {
181        let mut metadata = Self::load_metadata(storage, path)
182            .await?
183            .ok_or(EngineError::NotFound)?;
184        if let Some(version_meta) = metadata.versions.get_mut(&metadata.current_version) {
185            version_meta.deleted = true;
186        }
187        Self::save_metadata(storage, path, &metadata).await
188    }
189
190    async fn list(&self, storage: &dyn StorageBackend, prefix: &str) -> EngineResult<Vec<String>> {
191        let full_prefix = format!("{METADATA_PREFIX}{prefix}");
192        let keys = storage.list(&full_prefix).await?;
193        Ok(keys
194            .into_iter()
195            .filter_map(|k| k.strip_prefix(METADATA_PREFIX).map(|s| s.to_string()))
196            .collect())
197    }
198}
199
200#[cfg(test)]
201mod tests {
202    use super::*;
203    use secrets_core::storage::StorageResult;
204    use std::collections::HashMap;
205    use std::sync::Mutex;
206
207    #[derive(Default)]
208    struct MemStorage(Mutex<HashMap<String, StorageEntry>>);
209
210    #[async_trait]
211    impl StorageBackend for MemStorage {
212        async fn get(&self, path: &str) -> StorageResult<Option<StorageEntry>> {
213            Ok(self.0.lock().unwrap().get(path).cloned())
214        }
215        async fn put(&self, path: &str, entry: StorageEntry) -> StorageResult<()> {
216            self.0.lock().unwrap().insert(path.to_string(), entry);
217            Ok(())
218        }
219        async fn delete(&self, path: &str) -> StorageResult<()> {
220            self.0.lock().unwrap().remove(path);
221            Ok(())
222        }
223        async fn list(&self, prefix: &str) -> StorageResult<Vec<String>> {
224            Ok(self
225                .0
226                .lock()
227                .unwrap()
228                .keys()
229                .filter(|k| k.starts_with(prefix))
230                .cloned()
231                .collect())
232        }
233    }
234
235    #[tokio::test]
236    async fn write_then_read_round_trip() {
237        let storage = MemStorage::default();
238        let engine = KvEngine::new();
239        engine
240            .write(&storage, "app/db", json!({"password": "hunter2"}))
241            .await
242            .unwrap();
243        let result = engine.read(&storage, "app/db").await.unwrap();
244        assert_eq!(result["data"]["password"], "hunter2");
245        assert_eq!(result["metadata"]["version"], 1);
246    }
247
248    #[tokio::test]
249    async fn write_bumps_version() {
250        let storage = MemStorage::default();
251        let engine = KvEngine::new();
252        engine.write(&storage, "app/db", json!({"v": 1})).await.unwrap();
253        engine.write(&storage, "app/db", json!({"v": 2})).await.unwrap();
254        let result = engine.read(&storage, "app/db").await.unwrap();
255        assert_eq!(result["data"]["v"], 2);
256        assert_eq!(result["metadata"]["version"], 2);
257    }
258
259    #[tokio::test]
260    async fn soft_delete_hides_current_version() {
261        let storage = MemStorage::default();
262        let engine = KvEngine::new();
263        engine.write(&storage, "app/db", json!({"v": 1})).await.unwrap();
264        engine.delete(&storage, "app/db").await.unwrap();
265        assert!(matches!(
266            engine.read(&storage, "app/db").await,
267            Err(EngineError::NotFound)
268        ));
269    }
270
271    #[tokio::test]
272    async fn read_missing_returns_not_found() {
273        let storage = MemStorage::default();
274        let engine = KvEngine::new();
275        assert!(matches!(
276            engine.read(&storage, "nope").await,
277            Err(EngineError::NotFound)
278        ));
279    }
280
281    #[tokio::test]
282    async fn list_returns_paths_under_prefix() {
283        let storage = MemStorage::default();
284        let engine = KvEngine::new();
285        engine.write(&storage, "app/db", json!({})).await.unwrap();
286        engine.write(&storage, "app/api", json!({})).await.unwrap();
287        engine.write(&storage, "other/x", json!({})).await.unwrap();
288        let mut listed = engine.list(&storage, "app/").await.unwrap();
289        listed.sort();
290        assert_eq!(listed, vec!["app/api".to_string(), "app/db".to_string()]);
291    }
292}