Skip to main content

meerkat_mobkit/
blob_store.rs

1//! MobKit-owned binary blob storage.
2//!
3//! Published Meerkat 0.6 exposes a `BlobStore` trait whose payload is base64
4//! text. MobKit keeps that compatibility seam for the runtime while storing
5//! local blobs as raw bytes for HTTP serving and multipart upload paths.
6
7use 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}