feat(at-blob): InMemoryBlobStore echt implementieren
Der Store war ein Platzhalter aus vier unimplemented!() — jeder Codepfad, der ihn statt S3BlobStore benutzt hätte, wäre gepaniced. Backing store ist eine RwLock<HashMap<key, (Bytes, mime)>>, damit der Typ Send + Sync bleibt und hinter Arc<dyn BlobStore> funktioniert. put() berechnet die CID identisch zu s3.rs (sha256 → cid_for_raw(0x55, hash)), get() liefert None statt Fehler, delete() ist idempotent, public_url() zeigt auf die vorhandene PDS-Route /blob/:cid. 5 Tests: Roundtrip, fehlender Key, delete-dann-get, delete auf Unbekanntes, gleiche Bytes → gleiche CID. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_013HC9HLrUU1LNwkzp8nkDLX
This commit is contained in:
co-authored by
Claude Opus 5
parent
3c6f4dd67c
commit
9d009bfcba
+108
-9
@@ -1,7 +1,10 @@
|
||||
use anyhow::Result;
|
||||
use async_trait::async_trait;
|
||||
use at_crypto::cid::{cid_for_raw, sha256};
|
||||
use bytes::Bytes;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::collections::HashMap;
|
||||
use tokio::sync::RwLock;
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct BlobInfo {
|
||||
@@ -19,20 +22,116 @@ pub trait BlobStore: Send + Sync {
|
||||
async fn public_url(&self, key: &str) -> Result<String>;
|
||||
}
|
||||
|
||||
pub struct InMemoryBlobStore;
|
||||
/// In-process blob store backed by a `HashMap` guarded by a
|
||||
/// [`tokio::sync::RwLock`]. Used in tests and any dev/single-node
|
||||
/// deployment that doesn't need a real object store.
|
||||
///
|
||||
/// Behaves like [`crate::s3::S3BlobStore`] with respect to `BlobInfo`
|
||||
/// (same CID computation over the raw bytes, same field semantics) so
|
||||
/// tests can swap one for the other transparently. Data does not
|
||||
/// survive process restarts.
|
||||
#[derive(Default)]
|
||||
pub struct InMemoryBlobStore {
|
||||
blobs: RwLock<HashMap<String, (Bytes, String /* mime */)>>,
|
||||
}
|
||||
|
||||
impl InMemoryBlobStore {
|
||||
pub fn new() -> Self {
|
||||
Self::default()
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl BlobStore for InMemoryBlobStore {
|
||||
async fn put(&self, _key: &str, _data: Bytes, _mime: &str) -> Result<BlobInfo> {
|
||||
unimplemented!("in-memory blob store placeholder")
|
||||
async fn put(&self, key: &str, data: Bytes, mime: &str) -> Result<BlobInfo> {
|
||||
let hash = sha256(&data);
|
||||
let cid = cid_for_raw(0x55, hash)?;
|
||||
let size = data.len() as u64;
|
||||
self.blobs
|
||||
.write()
|
||||
.await
|
||||
.insert(key.to_string(), (data, mime.to_string()));
|
||||
Ok(BlobInfo {
|
||||
cid: cid.to_string(),
|
||||
mime_type: mime.to_string(),
|
||||
size,
|
||||
storage_key: key.to_string(),
|
||||
})
|
||||
}
|
||||
async fn get(&self, _key: &str) -> Result<Option<Bytes>> {
|
||||
unimplemented!()
|
||||
|
||||
async fn get(&self, key: &str) -> Result<Option<Bytes>> {
|
||||
Ok(self
|
||||
.blobs
|
||||
.read()
|
||||
.await
|
||||
.get(key)
|
||||
.map(|(data, _mime)| data.clone()))
|
||||
}
|
||||
async fn delete(&self, _key: &str) -> Result<()> {
|
||||
unimplemented!()
|
||||
|
||||
async fn delete(&self, key: &str) -> Result<()> {
|
||||
self.blobs.write().await.remove(key);
|
||||
Ok(())
|
||||
}
|
||||
async fn public_url(&self, _key: &str) -> Result<String> {
|
||||
unimplemented!()
|
||||
|
||||
async fn public_url(&self, key: &str) -> Result<String> {
|
||||
Ok(format!("/blob/{key}"))
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[tokio::test]
|
||||
async fn put_get_roundtrip() {
|
||||
let store = InMemoryBlobStore::new();
|
||||
let data = Bytes::from_static(b"hello world");
|
||||
let info = store.put("k1", data.clone(), "text/plain").await.unwrap();
|
||||
assert_eq!(info.storage_key, "k1");
|
||||
assert_eq!(info.mime_type, "text/plain");
|
||||
assert_eq!(info.size, data.len() as u64);
|
||||
|
||||
let got = store.get("k1").await.unwrap();
|
||||
assert_eq!(got, Some(data));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn get_missing_key_returns_none() {
|
||||
let store = InMemoryBlobStore::new();
|
||||
let got = store.get("does-not-exist").await.unwrap();
|
||||
assert_eq!(got, None);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn delete_then_get_returns_none() {
|
||||
let store = InMemoryBlobStore::new();
|
||||
store
|
||||
.put("k2", Bytes::from_static(b"data"), "application/octet-stream")
|
||||
.await
|
||||
.unwrap();
|
||||
store.delete("k2").await.unwrap();
|
||||
let got = store.get("k2").await.unwrap();
|
||||
assert_eq!(got, None);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn delete_missing_key_is_ok() {
|
||||
let store = InMemoryBlobStore::new();
|
||||
store.delete("never-existed").await.unwrap();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn same_bytes_produce_same_cid() {
|
||||
let store = InMemoryBlobStore::new();
|
||||
let data = Bytes::from_static(b"identical payload");
|
||||
let info_a = store
|
||||
.put("key-a", data.clone(), "application/octet-stream")
|
||||
.await
|
||||
.unwrap();
|
||||
let info_b = store
|
||||
.put("key-b", data.clone(), "application/octet-stream")
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(info_a.cid, info_b.cid);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user