1use std::sync::Arc;
9
10use async_trait::async_trait;
11use corium_core::{Datom, IndexOrder, encoding::DecodeError};
12use corium_db::Db;
13use corium_protocol::codec::{self, CodecError};
14use corium_store::{
15 BlobId, BlobStore, DbRoot, FORMAT_VERSION, RootStore, SegmentCache, SegmentCacheConfig,
16 SegmentCacheMetrics, SegmentReader, StoreError, db_root_name, decode_index_manifest,
17 decode_segment_keys, is_index_manifest, meta_root_name,
18};
19use thiserror::Error;
20
21#[async_trait]
26pub trait PeerStorage: Send + Sync {
27 async fn get_blob(&self, id: &BlobId) -> Result<Option<Vec<u8>>, StoreError>;
29 async fn get_peer_root(&self, name: &str) -> Result<Option<Vec<u8>>, StoreError>;
31}
32
33pub struct CachedPeerStorage {
35 storage: Arc<dyn PeerStorage>,
36 cache: SegmentCache,
37}
38
39impl CachedPeerStorage {
40 pub fn open(
46 storage: Arc<dyn PeerStorage>,
47 config: &SegmentCacheConfig,
48 ) -> std::io::Result<Self> {
49 Ok(Self {
50 storage,
51 cache: SegmentCache::open(config)?,
52 })
53 }
54
55 #[must_use]
57 pub fn metrics(&self) -> Arc<SegmentCacheMetrics> {
58 self.cache.metrics()
59 }
60}
61
62#[async_trait]
63impl SegmentReader for CachedPeerStorage {
64 async fn read_segment(&self, id: &BlobId) -> Result<Option<Vec<u8>>, StoreError> {
65 self.storage.get_blob(id).await
66 }
67}
68
69#[async_trait]
70impl PeerStorage for CachedPeerStorage {
71 async fn get_blob(&self, id: &BlobId) -> Result<Option<Vec<u8>>, StoreError> {
72 Ok(self
73 .cache
74 .get_or_load(self, id)
75 .await?
76 .map(|bytes| bytes.to_vec()))
77 }
78 async fn get_peer_root(&self, name: &str) -> Result<Option<Vec<u8>>, StoreError> {
79 self.storage.get_peer_root(name).await
80 }
81}
82
83#[async_trait]
84impl<S> PeerStorage for S
85where
86 S: BlobStore + RootStore + Send + Sync,
87{
88 async fn get_blob(&self, id: &BlobId) -> Result<Option<Vec<u8>>, StoreError> {
89 BlobStore::get(self, id).await
90 }
91
92 async fn get_peer_root(&self, name: &str) -> Result<Option<Vec<u8>>, StoreError> {
93 RootStore::get_root(self, name).await
94 }
95}
96
97#[derive(Debug, Error)]
99pub enum SnapshotError {
100 #[error(transparent)]
102 Store(#[from] StoreError),
103 #[error(transparent)]
105 Codec(#[from] CodecError),
106 #[error(transparent)]
108 Key(#[from] DecodeError),
109 #[error("malformed published root for database {0:?}")]
111 MalformedRoot(String),
112 #[error("storage format {found} is newer than supported format {supported}")]
114 UnsupportedFormat {
115 found: u32,
117 supported: u32,
119 },
120 #[error("published snapshot for database {0:?} has no metadata root")]
122 MissingMetadata(String),
123}
124
125pub async fn load_current_snapshot(
134 store: &dyn PeerStorage,
135 db: &str,
136) -> Result<Option<Db>, SnapshotError> {
137 let Some(root_bytes) = store.get_peer_root(&db_root_name(db)).await? else {
138 return Ok(None);
139 };
140 let root =
141 DbRoot::decode(&root_bytes).ok_or_else(|| SnapshotError::MalformedRoot(db.into()))?;
142 if root.format_version > FORMAT_VERSION {
143 return Err(SnapshotError::UnsupportedFormat {
144 found: root.format_version,
145 supported: FORMAT_VERSION,
146 });
147 }
148 let Some(roots) = root.roots else {
149 return Ok(None);
150 };
151 let Some(metadata) = store.get_peer_root(&meta_root_name(db)).await? else {
152 return Err(SnapshotError::MissingMetadata(db.into()));
153 };
154 let (schema, idents, interner) = codec::decode_metadata(&metadata)?;
155 let datoms = load_index_keys(store, &roots[0])
156 .await?
157 .into_iter()
158 .map(|key| Datom::from_key(IndexOrder::Eavt, &key))
159 .collect::<Result<Vec<_>, _>>()?;
160 Ok(Some(Db::from_current_snapshot(
161 root.index_basis_t,
162 schema,
163 idents,
164 interner,
165 datoms,
166 )))
167}
168
169async fn load_index_keys(store: &dyn PeerStorage, id: &BlobId) -> Result<Vec<Vec<u8>>, StoreError> {
172 let blob = store
173 .get_blob(id)
174 .await?
175 .ok_or_else(|| StoreError::MissingBlob(id.clone()))?;
176 if !is_index_manifest(&blob) {
177 return decode_segment_keys(&blob);
178 }
179 let mut keys = Vec::new();
180 for child in decode_index_manifest(&blob)? {
181 let chunk = store
182 .get_blob(&child)
183 .await?
184 .ok_or_else(|| StoreError::MissingBlob(child.clone()))?;
185 keys.extend(decode_segment_keys(&chunk)?);
186 }
187 Ok(keys)
188}
189
190pub struct SegmentSource<S> {
192 store: Arc<S>,
193 cache: SegmentCache,
194}
195
196impl<S: BlobStore + RootStore> SegmentSource<S> {
197 #[must_use]
199 pub fn new(store: Arc<S>) -> Self {
200 Self {
201 store,
202 cache: SegmentCache::default(),
203 }
204 }
205
206 pub async fn index_root(&self, db: &str) -> Result<Option<DbRoot>, StoreError> {
211 Ok(self
212 .store
213 .get_root(&db_root_name(db))
214 .await?
215 .as_deref()
216 .and_then(DbRoot::decode))
217 }
218
219 pub async fn lease_holder_endpoint(&self, db: &str) -> Result<Option<String>, StoreError> {
227 Ok(self
228 .index_root(db)
229 .await?
230 .and_then(|root| (!root.owner_endpoint.is_empty()).then_some(root.owner_endpoint)))
231 }
232
233 pub async fn segment(
240 &self,
241 root: &DbRoot,
242 order: IndexOrder,
243 ) -> Result<Option<Arc<[u8]>>, StoreError> {
244 let Some(roots) = &root.roots else {
245 return Ok(None);
246 };
247 let slot = match order {
248 IndexOrder::Eavt => 0,
249 IndexOrder::Aevt => 1,
250 IndexOrder::Avet => 2,
251 IndexOrder::Vaet => 3,
252 };
253 let Some(blob) = self
254 .cache
255 .get_or_load(self.store.as_ref(), &roots[slot])
256 .await?
257 else {
258 return Ok(None);
259 };
260 if !is_index_manifest(&blob) {
261 return Ok(Some(blob));
262 }
263 let mut bytes = Vec::new();
264 for child in decode_index_manifest(&blob)? {
265 let chunk = self
266 .cache
267 .get_or_load(self.store.as_ref(), &child)
268 .await?
269 .ok_or_else(|| StoreError::MissingBlob(child.clone()))?;
270 bytes.extend_from_slice(&chunk);
271 }
272 Ok(Some(bytes.into()))
273 }
274
275 pub fn segment_keys(bytes: &[u8]) -> Result<Vec<Vec<u8>>, StoreError> {
281 decode_segment_keys(bytes)
282 }
283}
284
285#[cfg(test)]
286mod tests {
287 use corium_core::{EntityId, KeywordInterner, Value};
288 use corium_db::Idents;
289 use corium_store::{MemoryStore, RootStore};
290
291 use super::*;
292
293 #[tokio::test]
294 async fn loads_published_eavt_snapshot() {
295 let store = MemoryStore::default();
296 let datom = Datom {
297 e: EntityId::from_raw(1_001),
298 a: EntityId::from_raw(101),
299 v: Value::Str("snapshot".into()),
300 tx: EntityId::from_raw(37),
301 added: true,
302 };
303 let key = datom.key(IndexOrder::Eavt);
304 let mut segment = Vec::new();
305 segment.extend_from_slice(&(key.len() as u64).to_be_bytes());
306 segment.extend_from_slice(&key);
307 let id = store.put(&segment).await.expect("put segment");
308 let root = DbRoot {
309 format_version: FORMAT_VERSION,
310 lease_version: 1,
311 owner: "test".into(),
312 lease_expires_unix_ms: 0,
313 owner_endpoint: String::new(),
314 index_basis_t: 37,
315 roots: Some([id.clone(), id.clone(), id.clone(), id]),
316 next_entity_id: 1_005,
317 last_tx_instant: 0,
318 key_manifest_version: 0,
319 };
320 RootStore::cas_root(&store, &db_root_name("music"), None, &root.encode())
321 .await
322 .expect("put root");
323 let metadata = codec::encode_metadata(
324 &corium_core::Schema::default(),
325 &Idents::default(),
326 &KeywordInterner::default(),
327 );
328 RootStore::cas_root(&store, &meta_root_name("music"), None, &metadata)
329 .await
330 .expect("put metadata");
331
332 let db = load_current_snapshot(&store, "music")
333 .await
334 .expect("load snapshot")
335 .expect("published snapshot");
336 assert_eq!(db.basis_t(), 37);
337 assert_eq!(db.datoms(), vec![datom]);
338 }
339
340 #[tokio::test]
341 async fn loads_chunked_manifest_snapshot() {
342 let store = MemoryStore::default();
343 let datoms: Vec<Datom> = (0..4u64)
344 .map(|n| Datom {
345 e: EntityId::from_raw(1_001 + n),
346 a: EntityId::from_raw(101),
347 v: Value::Long(i64::try_from(n).unwrap()),
348 tx: EntityId::from_raw(37),
349 added: true,
350 })
351 .collect();
352 let mut chunk_ids = Vec::new();
354 for pair in datoms.chunks(2) {
355 let mut chunk = Vec::new();
356 for datom in pair {
357 let key = datom.key(IndexOrder::Eavt);
358 chunk.extend_from_slice(&(key.len() as u64).to_be_bytes());
359 chunk.extend_from_slice(&key);
360 }
361 chunk_ids.push(store.put(&chunk).await.expect("put chunk"));
362 }
363 let manifest = corium_store::encode_index_manifest(&chunk_ids);
364 let id = store.put(&manifest).await.expect("put manifest");
365 let root = DbRoot {
366 format_version: FORMAT_VERSION,
367 lease_version: 1,
368 owner: "test".into(),
369 lease_expires_unix_ms: 0,
370 owner_endpoint: String::new(),
371 index_basis_t: 37,
372 roots: Some([id.clone(), id.clone(), id.clone(), id]),
373 next_entity_id: 1_005,
374 last_tx_instant: 0,
375 key_manifest_version: 0,
376 };
377 RootStore::cas_root(&store, &db_root_name("music"), None, &root.encode())
378 .await
379 .expect("put root");
380 let metadata = codec::encode_metadata(
381 &corium_core::Schema::default(),
382 &Idents::default(),
383 &KeywordInterner::default(),
384 );
385 RootStore::cas_root(&store, &meta_root_name("music"), None, &metadata)
386 .await
387 .expect("put metadata");
388
389 let db = load_current_snapshot(&store, "music")
390 .await
391 .expect("load snapshot")
392 .expect("published snapshot");
393 assert_eq!(db.basis_t(), 37);
394 assert_eq!(db.datoms(), datoms);
395 }
396
397 #[tokio::test]
398 async fn absent_publication_falls_back_to_log_replay() {
399 assert!(
400 load_current_snapshot(&MemoryStore::default(), "music")
401 .await
402 .expect("load snapshot")
403 .is_none()
404 );
405 }
406}