Skip to main content

mkit_server/memory/
blob.rs

1//! `MemoryBlobStore`: the reference [`BlobStore`].
2
3use std::collections::BTreeMap;
4use std::pin::Pin;
5use std::sync::atomic::{AtomicU64, Ordering};
6use std::sync::{Arc, Mutex};
7use std::task::{Context, Poll};
8
9use bytes::Bytes;
10use futures_core::Stream;
11use mkit_core::hash::{Hash, Hasher};
12use mkit_core::upload_parts::{PartHasher, PartPlan, merge_to_root};
13
14use super::{MemoryFault, lock, take_fault};
15use crate::store::{
16    BlobBody, BlobKey, BlobMeta, BlobStore, ByteRange, CommitOutcome, MAX_BLOB_PIECE_BYTES,
17    MultipartBlobStore, PackSink, PartRef, PartSink, StoreError,
18};
19
20/// Bodies longer than this are streamed in pieces of this size
21/// (`BlobStore::get`).
22const STREAM_CHUNK: usize = MAX_BLOB_PIECE_BYTES;
23
24/// `bytes` as a body: whole, or streamed when longer than [`STREAM_CHUNK`].
25fn body(bytes: Bytes) -> BlobBody {
26    if bytes.len() <= STREAM_CHUNK {
27        return BlobBody::Bytes(bytes);
28    }
29    BlobBody::Stream {
30        len: bytes.len() as u64,
31        stream: Box::pin(Chunks(bytes)),
32    }
33}
34
35/// The remaining bytes of a streamed body.
36struct Chunks(Bytes);
37
38impl Stream for Chunks {
39    type Item = Result<Bytes, StoreError>;
40
41    fn poll_next(self: Pin<&mut Self>, _: &mut Context<'_>) -> Poll<Option<Self::Item>> {
42        let rest = &mut self.get_mut().0;
43        if rest.is_empty() {
44            return Poll::Ready(None);
45        }
46        let n = rest.len().min(STREAM_CHUNK);
47        Poll::Ready(Some(Ok(rest.split_to(n))))
48    }
49}
50
51#[derive(Debug, Default)]
52struct Shared {
53    #[cfg(test)]
54    read_calls: AtomicU64,
55    blobs: Mutex<BTreeMap<BlobKey, Bytes>>,
56    sessions: Mutex<BTreeMap<Vec<u8>, MemoryMultipart>>,
57    next_session: AtomicU64,
58    fault: Mutex<Option<MemoryFault>>,
59}
60
61#[derive(Debug)]
62struct MemoryMultipart {
63    key: BlobKey,
64    len: u64,
65    part_size: u64,
66    parts: BTreeMap<u32, MemoryPart>,
67}
68
69#[derive(Debug)]
70struct MemoryPart {
71    bytes: Bytes,
72    cv: [u8; 32],
73}
74
75/// An in-memory [`BlobStore`] for one keyspace. Clones share the blobs.
76#[derive(Debug, Clone)]
77pub struct MemoryBlobStore {
78    keyspace: String,
79    shared: Arc<Shared>,
80    single_put_limit: Option<u64>,
81}
82
83/// The `packs` keyspace.
84impl Default for MemoryBlobStore {
85    fn default() -> Self {
86        Self::new("packs")
87    }
88}
89
90impl MemoryBlobStore {
91    #[cfg(test)]
92    pub(crate) fn read_calls(&self) -> u64 {
93        self.shared.read_calls.load(Ordering::SeqCst)
94    }
95    /// Open multipart sessions, observable only for test-faults conformance.
96    #[cfg(any(test, feature = "test-faults"))]
97    #[doc(hidden)]
98    #[must_use]
99    pub fn multipart_session_count(&self) -> usize {
100        lock(&self.shared.sessions).len()
101    }
102
103    /// An empty store for `keyspace` (`packs` for pack uploads).
104    #[must_use]
105    pub fn new(keyspace: impl Into<String>) -> Self {
106        Self {
107            keyspace: keyspace.into(),
108            shared: Arc::default(),
109            single_put_limit: None,
110        }
111    }
112
113    /// The keyspace this store serves.
114    #[must_use]
115    pub fn keyspace(&self) -> &str {
116        &self.keyspace
117    }
118
119    /// Cap what one `begin` upload may carry, as a bounded backend does, so
120    /// larger objects take the multipart path. For tests.
121    #[must_use]
122    pub fn with_single_put_limit(mut self, limit: u64) -> Self {
123        self.single_put_limit = Some(limit);
124        self
125    }
126
127    /// Arm a one-shot [`MemoryFault`] (`BlobWrite` or `BlobCommit`).
128    #[must_use]
129    pub fn with_fault(self, fault: MemoryFault) -> Self {
130        *lock(&self.shared.fault) = Some(fault);
131        self
132    }
133}
134
135impl MemoryBlobStore {
136    /// Complete a multipart upload against the key's hash (`None`) or an
137    /// object's content root.
138    fn complete_with(
139        &self,
140        key: BlobKey,
141        session: &[u8],
142        plan: &PartPlan,
143        parts: &[PartRef],
144        root: Option<Hash>,
145    ) -> Result<CommitOutcome, StoreError> {
146        let expected = key.expected_root(root)?;
147        let mut sessions = lock(&self.shared.sessions);
148        let upload = sessions.get(session).ok_or(StoreError::SessionGone)?;
149        if upload.key != key || upload.len != plan.total() || upload.part_size != plan.part_size() {
150            return Err(StoreError::SessionGone);
151        }
152        if parts.len() != plan.count() as usize {
153            return Err(StoreError::Invalid("wrong number of parts".into()));
154        }
155        let mut cvs = Vec::with_capacity(parts.len());
156        let mut bytes = Vec::new();
157        for (i, part) in parts.iter().enumerate() {
158            let index = u32::try_from(i).map_err(|_| StoreError::Invalid("part index".into()))?;
159            let stored = upload.parts.get(&index).ok_or(StoreError::SessionGone)?;
160            if part.index != index
161                || part.len
162                    != plan
163                        .expected_len(index)
164                        .map_err(|e| StoreError::Invalid(e.to_string().into()))?
165                || part.tag.as_slice() != stored.cv
166                || part.len != stored.bytes.len() as u64
167            {
168                return Err(StoreError::Invalid(
169                    "part reference does not match stored part".into(),
170                ));
171            }
172            cvs.push(stored.cv);
173            bytes.extend_from_slice(&stored.bytes);
174        }
175        if bytes.len() as u64 != plan.total()
176            || merge_to_root(plan, &cvs).map_err(|e| StoreError::Invalid(e.to_string().into()))?
177                != expected
178        {
179            return Err(StoreError::Invalid(
180                "merged part root does not match key".into(),
181            ));
182        }
183        let mut blobs = lock(&self.shared.blobs);
184        let outcome = if let std::collections::btree_map::Entry::Vacant(entry) = blobs.entry(key) {
185            entry.insert(Bytes::from(bytes));
186            CommitOutcome::Created
187        } else {
188            CommitOutcome::AlreadyPresent
189        };
190        sessions.remove(session);
191        Ok(outcome)
192    }
193}
194
195/// The upload handle of [`MemoryBlobStore`]. It hashes incrementally and
196/// buffers the blob: the one exception to `PackSink`'s one-part memory
197/// bound, since this backend holds every blob in memory anyway.
198#[derive(Debug)]
199pub struct MemoryPackSink {
200    shared: Arc<Shared>,
201    key: BlobKey,
202    len: u64,
203    hasher: Hasher,
204    buf: Vec<u8>,
205    writes: u32,
206}
207
208/// One attempted memory part. Its bytes are staged until its CV verifies.
209#[derive(Debug)]
210pub struct MemoryPartSink {
211    shared: Arc<Shared>,
212    key: BlobKey,
213    session: Vec<u8>,
214    index: u32,
215    expected_cv: [u8; 32],
216    hasher: PartHasher,
217    bytes: Vec<u8>,
218}
219
220impl MultipartBlobStore for MemoryBlobStore {
221    type PartSink = MemoryPartSink;
222    const MAX_PARTS: u32 = u32::MAX;
223
224    fn supports_multipart(&self) -> bool {
225        true
226    }
227
228    async fn begin_multipart(
229        &self,
230        key: BlobKey,
231        len: u64,
232        part_size: u64,
233    ) -> Result<Vec<u8>, StoreError> {
234        PartPlan::new(len, part_size, Self::MAX_PARTS)
235            .map_err(|e| StoreError::Invalid(e.to_string().into()))?;
236        let session = self
237            .shared
238            .next_session
239            .fetch_add(1, Ordering::Relaxed)
240            .to_be_bytes()
241            .to_vec();
242        lock(&self.shared.sessions).insert(
243            session.clone(),
244            MemoryMultipart {
245                key,
246                len,
247                part_size,
248                parts: BTreeMap::new(),
249            },
250        );
251        Ok(session)
252    }
253
254    async fn begin_part(
255        &self,
256        key: BlobKey,
257        session: &[u8],
258        plan: &PartPlan,
259        index: u32,
260        expected_cv: [u8; 32],
261    ) -> Result<MemoryPartSink, StoreError> {
262        let sessions = lock(&self.shared.sessions);
263        let upload = sessions.get(session).ok_or(StoreError::SessionGone)?;
264        if upload.key != key || upload.len != plan.total() || upload.part_size != plan.part_size() {
265            return Err(StoreError::SessionGone);
266        }
267        let hasher =
268            PartHasher::new(plan, index).map_err(|e| StoreError::Invalid(e.to_string().into()))?;
269        Ok(MemoryPartSink {
270            shared: Arc::clone(&self.shared),
271            key,
272            session: session.to_vec(),
273            index,
274            expected_cv,
275            hasher,
276            bytes: Vec::new(),
277        })
278    }
279
280    async fn complete(
281        &self,
282        key: BlobKey,
283        session: &[u8],
284        plan: &PartPlan,
285        parts: &[PartRef],
286    ) -> Result<CommitOutcome, StoreError> {
287        self.complete_with(key, session, plan, parts, None)
288    }
289
290    async fn complete_with_root(
291        &self,
292        key: BlobKey,
293        session: &[u8],
294        plan: &PartPlan,
295        parts: &[PartRef],
296        content_root: Hash,
297    ) -> Result<CommitOutcome, StoreError> {
298        self.complete_with(key, session, plan, parts, Some(content_root))
299    }
300
301    fn single_put_limit(&self) -> Option<u64> {
302        self.single_put_limit
303    }
304
305    async fn abort(&self, key: BlobKey, session: &[u8]) -> Result<(), StoreError> {
306        let mut sessions = lock(&self.shared.sessions);
307        if sessions
308            .get(session)
309            .is_some_and(|upload| upload.key == key)
310        {
311            sessions.remove(session);
312        }
313        Ok(())
314    }
315}
316
317impl PartSink for MemoryPartSink {
318    async fn write(&mut self, chunk: Bytes) -> Result<(), StoreError> {
319        if chunk.is_empty() {
320            return Err(StoreError::Invalid("empty part chunk".into()));
321        }
322        self.hasher
323            .update(&chunk)
324            .map_err(|e| StoreError::Invalid(e.to_string().into()))?;
325        self.bytes.extend_from_slice(&chunk);
326        Ok(())
327    }
328
329    async fn commit(self) -> Result<Vec<u8>, StoreError> {
330        let cv = self
331            .hasher
332            .finalize()
333            .map_err(|e| StoreError::Invalid(e.to_string().into()))?;
334        if cv != self.expected_cv {
335            return Err(StoreError::PartSubtreeMismatch);
336        }
337        let mut sessions = lock(&self.shared.sessions);
338        let upload = sessions
339            .get_mut(&self.session)
340            .ok_or(StoreError::SessionGone)?;
341        if upload.key != self.key {
342            return Err(StoreError::SessionGone);
343        }
344        upload.parts.insert(
345            self.index,
346            MemoryPart {
347                bytes: Bytes::from(self.bytes),
348                cv,
349            },
350        );
351        Ok(cv.to_vec())
352    }
353
354    async fn abort(self) {}
355}
356
357impl BlobStore for MemoryBlobStore {
358    type Sink = MemoryPackSink;
359
360    async fn begin(&self, key: BlobKey, len: u64) -> Result<MemoryPackSink, StoreError> {
361        Ok(MemoryPackSink {
362            shared: self.shared.clone(),
363            key,
364            len,
365            hasher: Hasher::new(),
366            buf: Vec::new(),
367            writes: 0,
368        })
369    }
370
371    async fn get(
372        &self,
373        key: &BlobKey,
374        range: Option<ByteRange>,
375    ) -> Result<Option<BlobBody>, StoreError> {
376        #[cfg(test)]
377        self.shared.read_calls.fetch_add(1, Ordering::SeqCst);
378        let Some(blob) = lock(&self.shared.blobs).get(key).cloned() else {
379            return Ok(None);
380        };
381        let Some(range) = range else {
382            return Ok(Some(body(blob)));
383        };
384        let span = range.resolve(blob.len() as u64)?;
385        // `resolve` bounds the span by the blob's length, which is a usize.
386        let index = |n: u64| usize::try_from(n).map_err(|_| StoreError::Invalid("range".into()));
387        Ok(Some(body(blob.slice(index(span.start)?..index(span.end)?))))
388    }
389
390    async fn head(&self, key: &BlobKey) -> Result<Option<BlobMeta>, StoreError> {
391        #[cfg(test)]
392        self.shared.read_calls.fetch_add(1, Ordering::SeqCst);
393        let blobs = lock(&self.shared.blobs);
394        Ok(blobs.get(key).map(|b| BlobMeta {
395            len: b.len() as u64,
396        }))
397    }
398
399    async fn probe(&self) -> Result<(), StoreError> {
400        Ok(())
401    }
402
403    async fn delete(&self, key: &BlobKey) -> Result<bool, StoreError> {
404        Ok(lock(&self.shared.blobs).remove(key).is_some())
405    }
406}
407
408impl MemoryPackSink {
409    /// Verify against `root` (or the key, for `None`) and publish.
410    fn finish(self, root: Option<Hash>) -> Result<CommitOutcome, StoreError> {
411        take_fault(&self.shared.fault, MemoryFault::BlobCommit)?;
412        let expected = self.key.expected_root(root)?;
413        if self.buf.len() as u64 != self.len {
414            return Err(StoreError::Invalid("blob length does not match".into()));
415        }
416        if self.hasher.finalize() != expected {
417            return Err(StoreError::Invalid(
418                "blob hash does not match its key".into(),
419            ));
420        }
421        let mut blobs = lock(&self.shared.blobs);
422        if blobs.contains_key(&self.key) {
423            return Ok(CommitOutcome::AlreadyPresent);
424        }
425        blobs.insert(self.key, Bytes::from(self.buf));
426        Ok(CommitOutcome::Created)
427    }
428}
429
430impl PackSink for MemoryPackSink {
431    async fn write(&mut self, chunk: Bytes) -> Result<(), StoreError> {
432        let nth = self.writes;
433        self.writes = self.writes.saturating_add(1);
434        take_fault(&self.shared.fault, MemoryFault::BlobWrite(nth))?;
435        if (self.buf.len() + chunk.len()) as u64 > self.len {
436            return Err(StoreError::Invalid("blob is longer than declared".into()));
437        }
438        self.hasher.update(&chunk);
439        self.buf.extend_from_slice(&chunk);
440        Ok(())
441    }
442
443    async fn commit(self) -> Result<CommitOutcome, StoreError> {
444        self.finish(None)
445    }
446
447    async fn commit_with_root(self, content_root: Hash) -> Result<CommitOutcome, StoreError> {
448        self.finish(Some(content_root))
449    }
450
451    async fn abort(self) {}
452}
453
454#[cfg(test)]
455mod tests {
456    use futures_executor::block_on;
457    use mkit_core::hash::hash;
458    use mkit_core::upload_parts::{MIN_PART_SIZE, PartPlan, part_subtree_cv};
459
460    use super::*;
461
462    fn key_of(bytes: &[u8]) -> BlobKey {
463        BlobKey::pack(hash(bytes))
464    }
465
466    fn put(
467        store: &MemoryBlobStore,
468        key: BlobKey,
469        len: u64,
470        chunks: &[&[u8]],
471    ) -> Result<CommitOutcome, StoreError> {
472        block_on(async {
473            let mut sink = store.begin(key, len).await?;
474            for chunk in chunks {
475                sink.write(Bytes::copy_from_slice(chunk)).await?;
476            }
477            sink.commit().await
478        })
479    }
480
481    fn get(store: &MemoryBlobStore, key: &BlobKey, range: Option<ByteRange>) -> Option<Bytes> {
482        match block_on(store.get(key, range)).unwrap()? {
483            BlobBody::Bytes(b) => Some(b),
484            BlobBody::Stream { .. } => panic!("a small memory blob is never streamed"),
485        }
486    }
487
488    #[test]
489    fn blob_commit_rejects_hash_and_len_mismatch_leaving_nothing() {
490        let store = MemoryBlobStore::default();
491        let key = key_of(b"hello");
492        for (k, len, chunks) in [
493            (key_of(b"other"), 5, &[&b"hello"[..]][..]),
494            (key, 6, &[&b"hello"[..]][..]),
495            (key, 4, &[&b"hel"[..], b"lo"][..]),
496        ] {
497            assert!(matches!(
498                put(&store, k, len, chunks),
499                Err(StoreError::Invalid(_))
500            ));
501        }
502        assert_eq!(block_on(store.head(&key)).unwrap(), None);
503        assert_eq!(get(&store, &key_of(b"other"), None), None);
504        // Aborting leaves nothing either.
505        block_on(async {
506            let mut sink = store.begin(key, 5).await.unwrap();
507            sink.write(Bytes::from_static(b"hello")).await.unwrap();
508            sink.abort().await;
509        });
510        assert_eq!(block_on(store.head(&key)).unwrap(), None);
511    }
512
513    #[test]
514    fn blob_commit_identical_bytes_is_already_present() {
515        let store = MemoryBlobStore::default();
516        let key = key_of(b"hello");
517        assert_eq!(
518            put(&store, key, 5, &[b"he", b"llo"]).unwrap(),
519            CommitOutcome::Created
520        );
521        assert_eq!(
522            put(&store, key, 5, &[b"hello"]).unwrap(),
523            CommitOutcome::AlreadyPresent
524        );
525        assert_eq!(
526            block_on(store.head(&key)).unwrap(),
527            Some(BlobMeta { len: 5 })
528        );
529        assert_eq!(store.keyspace(), "packs");
530        let empty = key_of(b"");
531        assert_eq!(put(&store, empty, 0, &[]).unwrap(), CommitOutcome::Created);
532        assert_eq!(get(&store, &empty, None).unwrap().len(), 0);
533    }
534
535    #[test]
536    fn blob_get_range() {
537        let store = MemoryBlobStore::default();
538        let key = key_of(b"0123456789");
539        put(&store, key, 10, &[b"0123456789"]).unwrap();
540        let range = |start, end_inclusive| {
541            Some(ByteRange {
542                start,
543                end_inclusive,
544            })
545        };
546        assert_eq!(get(&store, &key, range(0, 0)).unwrap(), "0");
547        assert_eq!(get(&store, &key, range(3, 5)).unwrap(), "345");
548        assert_eq!(get(&store, &key, range(8, 100)).unwrap(), "89");
549        assert!(matches!(
550            block_on(store.get(&key, range(10, 12))),
551            Err(StoreError::RangeNotSatisfiable { len: 10 })
552        ));
553        assert!(matches!(
554            block_on(store.get(&key, range(5, 4))),
555            Err(StoreError::Invalid(_))
556        ));
557    }
558
559    #[test]
560    fn blob_delete_then_get_none_and_second_delete_false() {
561        let store = MemoryBlobStore::default();
562        let key = key_of(b"x");
563        put(&store, key, 1, &[b"x"]).unwrap();
564        assert!(block_on(store.delete(&key)).unwrap());
565        assert_eq!(get(&store, &key, None), None);
566        assert!(!block_on(store.delete(&key)).unwrap());
567    }
568
569    #[test]
570    fn large_bodies_stream_in_bounded_pieces() {
571        let store = MemoryBlobStore::default();
572        let data: Vec<u8> = (0..=250_u8).cycle().take(STREAM_CHUNK * 2 + 7).collect();
573        let key = key_of(&data);
574        put(&store, key, data.len() as u64, &[&data]).unwrap();
575        let read = |range| {
576            let body = block_on(store.get(&key, range)).unwrap().unwrap();
577            let BlobBody::Stream { len, mut stream } = body else {
578                panic!("a large body is streamed");
579            };
580            let mut out = Vec::new();
581            while let Some(piece) =
582                block_on(core::future::poll_fn(|cx| stream.as_mut().poll_next(cx)))
583            {
584                let piece = piece.unwrap();
585                assert!(piece.len() <= MAX_BLOB_PIECE_BYTES);
586                out.extend_from_slice(&piece);
587            }
588            assert_eq!(out.len() as u64, len);
589            out
590        };
591        assert_eq!(read(None), data);
592        let range = ByteRange {
593            start: 3,
594            end_inclusive: STREAM_CHUNK as u64 + 3,
595        };
596        assert_eq!(read(Some(range)), &data[3..=STREAM_CHUNK + 3]);
597        // Exactly one piece is still one buffer.
598        let edge = &data[..STREAM_CHUNK];
599        let edge_key = key_of(edge);
600        put(&store, edge_key, edge.len() as u64, &[edge]).unwrap();
601        assert_eq!(get(&store, &edge_key, None).unwrap().len(), STREAM_CHUNK);
602    }
603
604    #[test]
605    fn poisoned_lock_recovers() {
606        let store = MemoryBlobStore::default();
607        let key = key_of(b"a");
608        put(&store, key, 1, &[b"a"]).unwrap();
609        let shared = store.shared.clone();
610        let panicked = std::thread::spawn(move || {
611            let _guard = shared.blobs.lock().unwrap();
612            panic!("poison the blob lock");
613        })
614        .join();
615        assert!(panicked.is_err() && store.shared.blobs.is_poisoned());
616        assert_eq!(get(&store, &key, None).unwrap(), "a");
617        let other = key_of(b"b");
618        assert_eq!(
619            put(&store, other, 1, &[b"b"]).unwrap(),
620            CommitOutcome::Created
621        );
622        assert!(block_on(store.delete(&key)).unwrap());
623    }
624
625    #[test]
626    fn blob_faults_fire_once_and_leave_nothing() {
627        let key = key_of(b"abc");
628        let store = MemoryBlobStore::default().with_fault(MemoryFault::BlobWrite(1));
629        assert!(matches!(
630            put(&store, key, 3, &[b"a", b"bc"]),
631            Err(StoreError::Unavailable(_))
632        ));
633        assert_eq!(
634            put(&store, key, 3, &[b"a", b"bc"]).unwrap(),
635            CommitOutcome::Created
636        );
637        let store = MemoryBlobStore::default().with_fault(MemoryFault::BlobCommit);
638        assert!(matches!(
639            put(&store, key, 3, &[b"abc"]),
640            Err(StoreError::Unavailable(_))
641        ));
642        assert_eq!(block_on(store.head(&key)).unwrap(), None);
643        assert_eq!(
644            put(&store, key, 3, &[b"abc"]).unwrap(),
645            CommitOutcome::Created
646        );
647    }
648
649    #[test]
650    fn multipart_parts_are_idempotent_and_bad_reupload_keeps_the_good_part() {
651        let store = MemoryBlobStore::default();
652        let mut data = vec![0x31; usize::try_from(MIN_PART_SIZE).unwrap()];
653        data.extend_from_slice(b"last part");
654        let key = key_of(&data);
655        let plan = PartPlan::new(data.len() as u64, MIN_PART_SIZE, 2).unwrap();
656        let session = block_on(store.begin_multipart(key, plan.total(), plan.part_size())).unwrap();
657        let first = &data[..usize::try_from(MIN_PART_SIZE).unwrap()];
658        let last = &data[usize::try_from(MIN_PART_SIZE).unwrap()..];
659        let cv0 = part_subtree_cv(&plan, 0, first).unwrap();
660        let cv1 = part_subtree_cv(&plan, 1, last).unwrap();
661
662        // Send the last part first, then the large part in bounded chunks.
663        let mut sink = block_on(store.begin_part(key, &session, &plan, 1, cv1)).unwrap();
664        block_on(sink.write(Bytes::copy_from_slice(last))).unwrap();
665        let tag1 = block_on(sink.commit()).unwrap();
666        let upload_first = || {
667            let mut sink = block_on(store.begin_part(key, &session, &plan, 0, cv0)).unwrap();
668            for chunk in first.chunks(MAX_BLOB_PIECE_BYTES) {
669                block_on(sink.write(Bytes::copy_from_slice(chunk))).unwrap();
670            }
671            block_on(sink.commit()).unwrap()
672        };
673        let tag0 = upload_first();
674        assert_eq!(upload_first(), tag0);
675
676        let mut bad = block_on(store.begin_part(key, &session, &plan, 0, cv0)).unwrap();
677        block_on(bad.write(Bytes::from(vec![
678            0x32;
679            usize::try_from(MIN_PART_SIZE).unwrap()
680        ])))
681        .unwrap();
682        assert!(matches!(
683            block_on(bad.commit()),
684            Err(StoreError::PartSubtreeMismatch)
685        ));
686        assert_eq!(block_on(store.head(&key)).unwrap(), None);
687
688        let parts = [
689            PartRef {
690                index: 0,
691                len: MIN_PART_SIZE,
692                tag: tag0,
693            },
694            PartRef {
695                index: 1,
696                len: last.len() as u64,
697                tag: tag1,
698            },
699        ];
700        assert_eq!(
701            block_on(store.complete(key, &session, &plan, &parts)).unwrap(),
702            CommitOutcome::Created
703        );
704        assert_eq!(
705            block_on(store.head(&key)).unwrap(),
706            Some(BlobMeta {
707                len: data.len() as u64
708            })
709        );
710        assert!(matches!(
711            block_on(store.complete(key, &session, &plan, &parts)),
712            Err(StoreError::SessionGone)
713        ));
714    }
715
716    #[test]
717    fn abort_and_unknown_sessions_are_gone() {
718        let store = MemoryBlobStore::default();
719        let plan = PartPlan::new(MIN_PART_SIZE + 1, MIN_PART_SIZE, 2).unwrap();
720        let key = key_of(b"absent");
721        let session = block_on(store.begin_multipart(key, plan.total(), plan.part_size())).unwrap();
722        assert!(matches!(
723            block_on(store.begin_part(key, b"unknown", &plan, 0, [0; 32])),
724            Err(StoreError::SessionGone)
725        ));
726        block_on(store.abort(key, &session)).unwrap();
727        block_on(store.abort(key, &session)).unwrap();
728        assert!(matches!(
729            block_on(store.begin_part(key, &session, &plan, 0, [0; 32])),
730            Err(StoreError::SessionGone)
731        ));
732        assert_eq!(block_on(store.head(&key)).unwrap(), None);
733    }
734
735    #[test]
736    fn marker_and_pack_hashes_have_separate_memory_keys() {
737        let store = MemoryBlobStore::default();
738        let content = b"marker";
739        let pack = key_of(content);
740        let marker = BlobKey::upload_marker(*pack.hash());
741        assert_ne!(pack, marker);
742        assert_eq!(
743            put(&store, marker, content.len() as u64, &[content]).unwrap(),
744            CommitOutcome::Created
745        );
746        assert_eq!(block_on(store.head(&pack)).unwrap(), None);
747        assert_eq!(
748            block_on(store.head(&marker)).unwrap(),
749            Some(BlobMeta {
750                len: content.len() as u64
751            })
752        );
753    }
754}