1use crate::WriteOptions;
4use futures::future::{BoxFuture, FutureExt as _, Shared};
5
6pub mod paged;
7mod read;
8mod tip;
9mod write;
10
11pub use read::Read;
12pub use write::Write;
13
14#[derive(Clone)]
20struct Completion(Shared<BoxFuture<'static, Result<(), crate::Error>>>);
21
22impl Completion {
23 fn handle(&self) -> crate::Handle<()> {
25 crate::Handle::from_future(self.0.clone())
26 }
27
28 async fn wait(&self) -> Result<(), crate::Error> {
30 self.0.clone().await
31 }
32}
33
34impl From<crate::Handle<()>> for Completion {
35 fn from(handle: crate::Handle<()>) -> Self {
36 Self(handle.boxed().shared())
37 }
38}
39
40enum SyncState {
52 Clean,
54 Dirty,
56 Pending(Completion),
58}
59
60impl SyncState {
61 const fn is_clean(&self) -> bool {
63 matches!(self, Self::Clean)
64 }
65
66 fn mark_dirty(&mut self) {
68 assert!(
69 !matches!(self, Self::Pending(_)),
70 "pending sync must be joined before marking dirty"
71 );
72 *self = Self::Dirty;
73 }
74
75 async fn wait_for_pending(&mut self) -> Result<(), crate::Error> {
77 let Self::Pending(pending) = self else {
78 return Ok(());
79 };
80 match pending.wait().await {
81 Ok(()) => {
82 *self = Self::Clean;
83 Ok(())
84 }
85 Err(err) => {
86 *self = Self::Dirty;
88 Err(err)
89 }
90 }
91 }
92
93 async fn write_at(
95 &mut self,
96 blob: &impl crate::Blob,
97 offset: u64,
98 bufs: impl Into<crate::IoBufs> + Send,
99 options: WriteOptions,
100 ) -> Result<(), crate::Error> {
101 self.wait_for_pending().await?;
102 let bufs = bufs.into();
103 if !options.contains(WriteOptions::SYNC) {
104 blob.write_at(offset, bufs, options).await?;
105 self.mark_dirty();
106 return Ok(());
107 }
108
109 match self {
110 Self::Dirty => {
111 blob.write_at(offset, bufs, options.without(WriteOptions::SYNC))
113 .await?;
114 blob.sync().await?;
115 *self = Self::Clean;
116 Ok(())
117 }
118 Self::Clean => {
119 self.mark_dirty();
121 blob.write_at(offset, bufs, options).await?;
122 *self = Self::Clean;
123 Ok(())
124 }
125 Self::Pending(_) => unreachable!("pending sync waited above"),
126 }
127 }
128
129 async fn resize(&mut self, blob: &impl crate::Blob, len: u64) -> Result<(), crate::Error> {
131 self.wait_for_pending().await?;
132 blob.resize(len).await?;
133 self.mark_dirty();
134 Ok(())
135 }
136
137 async fn sync(&mut self, blob: &impl crate::Blob) -> Result<(), crate::Error> {
139 self.wait_for_pending().await?;
140 if matches!(self, Self::Clean) {
141 return Ok(());
142 }
143 blob.sync().await?;
144 *self = Self::Clean;
145 Ok(())
146 }
147
148 async fn start_sync(&mut self, blob: &impl crate::Blob) -> crate::Handle<()> {
150 match self {
151 Self::Clean => crate::Handle::ready(Ok(())),
152 Self::Dirty => {
153 let pending = Completion::from(blob.start_sync().await);
155 let handle = pending.handle();
156 *self = Self::Pending(pending);
157 handle
158 }
159 Self::Pending(pending) => pending.handle(),
160 }
161 }
162}
163
164#[cfg(test)]
165mod tests {
166 use super::*;
167 use crate::{
168 Blob as _, BufferPool, BufferPoolConfig, Error, Handle, IoBufMut, IoBufs, IoBufsMut,
169 ReadOptions, Runner, Storage, WriteOptions, deterministic,
170 mocks::{DelayedSyncBlob, next_pending_sync},
171 telemetry::metrics::Registry,
172 };
173 use commonware_macros::test_traced;
174 use commonware_utils::{NZU32, NZUsize, sync::Mutex};
175 use futures::FutureExt;
176 use std::sync::Arc;
177
178 #[derive(Default)]
179 struct RangeSyncState {
180 data: Vec<u8>,
185
186 durable: Vec<u8>,
188
189 writes: usize,
191
192 full_syncs: usize,
194
195 range_syncs: usize,
197
198 uncached_writes: usize,
200
201 uncached_range_syncs: usize,
203 }
204
205 #[derive(Clone)]
212 pub struct SyncTrackingBlob {
213 state: Arc<Mutex<RangeSyncState>>,
214 }
215
216 impl SyncTrackingBlob {
217 pub fn new() -> Self {
218 Self {
219 state: Arc::new(Mutex::new(RangeSyncState::default())),
220 }
221 }
222
223 pub fn snapshot(&self) -> (Vec<u8>, usize, usize, usize) {
224 let state = self.state.lock();
225 (
226 state.durable.clone(),
227 state.writes,
228 state.full_syncs,
229 state.range_syncs,
230 )
231 }
232
233 pub fn size(&self) -> u64 {
234 self.state.lock().data.len() as u64
235 }
236
237 pub fn uncached_snapshot(&self) -> (usize, usize) {
238 let state = self.state.lock();
239 (state.uncached_writes, state.uncached_range_syncs)
240 }
241
242 fn write(data: &mut Vec<u8>, offset: u64, buf: &[u8]) -> Result<(), Error> {
243 let start = usize::try_from(offset).map_err(|_| Error::OffsetOverflow)?;
244 let end = start.checked_add(buf.len()).ok_or(Error::OffsetOverflow)?;
245 if end > data.len() {
246 data.resize(end, 0);
247 }
248 data[start..end].copy_from_slice(buf);
249 Ok(())
250 }
251 }
252
253 impl crate::Blob for SyncTrackingBlob {
254 async fn read_at(
255 &self,
256 offset: u64,
257 len: usize,
258 options: ReadOptions,
259 ) -> Result<IoBufsMut, Error> {
260 self.read_at_buf(offset, len, IoBufMut::with_capacity(len), options)
261 .await
262 }
263
264 async fn read_at_buf(
265 &self,
266 offset: u64,
267 len: usize,
268 buf: impl Into<IoBufsMut> + Send,
269 _options: ReadOptions,
270 ) -> Result<IoBufsMut, Error> {
271 let start = usize::try_from(offset).map_err(|_| Error::OffsetOverflow)?;
272 let end = start.checked_add(len).ok_or(Error::OffsetOverflow)?;
273 let state = self.state.lock();
274 if end > state.data.len() {
275 return Err(Error::BlobInsufficientLength);
276 }
277
278 let mut out = buf.into();
279 assert!(out.capacity() >= len);
280 unsafe { out.set_len(len) };
282 out.copy_from_slice(&state.data[start..end]);
283 Ok(out)
284 }
285
286 async fn write_at(
287 &self,
288 offset: u64,
289 buf: impl Into<IoBufs> + Send,
290 options: WriteOptions,
291 ) -> Result<(), Error> {
292 let buf = buf.into().coalesce();
293 let mut state = self.state.lock();
294 Self::write(&mut state.data, offset, buf.as_ref())?;
295 state.writes += 1;
296
297 let sync = options.contains(WriteOptions::SYNC);
298 if sync {
299 Self::write(&mut state.durable, offset, buf.as_ref())?;
300 state.range_syncs += 1;
301 }
302 if options.contains(WriteOptions::DONT_CACHE) {
303 if sync {
304 state.uncached_range_syncs += 1;
305 } else {
306 state.uncached_writes += 1;
307 }
308 }
309 Ok(())
310 }
311
312 async fn resize(&self, len: u64) -> Result<(), Error> {
313 let len = usize::try_from(len).map_err(|_| Error::OffsetOverflow)?;
314 self.state.lock().data.resize(len, 0);
315 Ok(())
316 }
317
318 async fn sync(&self) -> Result<(), Error> {
319 let mut state = self.state.lock();
320 state.durable = state.data.clone();
321 state.full_syncs += 1;
322 Ok(())
323 }
324
325 async fn start_sync(&self) -> Handle<()> {
326 Handle::ready(self.sync().await)
327 }
328 }
329
330 #[test_traced]
331 fn test_read_basic() {
332 let executor = deterministic::Runner::default();
333 executor.start(|context| async move {
334 let data = b"Hello, world! This is a test.";
336 let (blob, size) = context.open("partition", b"test").await.unwrap();
337 assert_eq!(size, 0);
338 blob.write_at(0, data, WriteOptions::default())
339 .await
340 .unwrap();
341 let size = data.len() as u64;
342
343 let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(10));
345
346 let read = reader.read(5).await.unwrap().coalesce();
348 assert_eq!(read.as_ref(), b"Hello");
349
350 let read = reader.read(14).await.unwrap().coalesce();
352 assert_eq!(read.as_ref(), b", world! This ");
353
354 assert_eq!(reader.position(), 19);
356
357 let read = reader.read(7).await.unwrap().coalesce();
359 assert_eq!(read.as_ref(), b"is a te");
360
361 let result = reader.read(5).await;
363 assert!(matches!(result, Err(Error::BlobInsufficientLength)));
364 });
365 }
366
367 #[test_traced]
368 fn test_read_allocates_lazily_on_first_refill() {
369 let executor = deterministic::Runner::default();
370 executor.start(|context| async move {
371 let data = b"hello world";
372 let (blob, _) = context.open("partition", b"lazy").await.unwrap();
373 blob.write_at(0, data, WriteOptions::default())
374 .await
375 .unwrap();
376
377 let mut registry = Registry::default();
382 let config = BufferPoolConfig::for_storage().with_max_per_class(NZU32!(1));
383 let pool = BufferPool::new(config, &mut registry);
384 let mut reader = Read::new(blob, data.len() as u64, NZUsize!(10), pool.clone());
385
386 let encoded = registry.encode();
391 assert!(
392 !encoded
393 .lines()
394 .any(|line| line.starts_with("buffer_pool_created{")),
395 "reader construction created a pool buffer: {encoded}"
396 );
397
398 let probe = pool
400 .try_alloc(1)
401 .expect("reader construction must not retain a buffer");
402 drop(probe);
403
404 let read = reader.read(5).await.unwrap().coalesce();
405 assert_eq!(read.as_ref(), b"hello");
406 });
407 }
408
409 #[test_traced]
410 fn test_read_cross_boundary() {
411 let executor = deterministic::Runner::default();
412 executor.start(|context| async move {
413 let data = b"ABCDEFGHIJKLMNOPQRSTUVWXYZ";
415 let (blob, size) = context.open("partition", b"test").await.unwrap();
416 assert_eq!(size, 0);
417 blob.write_at(0, data, WriteOptions::default())
418 .await
419 .unwrap();
420 let size = data.len() as u64;
421
422 let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(10));
424
425 let read = reader.read(15).await.unwrap().coalesce();
427 assert_eq!(read.as_ref(), b"ABCDEFGHIJKLMNO");
428
429 assert_eq!(reader.position(), 15);
431
432 let read = reader.read(11).await.unwrap().coalesce();
434 assert_eq!(read.as_ref(), b"PQRSTUVWXYZ");
435
436 assert_eq!(reader.position(), 26);
438 assert_eq!(reader.blob_remaining(), 0);
439 });
440 }
441
442 #[test_traced]
444 fn test_read_to_end_then_rewind_and_read_again() {
445 let executor = deterministic::Runner::default();
446 executor.start(|context| async move {
447 let data = b"ABCDEFGHIJKLMNOPQRSTUVWXYZ";
448 let (blob, size) = context.open("partition", b"test").await.unwrap();
449 assert_eq!(size, 0);
450 blob.write_at(0, data, WriteOptions::default())
451 .await
452 .unwrap();
453 let size = data.len() as u64;
454
455 let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(20));
456
457 let read = reader.read(21).await.unwrap().coalesce();
459 assert_eq!(read.as_ref(), b"ABCDEFGHIJKLMNOPQRSTU");
460
461 assert_eq!(reader.position(), 21);
463
464 let read = reader.read(5).await.unwrap().coalesce();
466 assert_eq!(read.as_ref(), b"VWXYZ");
467
468 reader.seek_to(0).unwrap();
470 let read = reader.read(21).await.unwrap().coalesce();
471 assert_eq!(read.as_ref(), b"ABCDEFGHIJKLMNOPQRSTU");
472 });
473 }
474
475 #[test_traced]
476 fn test_read_with_known_size() {
477 let executor = deterministic::Runner::default();
478 executor.start(|context| async move {
479 let data = b"This is a test with known size limitations.";
481 let (blob, size) = context.open("partition", b"test").await.unwrap();
482 assert_eq!(size, 0);
483 blob.write_at(0, data, WriteOptions::default())
484 .await
485 .unwrap();
486 let size = data.len() as u64;
487
488 let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(10));
490
491 assert_eq!(reader.blob_remaining(), size);
493
494 let read = reader.read(5).await.unwrap().coalesce();
496 assert_eq!(read.as_ref(), b"This ");
497
498 assert_eq!(reader.blob_remaining(), size - 5);
500
501 let read = reader.read((size - 5) as usize).await.unwrap().coalesce();
503 assert_eq!(read.as_ref(), b"is a test with known size limitations.");
504
505 assert_eq!(reader.blob_remaining(), 0);
507
508 let result = reader.read(1).await;
510 assert!(matches!(result, Err(Error::BlobInsufficientLength)));
511 });
512 }
513
514 #[test_traced]
515 fn test_read_oversized_request_does_not_consume_buffered_bytes() {
516 let executor = deterministic::Runner::default();
517 executor.start(|context| async move {
518 let data = b"abcdefghij";
519 let (blob, size) = context
520 .open("partition", b"double-count-regression")
521 .await
522 .unwrap();
523 assert_eq!(size, 0);
524 blob.write_at(0, data, WriteOptions::default())
525 .await
526 .unwrap();
527
528 let mut reader = Read::from_pooler(&context, blob, data.len() as u64, NZUsize!(8));
529
530 let first = reader.read(6).await.unwrap().coalesce();
532 assert_eq!(first.as_ref(), b"abcdef");
533 assert_eq!(reader.position(), 6);
534
535 let err = reader.read(5).await.unwrap_err();
537 assert!(matches!(err, Error::BlobInsufficientLength));
538 assert_eq!(reader.position(), 6);
539
540 let tail = reader.read(4).await.unwrap().coalesce();
542 assert_eq!(tail.as_ref(), b"ghij");
543 assert_eq!(reader.position(), 10);
544 });
545 }
546
547 #[test_traced]
548 fn test_read_large_data() {
549 let executor = deterministic::Runner::default();
550 executor.start(|context| async move {
551 let data_size = 1024 * 256; let data = vec![0x42; data_size];
554 let (blob, size) = context.open("partition", b"test").await.unwrap();
555 assert_eq!(size, 0);
556 blob.write_at(0, data.clone(), WriteOptions::default())
557 .await
558 .unwrap();
559 let size = data.len() as u64;
560
561 let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(64 * 1024));
563
564 let mut total_read = 0;
566 let chunk_size = 8 * 1024; while total_read < data_size {
569 let to_read = std::cmp::min(chunk_size, data_size - total_read);
570 let read = reader.read(to_read).await.unwrap().coalesce();
571
572 assert!(
574 read.as_ref().iter().all(|&b| b == 0x42),
575 "Data at position {total_read} is not correct"
576 );
577
578 total_read += to_read;
579 }
580
581 assert_eq!(total_read, data_size);
583
584 let result = reader.read(1).await;
586 assert!(matches!(result, Err(Error::BlobInsufficientLength)));
587 });
588 }
589
590 #[test_traced]
591 fn test_read_exact_size_reads() {
592 let executor = deterministic::Runner::default();
593 executor.start(|context| async move {
594 let buffer_size = 1024;
596 let data_size = buffer_size * 5 / 2; let data = vec![0x37; data_size];
598
599 let (blob, size) = context.open("partition", b"test").await.unwrap();
600 assert_eq!(size, 0);
601 blob.write_at(0, data.clone(), WriteOptions::default())
602 .await
603 .unwrap();
604 let size = data.len() as u64;
605
606 let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(buffer_size));
607
608 let read = reader.read(buffer_size).await.unwrap().coalesce();
610 assert!(read.as_ref().iter().all(|&b| b == 0x37));
611
612 let read = reader.read(buffer_size).await.unwrap().coalesce();
614 assert!(read.as_ref().iter().all(|&b| b == 0x37));
615
616 let half_buffer = buffer_size / 2;
618 let read = reader.read(half_buffer).await.unwrap().coalesce();
619 assert!(read.as_ref().iter().all(|&b| b == 0x37));
620
621 assert_eq!(reader.blob_remaining(), 0);
623 assert_eq!(reader.position(), size);
624 });
625 }
626
627 #[test_traced]
628 fn test_read_structure_single_vs_chunked() {
629 let executor = deterministic::Runner::default();
630 executor.start(|context| async move {
631 let data = b"ABCDEFGHIJKL";
632 let (blob, size) = context.open("partition", b"structural").await.unwrap();
633 assert_eq!(size, 0);
634 blob.write_at(0, data, WriteOptions::default())
635 .await
636 .unwrap();
637
638 let mut reader = Read::from_pooler(&context, blob, data.len() as u64, NZUsize!(5));
639
640 let first = reader.read(3).await.unwrap();
642 assert!(first.is_single());
643 assert_eq!(first.coalesce().as_ref(), b"ABC");
644
645 let second = reader.read(7).await.unwrap();
647 assert!(!second.is_single());
648 assert_eq!(second.coalesce().as_ref(), b"DEFGHIJ");
649 });
650 }
651
652 #[test_traced]
653 fn test_read_seek_to() {
654 let executor = deterministic::Runner::default();
655 executor.start(|context| async move {
656 let data = b"ABCDEFGHIJKLMNOPQRSTUVWXYZ";
658 let (blob, size) = context.open("partition", b"test").await.unwrap();
659 assert_eq!(size, 0);
660 blob.write_at(0, data, WriteOptions::default())
661 .await
662 .unwrap();
663 let size = data.len() as u64;
664
665 let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(10));
667
668 let read = reader.read(5).await.unwrap().coalesce();
670 assert_eq!(read.as_ref(), b"ABCDE");
671 assert_eq!(reader.position(), 5);
672
673 reader.seek_to(10).unwrap();
675 assert_eq!(reader.position(), 10);
676
677 let read = reader.read(5).await.unwrap().coalesce();
679 assert_eq!(read.as_ref(), b"KLMNO");
680
681 reader.seek_to(0).unwrap();
683 assert_eq!(reader.position(), 0);
684
685 let read = reader.read(5).await.unwrap().coalesce();
686 assert_eq!(read.as_ref(), b"ABCDE");
687
688 reader.seek_to(size).unwrap();
690 assert_eq!(reader.position(), size);
691
692 let result = reader.read(1).await;
694 assert!(matches!(result, Err(Error::BlobInsufficientLength)));
695
696 let result = reader.seek_to(size + 10);
698 assert!(matches!(result, Err(Error::BlobInsufficientLength)));
699 });
700 }
701
702 #[test_traced]
703 fn test_read_seek_with_refill() {
704 let executor = deterministic::Runner::default();
705 executor.start(|context| async move {
706 let data = vec![0x41; 1000]; let (blob, size) = context.open("partition", b"test").await.unwrap();
709 assert_eq!(size, 0);
710 blob.write_at(0, data.clone(), WriteOptions::default())
711 .await
712 .unwrap();
713 let size = data.len() as u64;
714
715 let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(10));
717
718 let _ = reader.read(5).await.unwrap().coalesce();
720
721 reader.seek_to(500).unwrap();
723
724 let read = reader.read(5).await.unwrap().coalesce();
726 assert_eq!(read.as_ref(), b"AAAAA"); assert_eq!(reader.position(), 505);
728
729 reader.seek_to(100).unwrap();
731
732 let _ = reader.read(5).await.unwrap().coalesce();
734 assert_eq!(reader.position(), 105);
735 });
736 }
737
738 #[test_traced]
739 fn test_read_seek_within_buffered_range() {
740 let executor = deterministic::Runner::default();
741 executor.start(|context| async move {
742 let data = b"ABCDEFGHIJKLMNOPQRSTUVWXYZ";
743 let (blob, size) = context.open("partition", b"test").await.unwrap();
744 assert_eq!(size, 0);
745 blob.write_at(0, data, WriteOptions::default())
746 .await
747 .unwrap();
748
749 let mut reader = Read::from_pooler(&context, blob, data.len() as u64, NZUsize!(10));
750
751 let read = reader.read(6).await.unwrap().coalesce();
753 assert_eq!(read.as_ref(), b"ABCDEF");
754 assert_eq!(reader.position(), 6);
755 assert_eq!(reader.buffer_remaining(), 4);
756
757 reader.seek_to(3).unwrap();
759 assert_eq!(reader.position(), 3);
760 assert_eq!(reader.buffer_remaining(), 7);
761
762 let read = reader.read(5).await.unwrap().coalesce();
763 assert_eq!(read.as_ref(), b"DEFGH");
764 assert_eq!(reader.position(), 8);
765 assert_eq!(reader.buffer_remaining(), 2);
766 });
767 }
768
769 #[test_traced]
770 fn test_read_seek_within_unread_buffer_does_not_refill() {
771 let executor = deterministic::Runner::default();
772 executor.start(|context| async move {
773 let data = b"ABCDEFGHIJKLMNOPQRSTUVWXYZ";
774 let (blob, size) = context
775 .open("partition", b"seek_unread_no_refill")
776 .await
777 .unwrap();
778 assert_eq!(size, 0);
779 blob.write_at(0, data, WriteOptions::default())
780 .await
781 .unwrap();
782
783 let mut reader = Read::from_pooler(&context, blob, data.len() as u64, NZUsize!(10));
784
785 let first = reader.read(6).await.unwrap();
787 assert_eq!(first.coalesce().as_ref(), b"ABCDEF");
788 assert_eq!(reader.position(), 6);
789 assert_eq!(reader.buffer_remaining(), 4);
790
791 reader.seek_to(7).unwrap();
793 assert_eq!(reader.position(), 7);
794 assert_eq!(reader.buffer_remaining(), 3);
795
796 let second = reader.read(3).await.unwrap();
798 assert_eq!(second.coalesce().as_ref(), b"HIJ");
799 assert_eq!(reader.position(), 10);
800 assert_eq!(reader.buffer_remaining(), 0);
801
802 let third = reader.read(1).await.unwrap();
804 assert_eq!(third.coalesce().as_ref(), b"K");
805 assert_eq!(reader.position(), 11);
806 assert_eq!(reader.buffer_remaining(), 9);
807 });
808 }
809
810 #[test_traced]
811 fn test_read_resize() {
812 let executor = deterministic::Runner::default();
813 executor.start(|context| async move {
814 let data = b"ABCDEFGHIJKLMNOPQRSTUVWXYZ";
816 let (blob, size) = context.open("partition", b"test").await.unwrap();
817 assert_eq!(size, 0);
818 blob.write_at(0, data, WriteOptions::default())
819 .await
820 .unwrap();
821 let data_len = data.len() as u64;
822
823 let reader = Read::from_pooler(&context, blob.clone(), data_len, NZUsize!(10));
825
826 let resize_len = data_len / 2;
828 reader.resize(resize_len).await.unwrap();
829
830 let (blob, size) = context.open("partition", b"test").await.unwrap();
832 assert_eq!(size, resize_len, "Blob should be resized to half size");
833
834 let mut new_reader = Read::from_pooler(&context, blob, size, NZUsize!(10));
836
837 let read = new_reader.read(size as usize).await.unwrap().coalesce();
839 assert_eq!(
840 read.as_ref(),
841 b"ABCDEFGHIJKLM",
842 "Resized content should match"
843 );
844
845 let result = new_reader.read(1).await;
847 assert!(matches!(result, Err(Error::BlobInsufficientLength)));
848
849 new_reader.resize(data_len * 2).await.unwrap();
851
852 let (blob, new_size) = context.open("partition", b"test").await.unwrap();
854 assert_eq!(new_size, data_len * 2);
855
856 let mut new_reader = Read::from_pooler(&context, blob, new_size, NZUsize!(10));
858 let read = new_reader.read(new_size as usize).await.unwrap().coalesce();
859 assert_eq!(&read.as_ref()[..size as usize], b"ABCDEFGHIJKLM");
860 assert_eq!(
861 &read.as_ref()[size as usize..],
862 vec![0u8; new_size as usize - size as usize]
863 );
864 });
865 }
866
867 #[test_traced]
868 fn test_read_resize_to_zero() {
869 let executor = deterministic::Runner::default();
870 executor.start(|context| async move {
871 let data = b"ABCDEFGHIJKLMNOPQRSTUVWXYZ";
873 let data_len = data.len() as u64;
874 let (blob, size) = context.open("partition", b"test").await.unwrap();
875 assert_eq!(size, 0);
876 blob.write_at(0, data, WriteOptions::default())
877 .await
878 .unwrap();
879
880 let reader = Read::from_pooler(&context, blob.clone(), data_len, NZUsize!(10));
882
883 reader.resize(0).await.unwrap();
885
886 let (blob, size) = context.open("partition", b"test").await.unwrap();
888 assert_eq!(size, 0, "Blob should be resized to zero");
889
890 let mut new_reader = Read::from_pooler(&context, blob, size, NZUsize!(10));
892
893 let result = new_reader.read(1).await;
895 assert!(matches!(result, Err(Error::BlobInsufficientLength)));
896 });
897 }
898
899 #[test_traced]
900 fn test_write_basic() {
901 let executor = deterministic::Runner::default();
902 executor.start(|context| async move {
903 let (blob, size) = context.open("partition", b"write_basic").await.unwrap();
905 assert_eq!(size, 0);
906
907 let mut writer = Write::from_pooler(&context, blob.clone(), size, NZUsize!(8));
908 writer.write_at(0, b"hello").await.unwrap();
909 assert_eq!(writer.size(), 5);
910 writer.sync().await.unwrap();
911 assert_eq!(writer.size(), 5);
912
913 let (blob, size) = context.open("partition", b"write_basic").await.unwrap();
915 assert_eq!(size, 5);
916 let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(8));
917 let read = reader.read(5).await.unwrap().coalesce();
918 assert_eq!(read.as_ref(), b"hello");
919 });
920 }
921
922 #[test_traced]
923 fn test_write_multiple_flushes() {
924 let executor = deterministic::Runner::default();
925 executor.start(|context| async move {
926 let (blob, size) = context.open("partition", b"write_multi").await.unwrap();
928 assert_eq!(size, 0);
929
930 let mut writer = Write::from_pooler(&context, blob.clone(), size, NZUsize!(4));
931 writer.write_at(0, b"abc").await.unwrap();
932 assert_eq!(writer.size(), 3);
933 writer.write_at(3, b"defg").await.unwrap();
934 assert_eq!(writer.size(), 7);
935 writer.sync().await.unwrap();
936
937 let (blob, size) = context.open("partition", b"write_multi").await.unwrap();
939 assert_eq!(size, 7);
940 let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(4));
941 let read = reader.read(7).await.unwrap().coalesce();
942 assert_eq!(read.as_ref(), b"abcdefg");
943 });
944 }
945
946 #[test_traced]
947 fn test_write_large_data() {
948 let executor = deterministic::Runner::default();
949 executor.start(|context| async move {
950 let (blob, size) = context.open("partition", b"write_large").await.unwrap();
952 assert_eq!(size, 0);
953
954 let mut writer = Write::from_pooler(&context, blob.clone(), size, NZUsize!(4));
955 writer.write_at(0, b"abc").await.unwrap();
956 assert_eq!(writer.size(), 3);
957 writer
958 .write_at(3, b"defghijklmnopqrstuvwxyz")
959 .await
960 .unwrap();
961 assert_eq!(writer.size(), 26);
962 writer.sync().await.unwrap();
963 assert_eq!(writer.size(), 26);
964
965 let (blob, size) = context.open("partition", b"write_large").await.unwrap();
967 assert_eq!(size, 26);
968 let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(4));
969 let read = reader.read(26).await.unwrap().coalesce();
970 assert_eq!(read.as_ref(), b"abcdefghijklmnopqrstuvwxyz");
971 });
972 }
973
974 #[test_traced]
975 fn test_write_append_to_buffer() {
976 let executor = deterministic::Runner::default();
977 executor.start(|context| async move {
978 let (blob, size) = context.open("partition", b"append_buf").await.unwrap();
980 let mut writer = Write::from_pooler(&context, blob.clone(), size, NZUsize!(10));
981
982 writer.write_at(0, b"hello").await.unwrap();
984 assert_eq!(writer.size(), 5);
985
986 writer.write_at(5, b" world").await.unwrap();
988 writer.sync().await.unwrap();
989 assert_eq!(writer.size(), 11);
990
991 let (blob, size) = context.open("partition", b"append_buf").await.unwrap();
993 assert_eq!(size, 11);
994 let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(10));
995 let read = reader.read(11).await.unwrap().coalesce();
996 assert_eq!(read.as_ref(), b"hello world");
997 });
998 }
999
1000 #[test_traced]
1001 fn test_write_into_middle_of_buffer() {
1002 let executor = deterministic::Runner::default();
1003 executor.start(|context| async move {
1004 let (blob, size) = context.open("partition", b"middle_buf").await.unwrap();
1006 let mut writer = Write::from_pooler(&context, blob.clone(), size, NZUsize!(20));
1007
1008 writer.write_at(0, b"abcdefghij").await.unwrap();
1010 assert_eq!(writer.size(), 10);
1011
1012 writer.write_at(2, b"01234").await.unwrap();
1014 assert_eq!(writer.size(), 10);
1015 writer.sync().await.unwrap();
1016
1017 let (blob, size) = context.open("partition", b"middle_buf").await.unwrap();
1019 assert_eq!(size, 10);
1020 let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(10));
1021 let read = reader.read(10).await.unwrap().coalesce();
1022 assert_eq!(read.as_ref(), b"ab01234hij");
1023
1024 writer.write_at(10, b"klmnopqrst").await.unwrap();
1026 assert_eq!(writer.size(), 20);
1027 writer.write_at(9, b"wxyz").await.unwrap();
1028 assert_eq!(writer.size(), 20);
1029 writer.sync().await.unwrap();
1030
1031 let (blob, size) = context.open("partition", b"middle_buf").await.unwrap();
1033 assert_eq!(size, 20);
1034 let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(20));
1035 let read = reader.read(20).await.unwrap().coalesce();
1036 assert_eq!(read.as_ref(), b"ab01234hiwxyznopqrst");
1037 });
1038 }
1039
1040 #[test_traced]
1041 fn test_write_before_buffer() {
1042 let executor = deterministic::Runner::default();
1043 executor.start(|context| async move {
1044 let (blob, size) = context.open("partition", b"before_buf").await.unwrap();
1046 let mut writer = Write::from_pooler(&context, blob.clone(), size, NZUsize!(10));
1047
1048 writer.write_at(10, b"0123456789").await.unwrap();
1050 assert_eq!(writer.size(), 20);
1051
1052 writer.write_at(0, b"abcde").await.unwrap();
1054 assert_eq!(writer.size(), 20);
1055 writer.sync().await.unwrap();
1056
1057 let (blob, size) = context.open("partition", b"before_buf").await.unwrap();
1059 assert_eq!(size, 20);
1060 let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(20));
1061 let read = reader.read(20).await.unwrap().coalesce();
1062 let mut expected = vec![0u8; 20];
1063 expected[0..5].copy_from_slice("abcde".as_bytes());
1064 expected[10..20].copy_from_slice("0123456789".as_bytes());
1065 assert_eq!(read.as_ref(), expected.as_slice());
1066
1067 writer.write_at(5, b"fghij").await.unwrap();
1069 assert_eq!(writer.size(), 20);
1070 writer.sync().await.unwrap();
1071 assert_eq!(writer.size(), 20);
1072
1073 let (blob, size) = context.open("partition", b"before_buf").await.unwrap();
1075 assert_eq!(size, 20);
1076 let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(20));
1077 let read = reader.read(20).await.unwrap().coalesce();
1078 expected[0..10].copy_from_slice("abcdefghij".as_bytes());
1079 assert_eq!(read.as_ref(), expected.as_slice());
1080 });
1081 }
1082
1083 #[test_traced]
1084 fn test_write_resize() {
1085 let executor = deterministic::Runner::default();
1086 executor.start(|context| async move {
1087 let (blob, size) = context.open("partition", b"resize_write").await.unwrap();
1089 let mut writer = Write::from_pooler(&context, blob, size, NZUsize!(10));
1090
1091 writer.write_at(0, b"hello world").await.unwrap();
1093 assert_eq!(writer.size(), 11);
1094 writer.sync().await.unwrap();
1095 assert_eq!(writer.size(), 11);
1096
1097 let (blob_check, size_check) =
1098 context.open("partition", b"resize_write").await.unwrap();
1099 assert_eq!(size_check, 11);
1100 drop(blob_check);
1101
1102 writer.resize(5).await.unwrap();
1104 assert_eq!(writer.size(), 5);
1105 writer.sync().await.unwrap();
1106
1107 let (blob, size) = context.open("partition", b"resize_write").await.unwrap();
1109 assert_eq!(size, 5);
1110 let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(5));
1111 let read = reader.read(5).await.unwrap().coalesce();
1112 assert_eq!(read.as_ref(), b"hello");
1113
1114 writer.write_at(0, b"X").await.unwrap();
1116 assert_eq!(writer.size(), 5);
1117 writer.sync().await.unwrap();
1118
1119 let (blob, size) = context.open("partition", b"resize_write").await.unwrap();
1121 assert_eq!(size, 5);
1122 let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(5));
1123 let read = reader.read(5).await.unwrap().coalesce();
1124 assert_eq!(read.as_ref(), b"Xello");
1125
1126 writer.resize(10).await.unwrap();
1128 assert_eq!(writer.size(), 10);
1129 writer.sync().await.unwrap();
1130
1131 let (blob, size) = context.open("partition", b"resize_write").await.unwrap();
1133 assert_eq!(size, 10);
1134 let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(10));
1135 let read = reader.read(10).await.unwrap().coalesce();
1136 assert_eq!(&read.as_ref()[0..5], b"Xello");
1137 assert_eq!(&read.as_ref()[5..10], [0u8; 5]);
1138
1139 let (blob_zero, size) = context.open("partition", b"resize_zero").await.unwrap();
1141 let mut writer_zero =
1142 Write::from_pooler(&context, blob_zero.clone(), size, NZUsize!(10));
1143 writer_zero.write_at(0, b"some data").await.unwrap();
1144 assert_eq!(writer_zero.size(), 9);
1145 writer_zero.sync().await.unwrap();
1146 assert_eq!(writer_zero.size(), 9);
1147 writer_zero.resize(0).await.unwrap();
1148 assert_eq!(writer_zero.size(), 0);
1149 writer_zero.sync().await.unwrap();
1150 assert_eq!(writer_zero.size(), 0);
1151
1152 let (_, size_z) = context.open("partition", b"resize_zero").await.unwrap();
1154 assert_eq!(size_z, 0);
1155 });
1156 }
1157
1158 #[test_traced]
1159 fn test_write_resize_grow_flushes_buffered_data() {
1160 let executor = deterministic::Runner::default();
1161 executor.start(|context| async move {
1162 let (blob, size) = context.open("partition", b"resize_grow").await.unwrap();
1163 let mut writer = Write::from_pooler(&context, blob, size, NZUsize!(10));
1164
1165 writer.write_at(0, b"hello").await.unwrap();
1166 writer.resize(10).await.unwrap();
1167 writer.sync().await.unwrap();
1168
1169 let read = writer.read_at(0, 5).await.unwrap().coalesce();
1170 assert_eq!(read.as_ref(), b"hello");
1171 });
1172 }
1173
1174 #[test_traced]
1175 fn test_write_read_at_on_writer() {
1176 let executor = deterministic::Runner::default();
1177 executor.start(|context| async move {
1178 let (blob, size) = context.open("partition", b"read_at_writer").await.unwrap();
1180 let mut writer = Write::from_pooler(&context, blob.clone(), size, NZUsize!(10));
1181
1182 writer.write_at(0, b"buffered").await.unwrap();
1184 assert_eq!(writer.size(), 8);
1185
1186 let read_buf_vec = writer.read_at(0, 4).await.unwrap().coalesce();
1188 assert_eq!(read_buf_vec, b"buff");
1189
1190 let read_buf_vec = writer.read_at(4, 4).await.unwrap().coalesce();
1191 assert_eq!(read_buf_vec, b"ered");
1192
1193 assert!(writer.read_at(8, 1).await.is_err());
1195
1196 writer.write_at(8, b" and flushed").await.unwrap();
1198 assert_eq!(writer.size(), 20);
1199 writer.sync().await.unwrap();
1200 assert_eq!(writer.size(), 20);
1201
1202 let read_buf_vec_2 = writer.read_at(0, 4).await.unwrap().coalesce();
1204 assert_eq!(read_buf_vec_2, b"buff");
1205
1206 let read_buf_7_vec = writer.read_at(13, 7).await.unwrap().coalesce();
1207 assert_eq!(read_buf_7_vec, b"flushed");
1208
1209 writer.write_at(20, b" more data").await.unwrap();
1211 assert_eq!(writer.size(), 30);
1212
1213 let read_buf_vec_3 = writer.read_at(20, 5).await.unwrap().coalesce();
1215 assert_eq!(read_buf_vec_3, b" more");
1216
1217 let combo_read_buf_vec = writer.read_at(16, 12).await.unwrap();
1219 assert_eq!(combo_read_buf_vec.coalesce(), b"shed more da");
1220
1221 writer.sync().await.unwrap();
1223 assert_eq!(writer.size(), 30);
1224 let (final_blob, final_size) =
1225 context.open("partition", b"read_at_writer").await.unwrap();
1226 assert_eq!(final_size, 30);
1227 let mut final_reader =
1228 Read::from_pooler(&context, final_blob, final_size, NZUsize!(30));
1229 let read = final_reader.read(30).await.unwrap().coalesce();
1230 assert_eq!(read.as_ref(), b"buffered and flushed more data");
1231 });
1232 }
1233
1234 #[test_traced]
1235 fn test_write_zero_length_read_past_eof_errors() {
1236 let executor = deterministic::Runner::default();
1237 executor.start(|context| async move {
1238 let (blob, size) = context.open("partition", b"zero_len_probe").await.unwrap();
1239 let mut writer = Write::from_pooler(&context, blob, size, NZUsize!(8));
1240 writer.write_at(0, b"abc").await.unwrap();
1241
1242 let empty = writer.read_at(3, 0).await.unwrap();
1243 assert!(empty.is_empty());
1244
1245 let err = writer.read_at(4, 0).await.unwrap_err();
1246 assert!(matches!(err, Error::BlobInsufficientLength));
1247 });
1248 }
1249
1250 #[test_traced]
1251 fn test_write_straddling_non_mergeable() {
1252 let executor = deterministic::Runner::default();
1253 executor.start(|context| async move {
1254 let (blob, size) = context.open("partition", b"write_straddle").await.unwrap();
1256 let mut writer = Write::from_pooler(&context, blob.clone(), size, NZUsize!(10));
1257
1258 writer.write_at(0, b"0123456789").await.unwrap();
1260 assert_eq!(writer.size(), 10);
1261
1262 writer.write_at(15, b"abc").await.unwrap();
1264 assert_eq!(writer.size(), 18);
1265 writer.sync().await.unwrap();
1266 assert_eq!(writer.size(), 18);
1267
1268 let (blob_check, size_check) =
1270 context.open("partition", b"write_straddle").await.unwrap();
1271 assert_eq!(size_check, 18);
1272 let mut reader = Read::from_pooler(&context, blob_check, size_check, NZUsize!(20));
1273 let read = reader.read(18).await.unwrap().coalesce();
1274
1275 let mut expected = vec![0u8; 18];
1276 expected[0..10].copy_from_slice(b"0123456789");
1277 expected[15..18].copy_from_slice(b"abc");
1278 assert_eq!(read.as_ref(), expected.as_slice());
1279
1280 let (blob2, size) = context.open("partition", b"write_straddle2").await.unwrap();
1282 let mut writer2 = Write::from_pooler(&context, blob2.clone(), size, NZUsize!(10));
1283 writer2.write_at(0, b"0123456789").await.unwrap();
1284 assert_eq!(writer2.size(), 10);
1285
1286 writer2.write_at(5, b"ABCDEFGHIJKL").await.unwrap();
1288 assert_eq!(writer2.size(), 17);
1289 writer2.sync().await.unwrap();
1290 assert_eq!(writer2.size(), 17);
1291
1292 let (blob_check2, size_check2) =
1294 context.open("partition", b"write_straddle2").await.unwrap();
1295 assert_eq!(size_check2, 17);
1296 let mut reader2 = Read::from_pooler(&context, blob_check2, size_check2, NZUsize!(20));
1297 let read = reader2.read(17).await.unwrap().coalesce();
1298 assert_eq!(read.as_ref(), b"01234ABCDEFGHIJKL");
1299 });
1300 }
1301
1302 #[test_traced]
1303 fn test_write_close() {
1304 let executor = deterministic::Runner::default();
1305 executor.start(|context| async move {
1306 let (blob_orig, size) = context.open("partition", b"write_close").await.unwrap();
1308 let mut writer = Write::from_pooler(&context, blob_orig.clone(), size, NZUsize!(8));
1309 writer.write_at(0, b"pending").await.unwrap();
1310 assert_eq!(writer.size(), 7);
1311
1312 writer.sync().await.unwrap();
1314
1315 let (blob_check, size_check) = context.open("partition", b"write_close").await.unwrap();
1317 assert_eq!(size_check, 7);
1318 let mut reader = Read::from_pooler(&context, blob_check, size_check, NZUsize!(8));
1319 let read = reader.read(7).await.unwrap().coalesce();
1320 assert_eq!(read.as_ref(), b"pending");
1321 });
1322 }
1323
1324 #[test_traced]
1325 fn test_write_direct_due_to_size() {
1326 let executor = deterministic::Runner::default();
1327 executor.start(|context| async move {
1328 let (blob, size) = context
1330 .open("partition", b"write_direct_size")
1331 .await
1332 .unwrap();
1333 let mut writer = Write::from_pooler(&context, blob.clone(), size, NZUsize!(5));
1334
1335 let data_large = b"0123456789";
1337 writer.write_at(0, data_large).await.unwrap();
1338 assert_eq!(writer.size(), 10);
1339
1340 writer.sync().await.unwrap();
1342
1343 let (blob_check, size_check) = context
1345 .open("partition", b"write_direct_size")
1346 .await
1347 .unwrap();
1348 assert_eq!(size_check, 10);
1349 let mut reader = Read::from_pooler(&context, blob_check, size_check, NZUsize!(10));
1350 let read = reader.read(10).await.unwrap().coalesce();
1351 assert_eq!(read.as_ref(), data_large.as_slice());
1352
1353 writer.write_at(10, b"abc").await.unwrap();
1355 assert_eq!(writer.size(), 13);
1356
1357 let read_small_buf_vec = writer.read_at(10, 3).await.unwrap().coalesce();
1359 assert_eq!(read_small_buf_vec, b"abc");
1360
1361 writer.sync().await.unwrap();
1362
1363 let (blob_check2, size_check2) = context
1365 .open("partition", b"write_direct_size")
1366 .await
1367 .unwrap();
1368 assert_eq!(size_check2, 13);
1369 let mut reader2 = Read::from_pooler(&context, blob_check2, size_check2, NZUsize!(13));
1370 let read = reader2.read(13).await.unwrap().coalesce();
1371 assert_eq!(&read.as_ref()[10..], b"abc".as_slice());
1372 });
1373 }
1374
1375 #[test_traced]
1376 fn test_write_overwrite_and_extend_in_buffer() {
1377 let executor = deterministic::Runner::default();
1378 executor.start(|context| async move {
1379 let (blob, size) = context
1381 .open("partition", b"overwrite_extend_buf")
1382 .await
1383 .unwrap();
1384 let mut writer = Write::from_pooler(&context, blob.clone(), size, NZUsize!(15));
1385
1386 writer.write_at(0, b"0123456789").await.unwrap();
1388 assert_eq!(writer.size(), 10);
1389
1390 writer.write_at(5, b"ABCDEFGHIJ").await.unwrap();
1392 assert_eq!(writer.size(), 15);
1393
1394 let read_buf_vec = writer.read_at(0, 15).await.unwrap().coalesce();
1396 assert_eq!(read_buf_vec, b"01234ABCDEFGHIJ");
1397
1398 writer.sync().await.unwrap();
1399
1400 let (blob_check, size_check) = context
1402 .open("partition", b"overwrite_extend_buf")
1403 .await
1404 .unwrap();
1405 assert_eq!(size_check, 15);
1406 let mut reader = Read::from_pooler(&context, blob_check, size_check, NZUsize!(15));
1407 let read = reader.read(15).await.unwrap().coalesce();
1408 assert_eq!(read.as_ref(), b"01234ABCDEFGHIJ".as_slice());
1409 });
1410 }
1411
1412 #[test_traced]
1413 fn test_write_at_size() {
1414 let executor = deterministic::Runner::default();
1415 executor.start(|context| async move {
1416 let (blob, size) = context.open("partition", b"write_end").await.unwrap();
1418 let mut writer = Write::from_pooler(&context, blob.clone(), size, NZUsize!(20));
1419
1420 writer.write_at(0, b"0123456789").await.unwrap();
1422 assert_eq!(writer.size(), 10);
1423 writer.sync().await.unwrap();
1424
1425 writer.write_at(writer.size(), b"abc").await.unwrap();
1427 assert_eq!(writer.size(), 13);
1428 writer.sync().await.unwrap();
1429
1430 let (blob_check, size_check) = context.open("partition", b"write_end").await.unwrap();
1432 assert_eq!(size_check, 13);
1433 let mut reader = Read::from_pooler(&context, blob_check, size_check, NZUsize!(13));
1434 let read = reader.read(13).await.unwrap().coalesce();
1435 assert_eq!(read.as_ref(), b"0123456789abc");
1436 });
1437 }
1438
1439 #[test_traced]
1440 fn test_write_at_size_multiple_appends() {
1441 let executor = deterministic::Runner::default();
1442 executor.start(|context| async move {
1443 let (blob, size) = context
1445 .open("partition", b"write_multiple_appends_at_size")
1446 .await
1447 .unwrap();
1448 let mut writer = Write::from_pooler(&context, blob.clone(), size, NZUsize!(5));
1449
1450 writer.write_at(0, b"AAA").await.unwrap();
1452 assert_eq!(writer.size(), 3);
1453 writer.sync().await.unwrap();
1454 assert_eq!(writer.size(), 3);
1455
1456 writer.write_at(writer.size(), b"BBB").await.unwrap();
1458 assert_eq!(writer.size(), 6); writer.sync().await.unwrap();
1460 assert_eq!(writer.size(), 6);
1461
1462 writer.write_at(writer.size(), b"CCC").await.unwrap();
1464 assert_eq!(writer.size(), 9); writer.sync().await.unwrap();
1466 assert_eq!(writer.size(), 9);
1467
1468 let (blob_check, size_check) = context
1470 .open("partition", b"write_multiple_appends_at_size")
1471 .await
1472 .unwrap();
1473 assert_eq!(size_check, 9);
1474 let mut reader = Read::from_pooler(&context, blob_check, size_check, NZUsize!(9));
1475 let read = reader.read(9).await.unwrap().coalesce();
1476 assert_eq!(read.as_ref(), b"AAABBBCCC");
1477 });
1478 }
1479
1480 #[test_traced]
1481 fn test_write_non_contiguous_then_append_at_size() {
1482 let executor = deterministic::Runner::default();
1483 executor.start(|context| async move {
1484 let (blob, size) = context
1486 .open("partition", b"write_non_contiguous_then_append")
1487 .await
1488 .unwrap();
1489 let mut writer = Write::from_pooler(&context, blob.clone(), size, NZUsize!(10));
1490
1491 writer.write_at(0, b"INITIAL").await.unwrap(); assert_eq!(writer.size(), 7);
1494 writer.write_at(20, b"NONCONTIG").await.unwrap();
1498 assert_eq!(writer.size(), 29);
1499 writer.sync().await.unwrap();
1500 assert_eq!(writer.size(), 29);
1501
1502 writer.write_at(writer.size(), b"APPEND").await.unwrap();
1504 assert_eq!(writer.size(), 35); writer.sync().await.unwrap();
1506 assert_eq!(writer.size(), 35);
1507
1508 let (blob_check, size_check) = context
1510 .open("partition", b"write_non_contiguous_then_append")
1511 .await
1512 .unwrap();
1513 assert_eq!(size_check, 35);
1514 let mut reader = Read::from_pooler(&context, blob_check, size_check, NZUsize!(35));
1515 let read = reader.read(35).await.unwrap().coalesce();
1516
1517 let mut expected = vec![0u8; 35];
1518 expected[0..7].copy_from_slice(b"INITIAL");
1519 expected[20..29].copy_from_slice(b"NONCONTIG");
1520 expected[29..35].copy_from_slice(b"APPEND");
1521 assert_eq!(read.as_ref(), expected.as_slice());
1522 });
1523 }
1524
1525 #[test_traced]
1526 fn test_write_resize_then_append_at_size() {
1527 let executor = deterministic::Runner::default();
1528 executor.start(|context| async move {
1529 let (blob, size) = context
1531 .open("partition", b"resize_then_append_at_size")
1532 .await
1533 .unwrap();
1534 let mut writer = Write::from_pooler(&context, blob.clone(), size, NZUsize!(10));
1535
1536 writer.write_at(0, b"0123456789ABCDEF").await.unwrap(); assert_eq!(writer.size(), 16);
1539 writer.sync().await.unwrap(); assert_eq!(writer.size(), 16);
1541
1542 let resize_to = 5;
1544 writer.resize(resize_to).await.unwrap();
1545 assert_eq!(writer.size(), resize_to);
1548 writer.sync().await.unwrap(); assert_eq!(writer.size(), resize_to);
1550
1551 writer.write_at(writer.size(), b"XXXXX").await.unwrap(); assert_eq!(writer.size(), 10); writer.sync().await.unwrap();
1556 assert_eq!(writer.size(), 10);
1557
1558 let (blob_check, size_check) = context
1560 .open("partition", b"resize_then_append_at_size")
1561 .await
1562 .unwrap();
1563 assert_eq!(size_check, 10);
1564 let mut reader = Read::from_pooler(&context, blob_check, size_check, NZUsize!(10));
1565 let read = reader.read(10).await.unwrap().coalesce();
1566 assert_eq!(read.as_ref(), b"01234XXXXX");
1567 });
1568 }
1569
1570 #[test_traced]
1572 fn test_write_start_sync_persists_and_marks_clean() {
1573 let executor = deterministic::Runner::default();
1574 executor.start(|context| async move {
1575 let blob = SyncTrackingBlob::new();
1576 let mut writer = Write::from_pooler(&context, blob.clone(), 0, NZUsize!(8));
1577
1578 writer.write_at(0, b"abc").await.unwrap();
1580 let handle = writer.start_sync().await;
1581 handle.await.unwrap();
1582
1583 let (durable, writes, full_syncs, range_syncs) = blob.snapshot();
1585 assert_eq!(durable.as_slice(), b"abc");
1586 assert_eq!(writes, 1);
1587 assert_eq!(full_syncs, 1);
1588 assert_eq!(range_syncs, 0);
1589
1590 writer.write_at(3, b"d").await.unwrap();
1593 writer.sync().await.unwrap();
1594 let (durable, writes, full_syncs, range_syncs) = blob.snapshot();
1595 assert_eq!(durable.as_slice(), b"abcd");
1596 assert_eq!(writes, 2);
1597 assert_eq!(full_syncs, 1);
1598 assert_eq!(range_syncs, 1);
1599
1600 let handle = writer.start_sync().await;
1602 handle.await.unwrap();
1603 let (_, _, full_syncs, range_syncs) = blob.snapshot();
1604 assert_eq!(full_syncs, 1);
1605 assert_eq!(range_syncs, 1);
1606 });
1607 }
1608
1609 #[test_traced]
1611 fn test_write_sync_waits_for_outstanding_start_sync() {
1612 let executor = deterministic::Runner::default();
1613 executor.start(|context| async move {
1614 let inner = SyncTrackingBlob::new();
1615 let (blob, pending) = DelayedSyncBlob::new(inner.clone());
1616 let mut writer = Write::from_pooler(&context, blob, 0, NZUsize!(8));
1617
1618 let handle = writer.start_sync().await;
1620 let deferred = next_pending_sync(&pending);
1621
1622 let mut sync = Box::pin(writer.sync());
1624 assert!(
1625 sync.as_mut().now_or_never().is_none(),
1626 "sync must wait for the outstanding start_sync handle"
1627 );
1628 deferred
1629 .blocked
1630 .await
1631 .expect("sync never waited on start_sync");
1632
1633 let (_, _, full_syncs, range_syncs) = inner.snapshot();
1634 assert_eq!(full_syncs, 0);
1635 assert_eq!(range_syncs, 0);
1636
1637 deferred.release.send(Ok(())).unwrap();
1639 sync.await.unwrap();
1640 handle.await.unwrap();
1641 let (_, _, full_syncs, range_syncs) = inner.snapshot();
1642 assert_eq!(full_syncs, 1);
1643 assert_eq!(range_syncs, 0);
1644 });
1645 }
1646
1647 #[test_traced]
1649 fn test_write_sync_after_start_sync_and_small_write_waits_before_range_sync() {
1650 let executor = deterministic::Runner::default();
1651 executor.start(|context| async move {
1652 let inner = SyncTrackingBlob::new();
1653 let (blob, pending) = DelayedSyncBlob::new(inner.clone());
1654 let mut writer = Write::from_pooler(&context, blob, 0, NZUsize!(8));
1655
1656 let handle = writer.start_sync().await;
1658 let deferred = next_pending_sync(&pending);
1659
1660 writer.write_at(0, b"abc").await.unwrap();
1662 let mut sync = Box::pin(writer.sync());
1663 assert!(
1664 sync.as_mut().now_or_never().is_none(),
1665 "sync must wait for the outstanding start_sync before flushing the small write"
1666 );
1667 deferred
1668 .blocked
1669 .await
1670 .expect("sync never waited on start_sync");
1671
1672 let (_, writes, full_syncs, range_syncs) = inner.snapshot();
1673 assert_eq!(writes, 0);
1674 assert_eq!(full_syncs, 0);
1675 assert_eq!(range_syncs, 0);
1676
1677 deferred.release.send(Ok(())).unwrap();
1679 sync.await.unwrap();
1680 handle.await.unwrap();
1681
1682 let (durable, writes, full_syncs, range_syncs) = inner.snapshot();
1683 assert_eq!(durable.as_slice(), b"abc");
1684 assert_eq!(writes, 1);
1685 assert_eq!(full_syncs, 1);
1686 assert_eq!(range_syncs, 1);
1687 });
1688 }
1689
1690 #[test_traced]
1692 fn test_write_at_overlap_flush_waits_for_outstanding_start_sync() {
1693 let executor = deterministic::Runner::default();
1694 executor.start(|context| async move {
1695 let inner = SyncTrackingBlob::new();
1696 inner
1697 .write_at(0, b"xxx", WriteOptions::default())
1698 .await
1699 .unwrap();
1700
1701 let (blob, pending) = DelayedSyncBlob::new(inner.clone());
1702 let mut writer = Write::from_pooler(&context, blob, inner.size(), NZUsize!(8));
1703
1704 let handle = writer.start_sync().await;
1705 let deferred = next_pending_sync(&pending);
1706
1707 writer.write_at(3, b"abc").await.unwrap();
1709
1710 let mut write = Box::pin(writer.write_at(2, b"ZZ"));
1712 assert!(
1713 write.as_mut().now_or_never().is_none(),
1714 "overlapping write must wait for the outstanding start_sync before flushing"
1715 );
1716 deferred
1717 .blocked
1718 .await
1719 .expect("write never waited on start_sync");
1720
1721 let (_, writes, full_syncs, range_syncs) = inner.snapshot();
1722 assert_eq!(writes, 1);
1723 assert_eq!(full_syncs, 0);
1724 assert_eq!(range_syncs, 0);
1725
1726 deferred.release.send(Ok(())).unwrap();
1728 write.await.unwrap();
1729 handle.await.unwrap();
1730
1731 let (_, writes, full_syncs, range_syncs) = inner.snapshot();
1732 assert_eq!(writes, 3);
1733 assert_eq!(full_syncs, 1);
1734 assert_eq!(range_syncs, 0);
1735 });
1736 }
1737
1738 #[test_traced]
1740 fn test_write_resize_waits_for_outstanding_start_sync_before_resizing() {
1741 let executor = deterministic::Runner::default();
1742 executor.start(|context| async move {
1743 let inner = SyncTrackingBlob::new();
1744 inner
1745 .write_at(0, b"abcdef", WriteOptions::default())
1746 .await
1747 .unwrap();
1748
1749 let (blob, pending) = DelayedSyncBlob::new(inner.clone());
1750 let mut writer = Write::from_pooler(&context, blob, inner.size(), NZUsize!(8));
1751
1752 let handle = writer.start_sync().await;
1753 let deferred = next_pending_sync(&pending);
1754 let original_size = inner.size();
1755
1756 let mut resize = Box::pin(writer.resize(3));
1758 assert!(
1759 resize.as_mut().now_or_never().is_none(),
1760 "resize must wait for the outstanding start_sync handle"
1761 );
1762 deferred
1763 .blocked
1764 .await
1765 .expect("resize never waited on start_sync");
1766 assert_eq!(
1767 inner.size(),
1768 original_size,
1769 "resize must not mutate the blob before the pending sync finishes"
1770 );
1771
1772 deferred.release.send(Ok(())).unwrap();
1774 resize.await.unwrap();
1775 handle.await.unwrap();
1776 assert_eq!(writer.size(), 3);
1777 assert_eq!(inner.size(), 3);
1778 });
1779 }
1780
1781 #[test_traced]
1782 fn test_write_sync_uses_range_sync_for_buffer_only_write() {
1783 let executor = deterministic::Runner::default();
1784 executor.start(|context| async move {
1785 let blob = SyncTrackingBlob::new();
1786 let mut writer = Write::from_pooler(&context, blob.clone(), 0, NZUsize!(8));
1787
1788 writer.sync().await.unwrap();
1790 let (durable, writes, full_syncs, range_syncs) = blob.snapshot();
1791 assert!(durable.is_empty());
1792 assert_eq!(writes, 0);
1793 assert_eq!(full_syncs, 1);
1794 assert_eq!(range_syncs, 0);
1795
1796 writer.write_at(0, b"abc").await.unwrap();
1798 writer.sync().await.unwrap();
1799
1800 let (durable, writes, full_syncs, range_syncs) = blob.snapshot();
1802 assert_eq!(durable.as_slice(), b"abc");
1803 assert_eq!(writes, 1);
1804 assert_eq!(full_syncs, 1);
1805 assert_eq!(range_syncs, 1);
1806
1807 writer.sync().await.unwrap();
1809 let (durable, writes, full_syncs, range_syncs) = blob.snapshot();
1810 assert_eq!(durable.as_slice(), b"abc");
1811 assert_eq!(writes, 1);
1812 assert_eq!(full_syncs, 1);
1813 assert_eq!(range_syncs, 1);
1814 });
1815 }
1816
1817 #[test_traced]
1818 fn test_write_sync_persists_pre_wrapped_blob_mutation() {
1819 let executor = deterministic::Runner::default();
1820 executor.start(|context| async move {
1821 let blob = SyncTrackingBlob::new();
1822
1823 blob.write_at(0, b"abc", WriteOptions::default())
1825 .await
1826 .unwrap();
1827
1828 let mut writer = Write::from_pooler(&context, blob.clone(), 3, NZUsize!(8));
1829 writer.sync().await.unwrap();
1830
1831 let (durable, writes, full_syncs, range_syncs) = blob.snapshot();
1833 assert_eq!(durable.as_slice(), b"abc");
1834 assert_eq!(writes, 1);
1835 assert_eq!(full_syncs, 1);
1836 assert_eq!(range_syncs, 0);
1837
1838 writer.write_at(3, b"d").await.unwrap();
1840 writer.sync().await.unwrap();
1841
1842 let (durable, writes, full_syncs, range_syncs) = blob.snapshot();
1843 assert_eq!(durable.as_slice(), b"abcd");
1844 assert_eq!(writes, 2);
1845 assert_eq!(full_syncs, 1);
1846 assert_eq!(range_syncs, 1);
1847 });
1848 }
1849
1850 #[test_traced]
1851 fn test_write_sync_failed_range_sync_does_not_mark_clean() {
1852 let executor = deterministic::Runner::default();
1853 executor.start(|context| async move {
1854 let name = b"failed_range_sync";
1855 let (blob, size) = context.open("partition", name).await.unwrap();
1856 let mut writer = Write::from_pooler(&context, blob, size, NZUsize!(8));
1857 writer.sync().await.unwrap();
1858
1859 writer.write_at(0, b"abc").await.unwrap();
1861
1862 context.remove("partition", Some(name)).await.unwrap();
1864 assert!(writer.sync().await.is_err());
1865
1866 assert!(writer.sync().await.is_err());
1869 });
1870 }
1871
1872 #[test_traced]
1873 fn test_write_sync_persists_prior_direct_flushes_with_buffered_tip() {
1874 let executor = deterministic::Runner::default();
1875 executor.start(|context| async move {
1876 let blob = SyncTrackingBlob::new();
1877 let mut writer = Write::from_pooler(&context, blob.clone(), 0, NZUsize!(4));
1878
1879 writer.write_at(0, b"abcdef").await.unwrap();
1881 writer.write_at(6, b"g").await.unwrap();
1882 writer.sync().await.unwrap();
1883
1884 let (durable, writes, full_syncs, range_syncs) = blob.snapshot();
1886 assert_eq!(durable.as_slice(), b"abcdefg");
1887 assert_eq!(writes, 2);
1888 assert_eq!(full_syncs, 1);
1889 assert_eq!(range_syncs, 0);
1890
1891 writer.sync().await.unwrap();
1893 let (durable, writes, full_syncs, range_syncs) = blob.snapshot();
1894 assert_eq!(durable.as_slice(), b"abcdefg");
1895 assert_eq!(writes, 2);
1896 assert_eq!(full_syncs, 1);
1897 assert_eq!(range_syncs, 0);
1898
1899 writer.write_at(7, b"h").await.unwrap();
1901 writer.sync().await.unwrap();
1902
1903 let (durable, writes, full_syncs, range_syncs) = blob.snapshot();
1904 assert_eq!(durable.as_slice(), b"abcdefgh");
1905 assert_eq!(writes, 3);
1906 assert_eq!(full_syncs, 1);
1907 assert_eq!(range_syncs, 1);
1908 });
1909 }
1910
1911 #[test_traced]
1912 fn test_write_sync_uses_full_sync_after_resize() {
1913 let executor = deterministic::Runner::default();
1914 executor.start(|context| async move {
1915 let blob = SyncTrackingBlob::new();
1916 let mut writer = Write::from_pooler(&context, blob.clone(), 0, NZUsize!(8));
1917 writer.sync().await.unwrap();
1918
1919 writer.write_at(0, b"abcdef").await.unwrap();
1921 writer.sync().await.unwrap();
1922
1923 writer.resize(4).await.unwrap();
1925 writer.sync().await.unwrap();
1926
1927 let (durable, writes, full_syncs, range_syncs) = blob.snapshot();
1929 assert_eq!(durable.as_slice(), b"abcd");
1930 assert_eq!(writes, 1);
1931 assert_eq!(full_syncs, 2);
1932 assert_eq!(range_syncs, 1);
1933 });
1934 }
1935}