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#[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}