Skip to main content

commonware_runtime/utils/buffer/
mod.rs

1//! Buffers for reading and writing to [crate::Blob]s.
2
3use 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/// A shared sync result.
15///
16/// Handles returned by [Completion::handle] are detached observers of the same shared result:
17/// dropping one neither cancels the underlying sync nor consumes its result, and every
18/// observer sees the same outcome.
19#[derive(Clone)]
20struct Completion(Shared<BoxFuture<'static, Result<(), crate::Error>>>);
21
22impl Completion {
23    /// Return a handle for the sync result.
24    fn handle(&self) -> crate::Handle<()> {
25        crate::Handle::from_future(self.0.clone())
26    }
27
28    /// Wait for the sync result.
29    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
40/// Tracks whether blob mutations still need a sync.
41///
42/// Callers rely on three properties:
43/// - Every operation that mutates the blob first waits for an in-flight sync, so a started
44///   sync's coverage is never disturbed by later writes.
45/// - [SyncState::start_sync] on a [SyncState::Pending] state returns the in-flight sync's
46///   handle (completed syncs resolve immediately), so re-requesting a sync is a cheap way to
47///   observe outstanding work.
48/// - A failure is never lost: every handle cloned from the shared completion reports it, and
49///   an unobserved failure surfaces from [SyncState::wait_for_pending] on the next operation,
50///   which also marks the state [SyncState::Dirty] since the mutations still need durability.
51enum SyncState {
52    // No unsynced mutations.
53    Clean,
54    // Unsynced mutations need a sync.
55    Dirty,
56    // A started sync is in flight.
57    Pending(Completion),
58}
59
60impl SyncState {
61    /// Whether no mutation still needs a sync.
62    const fn is_clean(&self) -> bool {
63        matches!(self, Self::Clean)
64    }
65
66    /// Mark a new unsynced mutation.
67    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    /// Wait for an in-flight sync before reusing or mutating the blob.
76    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                // The sync failed, so the pending mutations still need durability.
87                *self = Self::Dirty;
88                Err(err)
89            }
90        }
91    }
92
93    /// Write data with the provided options while tracking durability.
94    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                // Earlier mutations need a full durability barrier too.
112                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                // If this fails, a later sync must still cover the attempted write.
120                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    /// Resize the blob and require a later sync.
130    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    /// Make all pending mutations durable before returning.
138    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    /// Start making pending mutations durable and return a handle for completion.
149    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                // Store a shared completion so repeated calls observe the same sync.
154                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        /// All data currently stored in the blob.
181        ///
182        /// This includes every durable byte plus any newer bytes that have not
183        /// been made durable yet.
184        data: Vec<u8>,
185
186        /// Prefix/ranges of `data` that would survive a crash.
187        durable: Vec<u8>,
188
189        /// Number of write operations.
190        writes: usize,
191
192        /// Number of full sync barriers.
193        full_syncs: usize,
194
195        /// Number of range-scoped write syncs.
196        range_syncs: usize,
197
198        /// Number of writes carrying the uncached hint.
199        uncached_writes: usize,
200
201        /// Number of range-scoped write syncs carrying the uncached hint.
202        uncached_range_syncs: usize,
203    }
204
205    /// Test blob with separate visible and durable state.
206    ///
207    /// Writes and resizes only update `data`. `write_at(SYNC)` updates `data`
208    /// and then copies only that submitted range into `durable`. `sync` copies all
209    /// of `data` to `durable`. This lets tests assert that `Write::sync` uses range
210    /// sync only when no earlier unsynced mutation needs a full durability barrier.
211    #[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            // SAFETY: `len` bytes are filled by copy_from_slice below.
281            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            // Test basic buffered reading functionality with sequential reads
335            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            // Create a buffered reader with small buffer to test refilling
344            let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(10));
345
346            // Read some data
347            let read = reader.read(5).await.unwrap().coalesce();
348            assert_eq!(read.as_ref(), b"Hello");
349
350            // Read more data that requires a buffer refill
351            let read = reader.read(14).await.unwrap().coalesce();
352            assert_eq!(read.as_ref(), b", world! This ");
353
354            // Verify position tracking
355            assert_eq!(reader.position(), 19);
356
357            // Read the remaining data
358            let read = reader.read(7).await.unwrap().coalesce();
359            assert_eq!(read.as_ref(), b"is a te");
360
361            // Attempt to read beyond the end should fail
362            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            // A dedicated single-slot pool: reader construction must not touch
378            // it (an eagerly-allocated unwritten buffer would be checked out of
379            // the pool and, under empty-freeze semantics, released straight
380            // back).
381            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            // No size class may report a created buffer: construction must not
387            // check one out of the pool even transiently. The created gauge
388            // registers per-class label sets lazily, so a clean pool encodes no
389            // member lines at all.
390            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            // The single slot is also unclaimed at read time.
399            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            // Test reading data that spans multiple buffer refills
414            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            // Use a buffer smaller than the total data size
423            let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(10));
424
425            // Read data that crosses buffer boundaries
426            let read = reader.read(15).await.unwrap().coalesce();
427            assert_eq!(read.as_ref(), b"ABCDEFGHIJKLMNO");
428
429            // Verify position tracking
430            assert_eq!(reader.position(), 15);
431
432            // Read the remaining data
433            let read = reader.read(11).await.unwrap().coalesce();
434            assert_eq!(read.as_ref(), b"PQRSTUVWXYZ");
435
436            // Verify we're at the end
437            assert_eq!(reader.position(), 26);
438            assert_eq!(reader.blob_remaining(), 0);
439        });
440    }
441
442    // Regression test for https://github.com/commonwarexyz/monorepo/issues/1348
443    #[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            // Read data that crosses buffer boundaries
458            let read = reader.read(21).await.unwrap().coalesce();
459            assert_eq!(read.as_ref(), b"ABCDEFGHIJKLMNOPQRSTU");
460
461            // Verify position tracking
462            assert_eq!(reader.position(), 21);
463
464            // Read the remaining data
465            let read = reader.read(5).await.unwrap().coalesce();
466            assert_eq!(read.as_ref(), b"VWXYZ");
467
468            // Rewind and read again
469            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            // Test reader behavior with known blob size limits
480            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            // Create a buffered reader with buffer smaller than total data
489            let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(10));
490
491            // Check initial remaining bytes
492            assert_eq!(reader.blob_remaining(), size);
493
494            // Read partial data
495            let read = reader.read(5).await.unwrap().coalesce();
496            assert_eq!(read.as_ref(), b"This ");
497
498            // Check remaining bytes after partial read
499            assert_eq!(reader.blob_remaining(), size - 5);
500
501            // Read exactly up to the size limit
502            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            // Verify we're at the end
506            assert_eq!(reader.blob_remaining(), 0);
507
508            // Reading beyond the end should fail
509            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            // Fill the internal buffer and consume most of it (2 bytes remain buffered).
531            let first = reader.read(6).await.unwrap().coalesce();
532            assert_eq!(first.as_ref(), b"abcdef");
533            assert_eq!(reader.position(), 6);
534
535            // Only 4 bytes remain total, so this must fail without consuming anything.
536            let err = reader.read(5).await.unwrap_err();
537            assert!(matches!(err, Error::BlobInsufficientLength));
538            assert_eq!(reader.position(), 6);
539
540            // Remaining bytes should still be readable in full.
541            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            // Test reading large amounts of data in chunks
552            let data_size = 1024 * 256; // 256KB of data
553            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            // Use a buffer much smaller than the total data
562            let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(64 * 1024));
563
564            // Read all data in smaller chunks
565            let mut total_read = 0;
566            let chunk_size = 8 * 1024; // 8KB chunks
567
568            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                // Verify data integrity
573                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            // Verify we read everything
582            assert_eq!(total_read, data_size);
583
584            // Reading beyond the end should fail
585            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            // Create a blob with exactly 2.5 buffer sizes of data
595            let buffer_size = 1024;
596            let data_size = buffer_size * 5 / 2; // 2.5 buffers
597            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            // Read exactly one buffer size
609            let read = reader.read(buffer_size).await.unwrap().coalesce();
610            assert!(read.as_ref().iter().all(|&b| b == 0x37));
611
612            // Read exactly one buffer size more
613            let read = reader.read(buffer_size).await.unwrap().coalesce();
614            assert!(read.as_ref().iter().all(|&b| b == 0x37));
615
616            // Read the remaining half buffer
617            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            // Verify we're at the end
622            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            // First read fits in one fetched chunk.
641            let first = reader.read(3).await.unwrap();
642            assert!(first.is_single());
643            assert_eq!(first.coalesce().as_ref(), b"ABC");
644
645            // This read spans refill boundaries and should be represented as multiple chunks.
646            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            // Create a memory blob with some test data
657            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            // Create a buffer reader
666            let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(10));
667
668            // Read some data to advance the position
669            let read = reader.read(5).await.unwrap().coalesce();
670            assert_eq!(read.as_ref(), b"ABCDE");
671            assert_eq!(reader.position(), 5);
672
673            // Seek to a specific position
674            reader.seek_to(10).unwrap();
675            assert_eq!(reader.position(), 10);
676
677            // Read data from the new position
678            let read = reader.read(5).await.unwrap().coalesce();
679            assert_eq!(read.as_ref(), b"KLMNO");
680
681            // Seek to beginning
682            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            // Seek to end
689            reader.seek_to(size).unwrap();
690            assert_eq!(reader.position(), size);
691
692            // Trying to read should fail
693            let result = reader.read(1).await;
694            assert!(matches!(result, Err(Error::BlobInsufficientLength)));
695
696            // Seek beyond end should fail
697            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            // Create a memory blob with longer data
707            let data = vec![0x41; 1000]; // 1000 'A' characters
708            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            // Create a buffer reader with small buffer
716            let mut reader = Read::from_pooler(&context, blob, size, NZUsize!(10));
717
718            // Read some data
719            let _ = reader.read(5).await.unwrap().coalesce();
720
721            // Seek far ahead, past the current buffer
722            reader.seek_to(500).unwrap();
723
724            // Read data - should get data from position 500
725            let read = reader.read(5).await.unwrap().coalesce();
726            assert_eq!(read.as_ref(), b"AAAAA"); // Should still be 'A's);
727            assert_eq!(reader.position(), 505);
728
729            // Seek backwards
730            reader.seek_to(100).unwrap();
731
732            // Read again - should be at position 100
733            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            // Reads 0..=5, while the internal fetch cursor advances to 10.
752            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            // Seek back within [buffer_start, fetch_position).
758            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            // First read triggers a single refill of 10 bytes.
786            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            // Seek within the unread buffered window [6, 10).
792            reader.seek_to(7).unwrap();
793            assert_eq!(reader.position(), 7);
794            assert_eq!(reader.buffer_remaining(), 3);
795
796            // Consume only from the already buffered window.
797            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            // Refill should happen only now (at exhaustion), not at seek/read above.
803            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            // Create a memory blob with some test data
815            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            // Create a buffer reader
824            let reader = Read::from_pooler(&context, blob.clone(), data_len, NZUsize!(10));
825
826            // Resize the blob to half its size
827            let resize_len = data_len / 2;
828            reader.resize(resize_len).await.unwrap();
829
830            // Reopen to check truncation
831            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            // Create a new buffer and read to verify truncation
835            let mut new_reader = Read::from_pooler(&context, blob, size, NZUsize!(10));
836
837            // Read the content
838            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            // Reading beyond resized size should fail
846            let result = new_reader.read(1).await;
847            assert!(matches!(result, Err(Error::BlobInsufficientLength)));
848
849            // Test resize to larger size
850            new_reader.resize(data_len * 2).await.unwrap();
851
852            // Reopen to check resize
853            let (blob, new_size) = context.open("partition", b"test").await.unwrap();
854            assert_eq!(new_size, data_len * 2);
855
856            // Create a new buffer and read to verify resize
857            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            // Create a memory blob with some test data
872            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            // Create a buffer reader
881            let reader = Read::from_pooler(&context, blob.clone(), data_len, NZUsize!(10));
882
883            // Resize the blob to zero
884            reader.resize(0).await.unwrap();
885
886            // Reopen to check truncation
887            let (blob, size) = context.open("partition", b"test").await.unwrap();
888            assert_eq!(size, 0, "Blob should be resized to zero");
889
890            // Create a new buffer and try to read (should fail)
891            let mut new_reader = Read::from_pooler(&context, blob, size, NZUsize!(10));
892
893            // Reading from resized blob should fail
894            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            // Test basic buffered write and sync functionality
904            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            // Verify data was written correctly
914            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            // Test writes that cause buffer flushes due to capacity limits
927            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            // Verify the final result
938            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            // Test writing data larger than buffer capacity (direct write)
951            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            // Verify the complete data
966            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            // Test sequential appends that exceed buffer capacity
979            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            // Write data that fits in buffer
983            writer.write_at(0, b"hello").await.unwrap();
984            assert_eq!(writer.size(), 5);
985
986            // Append data that causes buffer flush
987            writer.write_at(5, b" world").await.unwrap();
988            writer.sync().await.unwrap();
989            assert_eq!(writer.size(), 11);
990
991            // Verify the complete result
992            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            // Test overwriting data within the buffer and extending it
1005            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            // Initial write
1009            writer.write_at(0, b"abcdefghij").await.unwrap();
1010            assert_eq!(writer.size(), 10);
1011
1012            // Overwrite middle section
1013            writer.write_at(2, b"01234").await.unwrap();
1014            assert_eq!(writer.size(), 10);
1015            writer.sync().await.unwrap();
1016
1017            // Verify overwrite result
1018            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            // Extend buffer and do partial overwrite
1025            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            // Verify final result
1032            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            // Test writing at offsets before the current buffer position
1045            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            // Write data at a later offset first
1049            writer.write_at(10, b"0123456789").await.unwrap();
1050            assert_eq!(writer.size(), 20);
1051
1052            // Write at an earlier offset (should flush buffer first)
1053            writer.write_at(0, b"abcde").await.unwrap();
1054            assert_eq!(writer.size(), 20);
1055            writer.sync().await.unwrap();
1056
1057            // Verify data placement with gap
1058            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            // Fill the gap between existing data
1068            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            // Verify gap is filled
1074            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            // Test blob resize functionality and subsequent writes
1088            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            // Write initial data
1092            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            // Resize to smaller size
1103            writer.resize(5).await.unwrap();
1104            assert_eq!(writer.size(), 5);
1105            writer.sync().await.unwrap();
1106
1107            // Verify resize
1108            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            // Write to resized blob
1115            writer.write_at(0, b"X").await.unwrap();
1116            assert_eq!(writer.size(), 5);
1117            writer.sync().await.unwrap();
1118
1119            // Verify overwrite
1120            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            // Test resize to larger size
1127            writer.resize(10).await.unwrap();
1128            assert_eq!(writer.size(), 10);
1129            writer.sync().await.unwrap();
1130
1131            // Verify resize
1132            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            // Test resize to zero
1140            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            // Ensure the blob is empty
1153            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            // Test reading through writer's read_at method (buffer + blob reads)
1179            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            // Write data that stays in buffer
1183            writer.write_at(0, b"buffered").await.unwrap();
1184            assert_eq!(writer.size(), 8);
1185
1186            // Read from buffer via writer
1187            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            // Reading past buffer end should fail
1194            assert!(writer.read_at(8, 1).await.is_err());
1195
1196            // Write large data that flushes buffer
1197            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            // Read from underlying blob through writer
1203            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            // Buffer new data at the end
1210            writer.write_at(20, b" more data").await.unwrap();
1211            assert_eq!(writer.size(), 30);
1212
1213            // Read newly buffered data
1214            let read_buf_vec_3 = writer.read_at(20, 5).await.unwrap().coalesce();
1215            assert_eq!(read_buf_vec_3, b" more");
1216
1217            // Read spanning both blob and buffer
1218            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            // Verify complete content by reopening
1222            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            // Test writes that cannot be merged into buffer (non-contiguous/too large)
1255            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            // Fill buffer completely
1259            writer.write_at(0, b"0123456789").await.unwrap();
1260            assert_eq!(writer.size(), 10);
1261
1262            // Write at non-contiguous offset (should flush then write directly)
1263            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            // Verify data with gap
1269            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            // Test write that exceeds buffer capacity
1281            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            // Write large data that exceeds capacity
1287            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            // Verify overwrite result
1293            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            // Test that closing writer flushes and persists buffered data
1307            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            // Sync writer to persist data
1313            writer.sync().await.unwrap();
1314
1315            // Verify data persistence
1316            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            // Test direct writes when data exceeds buffer capacity
1329            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            // Write data larger than buffer capacity (should write directly)
1336            let data_large = b"0123456789";
1337            writer.write_at(0, data_large).await.unwrap();
1338            assert_eq!(writer.size(), 10);
1339
1340            // Sync to ensure data is persisted
1341            writer.sync().await.unwrap();
1342
1343            // Verify direct write worked
1344            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            // Now write small data that should be buffered
1354            writer.write_at(10, b"abc").await.unwrap();
1355            assert_eq!(writer.size(), 13);
1356
1357            // Verify it's in buffer by reading through writer
1358            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            // Verify final state
1364            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            // Test complex buffer operations: overwrite and extend within capacity
1380            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            // Write initial data
1387            writer.write_at(0, b"0123456789").await.unwrap();
1388            assert_eq!(writer.size(), 10);
1389
1390            // Overwrite and extend within buffer capacity
1391            writer.write_at(5, b"ABCDEFGHIJ").await.unwrap();
1392            assert_eq!(writer.size(), 15);
1393
1394            // Verify buffer content through writer
1395            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            // Verify persisted result
1401            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            // Test writing at the current logical end of the blob
1417            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            // Write initial data
1421            writer.write_at(0, b"0123456789").await.unwrap();
1422            assert_eq!(writer.size(), 10);
1423            writer.sync().await.unwrap();
1424
1425            // Append at the current size (logical end)
1426            writer.write_at(writer.size(), b"abc").await.unwrap();
1427            assert_eq!(writer.size(), 13);
1428            writer.sync().await.unwrap();
1429
1430            // Verify complete result
1431            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            // Test multiple appends using writer.size()
1444            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            // First write
1451            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            // Append using size()
1457            writer.write_at(writer.size(), b"BBB").await.unwrap();
1458            assert_eq!(writer.size(), 6); // 3 (AAA) + 3 (BBB)
1459            writer.sync().await.unwrap();
1460            assert_eq!(writer.size(), 6);
1461
1462            // Append again using size()
1463            writer.write_at(writer.size(), b"CCC").await.unwrap();
1464            assert_eq!(writer.size(), 9); // 6 + 3 (CCC)
1465            writer.sync().await.unwrap();
1466            assert_eq!(writer.size(), 9);
1467
1468            // Verify final content
1469            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            // Test writing non-contiguously, then appending at the new size
1485            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            // Initial buffered write
1492            writer.write_at(0, b"INITIAL").await.unwrap(); // 7 bytes
1493            assert_eq!(writer.size(), 7);
1494            // Buffer contains "INITIAL", inner.position = 0
1495
1496            // Non-contiguous write, forces flush of "INITIAL" and direct write of "NONCONTIG"
1497            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            // Append at the new size
1503            writer.write_at(writer.size(), b"APPEND").await.unwrap();
1504            assert_eq!(writer.size(), 35); // 29 + 6
1505            writer.sync().await.unwrap();
1506            assert_eq!(writer.size(), 35);
1507
1508            // Verify final content
1509            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            // Test truncating, then appending at the new size
1530            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            // Write initial data and sync
1537            writer.write_at(0, b"0123456789ABCDEF").await.unwrap(); // 16 bytes
1538            assert_eq!(writer.size(), 16);
1539            writer.sync().await.unwrap(); // inner.position = 16, buffer empty
1540            assert_eq!(writer.size(), 16);
1541
1542            // Resize
1543            let resize_to = 5;
1544            writer.resize(resize_to).await.unwrap();
1545            // after resize, inner.position should be `resize_to` (5)
1546            // buffer should be empty
1547            assert_eq!(writer.size(), resize_to);
1548            writer.sync().await.unwrap(); // Ensure truncation is persisted for verify step
1549            assert_eq!(writer.size(), resize_to);
1550
1551            // Append at the new (resized) size
1552            writer.write_at(writer.size(), b"XXXXX").await.unwrap(); // 5 bytes
1553            // inner.buffer = "XXXXX", inner.position = 5
1554            assert_eq!(writer.size(), 10); // 5 (resized) + 5 (XXXXX)
1555            writer.sync().await.unwrap();
1556            assert_eq!(writer.size(), 10);
1557
1558            // Verify final content
1559            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    // Verifies start_sync flushes current bytes, completes durability, and marks the writer clean.
1571    #[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            // Start a sync for buffered bytes and wait for the returned handle.
1579            writer.write_at(0, b"abc").await.unwrap();
1580            let handle = writer.start_sync().await;
1581            handle.await.unwrap();
1582
1583            // The buffered write required a full sync because the fresh writer starts dirty.
1584            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            // The started sync marked the writer clean, so the next buffered write can use a
1591            // range-scoped sync.
1592            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            // Nothing left to sync.
1601            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    // Verifies sync waits for an outstanding start_sync instead of starting new disk work.
1610    #[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            // Hold the started sync open so a later sync cannot finish right away.
1619            let handle = writer.start_sync().await;
1620            let deferred = next_pending_sync(&pending);
1621
1622            // The attempted sync reaches the pending handle and cannot complete yet.
1623            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            // Releasing the original handle lets sync observe the completed disk sync.
1638            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    // Verifies writes made after start_sync wait before they are flushed.
1648    #[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            // Begin syncing the initial dirty state and keep that sync blocked.
1657            let handle = writer.start_sync().await;
1658            let deferred = next_pending_sync(&pending);
1659
1660            // The tip must not reach the blob while the earlier sync is pending.
1661            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            // After the earlier sync completes, the buffered write can be persisted.
1678            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    // Verifies overlapping writes wait before flushing buffered bytes while start_sync is pending.
1691    #[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            // This append is local while the earlier sync is pending.
1708            writer.write_at(3, b"abc").await.unwrap();
1709
1710            // The drained tip must not reach the blob while the earlier sync is pending.
1711            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            // Releasing the sync lets the parked write reach the blob.
1727            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    // Verifies resize does not mutate the blob before an outstanding start_sync completes.
1739    #[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            // Resize must not reach the blob while the earlier sync is pending.
1757            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            // Releasing the sync lets the resize apply.
1773            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            // A fresh writer preserves one sync barrier for mutations that predate wrapping.
1789            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            // The write remains entirely buffered, so sync can make just this range durable.
1797            writer.write_at(0, b"abc").await.unwrap();
1798            writer.sync().await.unwrap();
1799
1800            // No prior plain blob mutation required another full sync barrier.
1801            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            // The prior sync used a range-scoped write, so there is no pending full-sync barrier.
1808            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            // Simulate a plain blob mutation before the writer wraps it.
1824            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            // The first sync must use a full barrier to make the pre-wrapped write durable.
1832            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            // After the barrier is clear, a buffered tip-only write can use range sync again.
1839            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            // Keep the write buffered so sync attempts the clean range-scoped write path.
1860            writer.write_at(0, b"abc").await.unwrap();
1861
1862            // Removing the blob makes the range-sync flush fail.
1863            context.remove("partition", Some(name)).await.unwrap();
1864            assert!(writer.sync().await.is_err());
1865
1866            // The failed range-scoped write must leave a pending full-sync barrier, so a
1867            // later sync cannot report success.
1868            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            // This exceeds the buffer and forces a plain write before the final buffered tip.
1880            writer.write_at(0, b"abcdef").await.unwrap();
1881            writer.write_at(6, b"g").await.unwrap();
1882            writer.sync().await.unwrap();
1883
1884            // The final sync must cover both the prior plain write and the buffered tip.
1885            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            // With no new writes, sync has no work left.
1892            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            // After the full sync, the next buffer-only write can use range sync again.
1900            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            // Establish already-durable data with a range sync.
1920            writer.write_at(0, b"abcdef").await.unwrap();
1921            writer.sync().await.unwrap();
1922
1923            // Resize alone is an unsynced blob mutation.
1924            writer.resize(4).await.unwrap();
1925            writer.sync().await.unwrap();
1926
1927            // The resized contents require a full sync barrier to become durable.
1928            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}