1use std::path::PathBuf;
8use std::sync::Arc;
9
10use async_trait::async_trait;
11use base64::Engine;
12use bytes::Bytes;
13use meerkat_core::{BlobId, BlobPayload, BlobRef, BlobStore, BlobStoreError};
14use object_store::path::Path as ObjectPath;
15use object_store::{ObjectStore, ObjectStoreExt, PutPayload};
16use serde::{Deserialize, Serialize};
17use sha2::{Digest, Sha256};
18use tokio::sync::RwLock;
19
20#[derive(Debug, Clone, PartialEq, Eq)]
21pub struct BinaryBlobPayload {
22 pub blob_id: BlobId,
23 pub media_type: String,
24 pub size: u64,
25 pub data: Bytes,
26}
27
28#[async_trait]
29pub trait BinaryBlobStore: Send + Sync {
30 async fn put_bytes(&self, media_type: &str, data: Bytes) -> Result<BlobRef, BlobStoreError>;
31 async fn get_bytes(&self, blob_id: &BlobId) -> Result<BinaryBlobPayload, BlobStoreError>;
32 async fn delete(&self, blob_id: &BlobId) -> Result<(), BlobStoreError>;
33 fn is_persistent(&self) -> bool;
34}
35
36#[derive(Debug, Clone, Serialize, Deserialize)]
37struct BlobMetadata {
38 media_type: String,
39 size: u64,
40 created_at_ms: u64,
41 source: String,
42}
43
44#[derive(Debug, Clone, Serialize, Deserialize)]
45struct LegacyStoredBlob {
46 media_type: String,
47 data: String,
48}
49
50#[derive(Debug)]
51enum BlobObjectBackend {
52 ObjectStore {
53 store: Arc<dyn ObjectStore>,
54 persistent: bool,
55 },
56 Memory {
57 blobs: RwLock<std::collections::BTreeMap<BlobId, (String, Bytes)>>,
58 },
59}
60
61pub struct ObjectStoreBlobStore {
62 backend: BlobObjectBackend,
63 legacy_root: Option<PathBuf>,
64}
65
66impl ObjectStoreBlobStore {
67 pub fn local(root: PathBuf) -> Result<Self, BlobStoreError> {
68 std::fs::create_dir_all(&root).map_err(|err| BlobStoreError::Internal(err.to_string()))?;
69 let store = object_store::local::LocalFileSystem::new_with_prefix(&root)
70 .map_err(|err| BlobStoreError::Internal(err.to_string()))?;
71 Ok(Self {
72 backend: BlobObjectBackend::ObjectStore {
73 store: Arc::new(store),
74 persistent: true,
75 },
76 legacy_root: Some(root),
77 })
78 }
79
80 pub fn memory() -> Self {
81 Self {
82 backend: BlobObjectBackend::Memory {
83 blobs: RwLock::new(std::collections::BTreeMap::new()),
84 },
85 legacy_root: None,
86 }
87 }
88
89 fn object_path(blob_id: &BlobId) -> ObjectPath {
90 ObjectPath::from(format!("objects/{}.bin", storage_key(blob_id)))
91 }
92
93 fn meta_path(blob_id: &BlobId) -> ObjectPath {
94 ObjectPath::from(format!("meta/{}.json", storage_key(blob_id)))
95 }
96
97 fn legacy_path(&self, blob_id: &BlobId) -> Option<PathBuf> {
98 let root = self.legacy_root.as_ref()?;
99 let key = storage_key(blob_id);
100 let prefix = key.get(0..2).unwrap_or("xx");
101 Some(root.join(prefix).join(format!("{key}.json")))
102 }
103
104 async fn read_legacy(&self, blob_id: &BlobId) -> Result<BinaryBlobPayload, BlobStoreError> {
105 let path = self
106 .legacy_path(blob_id)
107 .ok_or_else(|| BlobStoreError::NotFound(blob_id.clone()))?;
108 let bytes = tokio::fs::read(&path).await.map_err(|err| {
109 if err.kind() == std::io::ErrorKind::NotFound {
110 BlobStoreError::NotFound(blob_id.clone())
111 } else {
112 BlobStoreError::ReadFailed(err.to_string())
113 }
114 })?;
115 let stored: LegacyStoredBlob = serde_json::from_slice(&bytes)
116 .map_err(|err| BlobStoreError::ReadFailed(err.to_string()))?;
117 let decoded = base64::engine::general_purpose::STANDARD
118 .decode(stored.data.as_bytes())
119 .map_err(|err| {
120 BlobStoreError::ReadFailed(format!("invalid legacy blob base64: {err}"))
121 })?;
122 Ok(BinaryBlobPayload {
123 blob_id: blob_id.clone(),
124 media_type: stored.media_type,
125 size: decoded.len() as u64,
126 data: Bytes::from(decoded),
127 })
128 }
129}
130
131#[async_trait]
132impl BinaryBlobStore for ObjectStoreBlobStore {
133 async fn put_bytes(&self, media_type: &str, data: Bytes) -> Result<BlobRef, BlobStoreError> {
134 let blob_id = compute_blob_id(media_type, &data);
135 match &self.backend {
136 BlobObjectBackend::ObjectStore { store, .. } => {
137 let meta = BlobMetadata {
138 media_type: media_type.to_string(),
139 size: data.len() as u64,
140 created_at_ms: current_time_ms(),
141 source: "mobkit".to_string(),
142 };
143 let meta_bytes = serde_json::to_vec(&meta)
144 .map_err(|err| BlobStoreError::WriteFailed(err.to_string()))?;
145 store
146 .put(&Self::object_path(&blob_id), PutPayload::from(data.clone()))
147 .await
148 .map_err(|err| BlobStoreError::WriteFailed(err.to_string()))?;
149 store
150 .put(
151 &Self::meta_path(&blob_id),
152 PutPayload::from(Bytes::from(meta_bytes)),
153 )
154 .await
155 .map_err(|err| BlobStoreError::WriteFailed(err.to_string()))?;
156 }
157 BlobObjectBackend::Memory { blobs } => {
158 blobs
159 .write()
160 .await
161 .entry(blob_id.clone())
162 .or_insert_with(|| (media_type.to_string(), data.clone()));
163 }
164 }
165 Ok(BlobRef {
166 blob_id,
167 media_type: media_type.to_string(),
168 })
169 }
170
171 async fn get_bytes(&self, blob_id: &BlobId) -> Result<BinaryBlobPayload, BlobStoreError> {
172 if !is_valid_blob_id(blob_id) {
173 return Err(BlobStoreError::NotFound(blob_id.clone()));
174 }
175 match &self.backend {
176 BlobObjectBackend::ObjectStore { store, .. } => {
177 let meta_bytes = match store.get(&Self::meta_path(blob_id)).await {
178 Ok(result) => result
179 .bytes()
180 .await
181 .map_err(|err| BlobStoreError::ReadFailed(err.to_string()))?,
182 Err(object_store::Error::NotFound { .. }) => {
183 return self.read_legacy(blob_id).await;
184 }
185 Err(err) => return Err(BlobStoreError::ReadFailed(err.to_string())),
186 };
187 let meta: BlobMetadata = serde_json::from_slice(&meta_bytes)
188 .map_err(|err| BlobStoreError::ReadFailed(err.to_string()))?;
189 let data = match store.get(&Self::object_path(blob_id)).await {
190 Ok(result) => result
191 .bytes()
192 .await
193 .map_err(|err| BlobStoreError::ReadFailed(err.to_string()))?,
194 Err(object_store::Error::NotFound { .. }) => {
195 return self.read_legacy(blob_id).await;
196 }
197 Err(err) => return Err(BlobStoreError::ReadFailed(err.to_string())),
198 };
199 Ok(BinaryBlobPayload {
200 blob_id: blob_id.clone(),
201 media_type: meta.media_type,
202 size: meta.size,
203 data,
204 })
205 }
206 BlobObjectBackend::Memory { blobs } => {
207 let blobs = blobs.read().await;
208 let (media_type, data) = blobs
209 .get(blob_id)
210 .ok_or_else(|| BlobStoreError::NotFound(blob_id.clone()))?;
211 Ok(BinaryBlobPayload {
212 blob_id: blob_id.clone(),
213 media_type: media_type.clone(),
214 size: data.len() as u64,
215 data: data.clone(),
216 })
217 }
218 }
219 }
220
221 async fn delete(&self, blob_id: &BlobId) -> Result<(), BlobStoreError> {
222 if !is_valid_blob_id(blob_id) {
223 return Err(BlobStoreError::NotFound(blob_id.clone()));
224 }
225 match &self.backend {
226 BlobObjectBackend::ObjectStore { store, .. } => {
227 for path in [Self::object_path(blob_id), Self::meta_path(blob_id)] {
228 match store.delete(&path).await {
229 Ok(()) => {}
230 Err(object_store::Error::NotFound { .. }) => {}
231 Err(err) => return Err(BlobStoreError::DeleteFailed(err.to_string())),
232 }
233 }
234 }
235 BlobObjectBackend::Memory { blobs } => {
236 blobs.write().await.remove(blob_id);
237 }
238 }
239 Ok(())
240 }
241
242 fn is_persistent(&self) -> bool {
243 match &self.backend {
244 BlobObjectBackend::ObjectStore { persistent, .. } => *persistent,
245 BlobObjectBackend::Memory { .. } => false,
246 }
247 }
248}
249
250pub struct Base64BlobStoreAdapter {
251 inner: Arc<dyn BinaryBlobStore>,
252}
253
254impl Base64BlobStoreAdapter {
255 pub fn new(inner: Arc<dyn BinaryBlobStore>) -> Self {
256 Self { inner }
257 }
258}
259
260#[async_trait]
261impl BlobStore for Base64BlobStoreAdapter {
262 async fn put_image(&self, media_type: &str, data: &str) -> Result<BlobRef, BlobStoreError> {
263 let bytes = base64::engine::general_purpose::STANDARD
264 .decode(data.as_bytes())
265 .map_err(|err| BlobStoreError::WriteFailed(format!("invalid blob base64: {err}")))?;
266 self.inner.put_bytes(media_type, Bytes::from(bytes)).await
267 }
268
269 async fn get(&self, blob_id: &BlobId) -> Result<BlobPayload, BlobStoreError> {
270 let payload = self.inner.get_bytes(blob_id).await?;
271 Ok(BlobPayload {
272 blob_id: payload.blob_id,
273 media_type: payload.media_type,
274 data: base64::engine::general_purpose::STANDARD.encode(payload.data.as_ref()),
275 })
276 }
277
278 async fn delete(&self, blob_id: &BlobId) -> Result<(), BlobStoreError> {
279 self.inner.delete(blob_id).await
280 }
281
282 fn is_persistent(&self) -> bool {
283 self.inner.is_persistent()
284 }
285}
286
287pub struct BinaryBlobStoreAdapter {
288 inner: Arc<dyn BlobStore>,
289}
290
291impl BinaryBlobStoreAdapter {
292 pub fn new(inner: Arc<dyn BlobStore>) -> Self {
293 Self { inner }
294 }
295}
296
297#[async_trait]
298impl BinaryBlobStore for BinaryBlobStoreAdapter {
299 async fn put_bytes(&self, media_type: &str, data: Bytes) -> Result<BlobRef, BlobStoreError> {
300 let encoded = base64::engine::general_purpose::STANDARD.encode(data.as_ref());
301 self.inner.put_image(media_type, &encoded).await
302 }
303
304 async fn get_bytes(&self, blob_id: &BlobId) -> Result<BinaryBlobPayload, BlobStoreError> {
305 let payload = self.inner.get(blob_id).await?;
306 let decoded = base64::engine::general_purpose::STANDARD
307 .decode(payload.data.as_bytes())
308 .map_err(|err| {
309 BlobStoreError::ReadFailed(format!("invalid blob base64 payload: {err}"))
310 })?;
311 Ok(BinaryBlobPayload {
312 blob_id: payload.blob_id,
313 media_type: payload.media_type,
314 size: decoded.len() as u64,
315 data: Bytes::from(decoded),
316 })
317 }
318
319 async fn delete(&self, blob_id: &BlobId) -> Result<(), BlobStoreError> {
320 self.inner.delete(blob_id).await
321 }
322
323 fn is_persistent(&self) -> bool {
324 self.inner.is_persistent()
325 }
326}
327
328fn compute_blob_id(media_type: &str, bytes: &[u8]) -> BlobId {
329 let mut hasher = Sha256::new();
330 hasher.update(media_type.as_bytes());
331 hasher.update([0]);
332 hasher.update(bytes);
333 BlobId::new(format!("sha256:{:x}", hasher.finalize()))
334}
335
336pub fn is_valid_blob_id(blob_id: &BlobId) -> bool {
337 is_valid_blob_id_value(blob_id.as_str())
338}
339
340pub fn is_valid_blob_id_value(value: &str) -> bool {
341 let Some(hex) = value.strip_prefix("sha256:") else {
342 return false;
343 };
344 hex.len() == 64
345 && hex
346 .bytes()
347 .all(|byte| matches!(byte, b'0'..=b'9' | b'a'..=b'f'))
348}
349
350fn storage_key(blob_id: &BlobId) -> &str {
351 blob_id
352 .as_str()
353 .strip_prefix("sha256:")
354 .unwrap_or(blob_id.as_str())
355}
356
357fn current_time_ms() -> u64 {
358 std::time::SystemTime::now()
359 .duration_since(std::time::UNIX_EPOCH)
360 .map(|duration| duration.as_millis() as u64)
361 .unwrap_or_default()
362}
363
364#[cfg(test)]
365#[allow(clippy::expect_used)]
366mod tests {
367 use super::*;
368
369 #[tokio::test]
370 async fn binary_blob_ids_hash_raw_bytes_not_base64_text() {
371 let store = ObjectStoreBlobStore::memory();
372 let first = store
373 .put_bytes("image/png", Bytes::from_static(b"abc"))
374 .await
375 .expect("put");
376 let second = store
377 .put_bytes("image/png", Bytes::from_static(b"abc"))
378 .await
379 .expect("put same");
380 let third = store
381 .put_bytes("image/jpeg", Bytes::from_static(b"abc"))
382 .await
383 .expect("put media");
384
385 assert_eq!(first.blob_id, second.blob_id);
386 assert_ne!(first.blob_id, third.blob_id);
387 assert_ne!(
388 first.blob_id,
389 compute_legacy_blob_id("image/png", "YWJj"),
390 "raw-byte hashing must not reuse the legacy base64-text hash"
391 );
392 }
393
394 #[tokio::test]
395 async fn base64_adapter_roundtrips_at_json_boundary() {
396 let binary: Arc<dyn BinaryBlobStore> = Arc::new(ObjectStoreBlobStore::memory());
397 let adapter = Base64BlobStoreAdapter::new(binary);
398 let blob = adapter
399 .put_image("image/png", "YWJj")
400 .await
401 .expect("put base64");
402 let payload = adapter.get(&blob.blob_id).await.expect("get base64");
403 assert_eq!(payload.media_type, "image/png");
404 assert_eq!(payload.data, "YWJj");
405 }
406
407 #[tokio::test]
408 async fn local_store_reads_legacy_fs_blob_layout() {
409 let temp = tempfile::tempdir().expect("tempdir");
410 let store = ObjectStoreBlobStore::local(temp.path().to_path_buf()).expect("local store");
411 let legacy_id = compute_legacy_blob_id("image/png", "YWJj");
412 let key = storage_key(&legacy_id);
413 let legacy_dir = temp.path().join(&key[0..2]);
414 tokio::fs::create_dir_all(&legacy_dir)
415 .await
416 .expect("legacy dir");
417 tokio::fs::write(
418 legacy_dir.join(format!("{key}.json")),
419 serde_json::to_vec(&LegacyStoredBlob {
420 media_type: "image/png".to_string(),
421 data: "YWJj".to_string(),
422 })
423 .expect("legacy json"),
424 )
425 .await
426 .expect("legacy write");
427
428 let payload = store.get_bytes(&legacy_id).await.expect("legacy fallback");
429 assert_eq!(payload.blob_id, legacy_id);
430 assert_eq!(payload.media_type, "image/png");
431 assert_eq!(payload.data, Bytes::from_static(b"abc"));
432 }
433
434 #[test]
435 fn local_store_creates_missing_root() {
436 let temp = tempfile::tempdir().expect("tempdir");
437 let root = temp.path().join("missing").join("blobs");
438
439 let store = ObjectStoreBlobStore::local(root.clone()).expect("local store opens");
440
441 assert!(root.is_dir());
442 assert!(store.is_persistent());
443 }
444
445 #[tokio::test]
446 async fn invalid_blob_ids_do_not_address_object_paths() {
447 let store = ObjectStoreBlobStore::memory();
448 let invalid = BlobId::from("../not-a-sha");
449
450 let err = store.get_bytes(&invalid).await.expect_err("invalid id");
451 assert!(matches!(err, BlobStoreError::NotFound(id) if id == invalid));
452 let err = store.delete(&invalid).await.expect_err("invalid delete");
453 assert!(matches!(err, BlobStoreError::NotFound(id) if id == invalid));
454 }
455
456 fn compute_legacy_blob_id(media_type: &str, data: &str) -> BlobId {
457 let mut hasher = Sha256::new();
458 hasher.update(media_type.as_bytes());
459 hasher.update([0]);
460 hasher.update(data.as_bytes());
461 BlobId::new(format!("sha256:{:x}", hasher.finalize()))
462 }
463}