Skip to main content

commonware_runtime/utils/buffer/paged/
read.rs

1use super::Checksum;
2use crate::{Blob, Buf, Error, IoBuf, ReadOptions};
3use commonware_codec::FixedSize;
4use std::{collections::VecDeque, num::NonZeroU16};
5use tracing::error;
6
7/// State for a single buffer of pages read from the blob.
8///
9/// Each fill produces one `BufferState` containing all pages read in that batch.
10/// Navigation skips CRCs by computing offsets rather than creating separate
11/// `Bytes` slices per page.
12pub(super) struct BufferState {
13    /// The raw physical buffer containing pages with interleaved CRCs.
14    buffer: IoBuf,
15    /// Number of pages in this buffer.
16    num_pages: usize,
17    /// Logical length of the last page (may be partial).
18    last_page_len: usize,
19}
20
21/// Async I/O component that prefetches pages and validates CRCs.
22///
23/// This handles reading batches of pages from the blob, validating their
24/// checksums, and producing `BufferState` for the sync buffering layer.
25pub(super) struct PageReader<B: Blob> {
26    /// The underlying blob to read from.
27    blob: B,
28    /// Physical page size (page_size + CHECKSUM_SIZE).
29    physical_page_size: usize,
30    /// Logical page size (data bytes per page, not including CRC).
31    page_size: usize,
32    /// The physical size of the blob.
33    physical_blob_size: u64,
34    /// The size of the blob.
35    logical_blob_size: u64,
36    /// Next page index to read from the blob.
37    blob_page: u64,
38    /// Number of pages to prefetch at once.
39    prefetch_count: usize,
40    /// Options applied to every blob read.
41    read_options: ReadOptions,
42}
43
44impl<B: Blob> PageReader<B> {
45    /// Creates a new PageReader.
46    ///
47    /// The `physical_blob_size` must already exclude any trailing invalid data
48    /// (e.g., junk pages from an interrupted write). Each physical page is the same
49    /// size on disk, but the CRC record indicates how much logical data it contains.
50    /// The last page may be logically partial (CRC length < logical page size), but
51    /// all preceding pages must be logically full. A logically partial non-last page
52    /// indicates corruption and will cause an `Error::InvalidChecksum`.
53    pub(super) fn new(
54        blob: B,
55        physical_blob_size: u64,
56        logical_blob_size: u64,
57        prefetch_count: usize,
58        page_size: NonZeroU16,
59        read_options: ReadOptions,
60    ) -> Self {
61        let page_size = page_size.get() as usize;
62        let physical_page_size = page_size + Checksum::SIZE;
63        let physical_pages = physical_blob_size / physical_page_size as u64;
64        let logical_pages = if logical_blob_size == 0 {
65            0
66        } else {
67            ((logical_blob_size - 1) / page_size as u64) + 1
68        };
69        assert_eq!(physical_blob_size % physical_page_size as u64, 0);
70        assert_eq!(physical_pages, logical_pages);
71
72        Self {
73            blob,
74            physical_page_size,
75            page_size,
76            physical_blob_size,
77            logical_blob_size,
78            blob_page: 0,
79            prefetch_count,
80            read_options,
81        }
82    }
83
84    /// Returns the size of the blob.
85    pub(super) const fn blob_size(&self) -> u64 {
86        self.logical_blob_size
87    }
88
89    /// Returns the physical page size.
90    pub(super) const fn physical_page_size(&self) -> usize {
91        self.physical_page_size
92    }
93
94    /// Returns the logical page size.
95    pub(super) const fn page_size(&self) -> usize {
96        self.page_size
97    }
98
99    /// Fills a buffer with the next batch of pages.
100    ///
101    /// Returns `Some((BufferState, logical_bytes))` if data was loaded,
102    /// `None` if no more data available.
103    pub(super) async fn fill(&mut self) -> Result<Option<(BufferState, usize)>, Error> {
104        // Calculate physical read offset
105        let start_offset = match self.blob_page.checked_mul(self.physical_page_size as u64) {
106            Some(o) => o,
107            None => return Err(Error::OffsetOverflow),
108        };
109        if start_offset >= self.physical_blob_size {
110            return Ok(None); // No more data
111        }
112
113        // Calculate how many pages to read
114        let remaining_physical = (self.physical_blob_size - start_offset) as usize;
115        let max_pages = remaining_physical / self.physical_page_size;
116        let pages_to_read = max_pages.min(self.prefetch_count);
117        if pages_to_read == 0 {
118            return Ok(None);
119        }
120        let bytes_to_read = pages_to_read * self.physical_page_size;
121
122        // Read physical data
123        let physical_buf = self
124            .blob
125            .read_at(start_offset, bytes_to_read, self.read_options)
126            .await?
127            .coalesce()
128            .freeze();
129
130        // Validate CRCs and compute total logical bytes
131        let mut total_logical = 0usize;
132        let mut last_len = 0usize;
133        let is_final_batch = pages_to_read == max_pages;
134        for page_idx in 0..pages_to_read {
135            let page_start = page_idx * self.physical_page_size;
136            let page_slice =
137                &physical_buf.as_ref()[page_start..page_start + self.physical_page_size];
138            let Some(checksum) = Checksum::validate_page(page_slice) else {
139                error!(page = self.blob_page + page_idx as u64, "CRC mismatch");
140                return Err(Error::InvalidChecksum);
141            };
142            let len = checksum.len as usize;
143
144            // Only the final page in the blob may have partial length
145            let is_last_page_in_blob = is_final_batch && page_idx + 1 == pages_to_read;
146            if !is_last_page_in_blob && len != self.page_size {
147                error!(
148                    page = self.blob_page + page_idx as u64,
149                    expected = self.page_size,
150                    actual = len,
151                    "non-last page has partial length"
152                );
153                return Err(Error::InvalidChecksum);
154            }
155
156            let logical_start = (self.blob_page + page_idx as u64)
157                .checked_mul(self.page_size as u64)
158                .ok_or(Error::OffsetOverflow)?;
159            let logical_remaining = self.logical_blob_size.saturating_sub(logical_start);
160            let logical_remaining_in_page = logical_remaining.min(self.page_size as u64) as usize;
161            let exposed_len = len.min(logical_remaining_in_page);
162
163            total_logical += exposed_len;
164            last_len = exposed_len;
165        }
166        self.blob_page += pages_to_read as u64;
167
168        let state = BufferState {
169            buffer: physical_buf,
170            num_pages: pages_to_read,
171            last_page_len: last_len,
172        };
173
174        Ok(Some((state, total_logical)))
175    }
176}
177
178/// Sync buffering component that implements the `Buf` trait.
179///
180/// This accumulates `BufferState` from multiple fills and provides navigation
181/// across pages while skipping CRCs. Consumed buffers are cleaned up in
182/// `advance()`.
183struct ReplayBuf {
184    /// Physical page size (page_size + CHECKSUM_SIZE).
185    physical_page_size: usize,
186    /// Logical page size (data bytes per page, not including CRC).
187    page_size: usize,
188    /// Accumulated buffers from fills.
189    buffers: VecDeque<BufferState>,
190    /// Current page index within the front buffer.
191    current_page: usize,
192    /// Current offset within the current page's logical data.
193    offset_in_page: usize,
194    /// Total remaining logical bytes across all buffers.
195    remaining: usize,
196}
197
198impl ReplayBuf {
199    /// Creates a new ReplayBuf.
200    const fn new(physical_page_size: usize, page_size: usize) -> Self {
201        Self {
202            physical_page_size,
203            page_size,
204            buffers: VecDeque::new(),
205            current_page: 0,
206            offset_in_page: 0,
207            remaining: 0,
208        }
209    }
210
211    /// Clears the buffer and resets the read offset to 0.
212    fn clear(&mut self) {
213        self.buffers.clear();
214        self.current_page = 0;
215        self.offset_in_page = 0;
216        self.remaining = 0;
217    }
218
219    /// Adds a buffer from a fill operation.
220    fn push(&mut self, state: BufferState, logical_bytes: usize) {
221        // If buffers is empty, this is the first fill after a seek.
222        // Skip bytes before the seek offset (offset_in_page).
223        let skip = if self.buffers.is_empty() {
224            self.offset_in_page
225        } else {
226            0
227        };
228        self.buffers.push_back(state);
229        self.remaining += logical_bytes.saturating_sub(skip);
230    }
231
232    /// Returns the logical length of the given page in the given buffer.
233    const fn page_len(buf: &BufferState, page_idx: usize, page_size: usize) -> usize {
234        if page_idx + 1 == buf.num_pages {
235            buf.last_page_len
236        } else {
237            page_size
238        }
239    }
240}
241
242impl Buf for ReplayBuf {
243    fn remaining(&self) -> usize {
244        self.remaining
245    }
246
247    fn chunk(&self) -> &[u8] {
248        let Some(buf) = self.buffers.front() else {
249            return &[];
250        };
251        if self.current_page >= buf.num_pages {
252            return &[];
253        }
254        let page_len = Self::page_len(buf, self.current_page, self.page_size);
255        let physical_start = self.current_page * self.physical_page_size + self.offset_in_page;
256        let physical_end = self.current_page * self.physical_page_size + page_len;
257        &buf.buffer.as_ref()[physical_start..physical_end]
258    }
259
260    fn advance(&mut self, mut cnt: usize) {
261        self.remaining = self.remaining.saturating_sub(cnt);
262
263        while cnt > 0 {
264            let Some(buf) = self.buffers.front() else {
265                break;
266            };
267
268            // Advance within current buffer
269            while cnt > 0 && self.current_page < buf.num_pages {
270                let page_len = Self::page_len(buf, self.current_page, self.page_size);
271                let available = page_len - self.offset_in_page;
272                if cnt < available {
273                    self.offset_in_page += cnt;
274                    return;
275                }
276                cnt -= available;
277                self.current_page += 1;
278                self.offset_in_page = 0;
279            }
280
281            // Current buffer exhausted, move to next
282            if self.current_page >= buf.num_pages {
283                self.buffers.pop_front();
284                self.current_page = 0;
285                self.offset_in_page = 0;
286            }
287        }
288    }
289}
290
291/// Replays logical data from a blob containing pages with interleaved CRCs.
292///
293/// This combines async I/O (`PageReader`) with sync buffering (`ReplayBuf`)
294/// to provide an `ensure(n)` + `Buf` interface for codec decoding.
295pub struct Replay<B: Blob> {
296    /// Async I/O component.
297    reader: PageReader<B>,
298    /// Sync buffering component.
299    buffer: ReplayBuf,
300    /// Whether the blob has been fully read.
301    exhausted: bool,
302}
303
304impl<B: Blob> Replay<B> {
305    /// Creates a new Replay from a PageReader.
306    pub(super) const fn new(reader: PageReader<B>) -> Self {
307        let physical_page_size = reader.physical_page_size();
308        let page_size = reader.page_size();
309        Self {
310            reader,
311            buffer: ReplayBuf::new(physical_page_size, page_size),
312            exhausted: false,
313        }
314    }
315
316    /// Returns the size of the blob.
317    pub const fn blob_size(&self) -> u64 {
318        self.reader.blob_size()
319    }
320
321    /// Returns true if the reader has been exhausted (no more pages to read).
322    ///
323    /// When exhausted, the buffer may still contain data that hasn't been consumed.
324    /// Callers should check `remaining()` to see if there's data left to process.
325    pub const fn is_exhausted(&self) -> bool {
326        self.exhausted
327    }
328
329    /// Ensures at least `n` bytes are available in the buffer.
330    ///
331    /// This method fills the buffer from the blob until either:
332    /// - At least `n` bytes are available (returns `Ok(true)`)
333    /// - The blob is exhausted with fewer than `n` bytes (returns `Ok(false)`)
334    /// - A read error occurs (returns `Err`)
335    ///
336    /// When `Ok(false)` is returned, callers should still attempt to process
337    /// the remaining bytes in the buffer (check `remaining()`), as they may
338    /// contain valid data that doesn't require the full `n` bytes.
339    pub async fn ensure(&mut self, n: usize) -> Result<bool, Error> {
340        while self.buffer.remaining < n && !self.exhausted {
341            match self.reader.fill().await? {
342                Some((state, logical_bytes)) => {
343                    self.buffer.push(state, logical_bytes);
344                }
345                None => {
346                    self.exhausted = true;
347                }
348            }
349        }
350        Ok(self.buffer.remaining >= n)
351    }
352
353    /// Seeks to `offset` in the blob, returning `Err(BlobInsufficientLength)` if `offset` exceeds
354    /// the blob size.
355    pub fn seek_to(&mut self, offset: u64) -> Result<(), Error> {
356        if offset > self.reader.blob_size() {
357            return Err(Error::BlobInsufficientLength);
358        }
359
360        self.buffer.clear();
361        self.exhausted = false;
362
363        let page_size = self.reader.page_size as u64;
364        self.reader.blob_page = offset / page_size;
365        self.buffer.current_page = 0;
366        self.buffer.offset_in_page = (offset % page_size) as usize;
367
368        Ok(())
369    }
370}
371
372impl<B: Blob> Buf for Replay<B> {
373    fn remaining(&self) -> usize {
374        self.buffer.remaining()
375    }
376
377    fn chunk(&self) -> &[u8] {
378        self.buffer.chunk()
379    }
380
381    fn advance(&mut self, cnt: usize) {
382        self.buffer.advance(cnt);
383    }
384}
385
386#[cfg(test)]
387mod tests {
388    use super::{super::writer::Writer, *};
389    use crate::{Runner as _, Storage as _, deterministic};
390    use commonware_macros::test_traced;
391    use commonware_utils::{NZU16, NZUsize};
392
393    const PAGE_SIZE: NonZeroU16 = NZU16!(103);
394    const BUFFER_PAGES: usize = 2;
395
396    #[test_traced("DEBUG")]
397    fn test_replay_basic() {
398        let executor = deterministic::Runner::default();
399        executor.start(|context: deterministic::Context| async move {
400            let (blob, blob_size) = context.open("test_partition", b"test_blob").await.unwrap();
401            assert_eq!(blob_size, 0);
402
403            let cache_ref =
404                super::super::CacheRef::from_pooler(&context, PAGE_SIZE, NZUsize!(BUFFER_PAGES));
405            let mut append = Writer::new(blob.clone(), blob_size, BUFFER_PAGES * 115, cache_ref)
406                .await
407                .unwrap();
408
409            // Write data spanning multiple pages
410            let data: Vec<u8> = (0u8..=255).cycle().take(300).collect();
411            append.append(&data).await.unwrap();
412            append.sync().await.unwrap();
413
414            // Create Replay
415            let mut replay = append
416                .replay(NZUsize!(BUFFER_PAGES), ReadOptions::default())
417                .await
418                .unwrap();
419
420            // Ensure all data is available
421            replay.ensure(300).await.unwrap();
422
423            // Verify we got all the data
424            assert_eq!(replay.remaining(), 300);
425
426            // Read all data via Buf interface
427            let mut collected = Vec::new();
428            while replay.remaining() > 0 {
429                let chunk = replay.chunk();
430                collected.extend_from_slice(chunk);
431                let len = chunk.len();
432                replay.advance(len);
433            }
434            assert_eq!(collected, data);
435        });
436    }
437
438    #[test_traced("DEBUG")]
439    fn test_replay_partial_page() {
440        let executor = deterministic::Runner::default();
441        executor.start(|context: deterministic::Context| async move {
442            let (blob, blob_size) = context.open("test_partition", b"test_blob").await.unwrap();
443
444            let cache_ref =
445                super::super::CacheRef::from_pooler(&context, PAGE_SIZE, NZUsize!(BUFFER_PAGES));
446            let mut append = Writer::new(blob.clone(), blob_size, BUFFER_PAGES * 115, cache_ref)
447                .await
448                .unwrap();
449
450            // Write data that doesn't fill the last page
451            let data: Vec<u8> = (1u8..=(PAGE_SIZE.get() + 10) as u8).collect();
452            append.append(&data).await.unwrap();
453            append.sync().await.unwrap();
454
455            let mut replay = append
456                .replay(NZUsize!(BUFFER_PAGES), ReadOptions::default())
457                .await
458                .unwrap();
459
460            // Ensure all data is available
461            replay.ensure(data.len()).await.unwrap();
462
463            assert_eq!(replay.remaining(), data.len());
464        });
465    }
466
467    #[test_traced("DEBUG")]
468    fn test_replay_cross_buffer_boundary() {
469        // Use prefetch_count=1 to force separate BufferStates per page.
470        // This tests navigation across multiple BufferStates in the VecDeque.
471        let executor = deterministic::Runner::default();
472        executor.start(|context: deterministic::Context| async move {
473            let (blob, blob_size) = context.open("test_partition", b"test_blob").await.unwrap();
474            assert_eq!(blob_size, 0);
475
476            let cache_ref =
477                super::super::CacheRef::from_pooler(&context, PAGE_SIZE, NZUsize!(BUFFER_PAGES));
478            let mut append = Writer::new(blob.clone(), blob_size, BUFFER_PAGES * 115, cache_ref)
479                .await
480                .unwrap();
481
482            // Write data spanning 4 pages (4 * 103 = 412 bytes, with last page partial)
483            let data: Vec<u8> = (0u8..=255).cycle().take(400).collect();
484            append.append(&data).await.unwrap();
485            append.sync().await.unwrap();
486
487            // Create Replay with buffer size that results in prefetch_count=1.
488            // Physical page size = 103 + 12 = 115 bytes.
489            // Buffer size of 115 gives prefetch_pages = 115/115 = 1.
490            let mut replay = append
491                .replay(NZUsize!(115), ReadOptions::default())
492                .await
493                .unwrap();
494
495            // Ensure all data - this requires 4 separate fill() calls (one per page).
496            // Each fill() creates a new BufferState, so we'll have 4 BufferStates.
497            assert!(replay.ensure(400).await.unwrap());
498            assert_eq!(replay.remaining(), 400);
499
500            // Read all data via Buf interface, verifying navigation across BufferStates.
501            let mut collected = Vec::new();
502            let mut chunks_read = 0;
503            while replay.remaining() > 0 {
504                let chunk = replay.chunk();
505                assert!(
506                    !chunk.is_empty(),
507                    "chunk() returned empty but remaining > 0"
508                );
509                collected.extend_from_slice(chunk);
510                let len = chunk.len();
511                replay.advance(len);
512                chunks_read += 1;
513            }
514
515            assert_eq!(collected, data);
516            // With prefetch_count=1 and 4 pages, we expect at least 4 chunks
517            // (one per page, though partial reads could result in more).
518            assert!(
519                chunks_read >= 4,
520                "Expected at least 4 chunks for 4 pages, got {}",
521                chunks_read
522            );
523        });
524    }
525
526    #[test_traced("DEBUG")]
527    fn test_replay_empty_blob() {
528        // Test that replaying an empty blob works correctly.
529        // ensure() should return Ok(false) when no data is available.
530        let executor = deterministic::Runner::default();
531        executor.start(|context: deterministic::Context| async move {
532            let (blob, blob_size) = context.open("test_partition", b"test_blob").await.unwrap();
533            assert_eq!(blob_size, 0);
534
535            let cache_ref =
536                super::super::CacheRef::from_pooler(&context, PAGE_SIZE, NZUsize!(BUFFER_PAGES));
537            let mut append = Writer::new(blob.clone(), blob_size, BUFFER_PAGES * 115, cache_ref)
538                .await
539                .unwrap();
540
541            // Don't write any data - blob remains empty
542            assert_eq!(append.size(), 0);
543
544            // Create Replay on empty blob
545            let mut replay = append
546                .replay(NZUsize!(BUFFER_PAGES), ReadOptions::default())
547                .await
548                .unwrap();
549
550            // Verify initial state - remaining is 0, but not yet marked exhausted
551            // (exhausted is set after first fill attempt)
552            assert_eq!(replay.remaining(), 0);
553
554            // ensure(0) should succeed (we have >= 0 bytes)
555            assert!(replay.ensure(0).await.unwrap());
556
557            // ensure(1) should return Ok(false) - not enough data, and marks exhausted
558            assert!(!replay.ensure(1).await.unwrap());
559
560            // Now should be marked as exhausted after the fill attempt
561            assert!(replay.is_exhausted());
562
563            // chunk() should return empty slice
564            assert!(replay.chunk().is_empty());
565
566            // remaining should still be 0
567            assert_eq!(replay.remaining(), 0);
568        });
569    }
570
571    #[test_traced("DEBUG")]
572    fn test_replay_seek_to() {
573        let executor = deterministic::Runner::default();
574        executor.start(|context: deterministic::Context| async move {
575            let (blob, blob_size) = context.open("test_partition", b"test_blob").await.unwrap();
576
577            let cache_ref =
578                super::super::CacheRef::from_pooler(&context, PAGE_SIZE, NZUsize!(BUFFER_PAGES));
579            let mut append = Writer::new(blob.clone(), blob_size, BUFFER_PAGES * 115, cache_ref)
580                .await
581                .unwrap();
582
583            // Write data spanning multiple pages
584            let data: Vec<u8> = (0u8..=255).cycle().take(300).collect();
585            append.append(&data).await.unwrap();
586            append.sync().await.unwrap();
587
588            let mut replay = append
589                .replay(NZUsize!(BUFFER_PAGES), ReadOptions::default())
590                .await
591                .unwrap();
592
593            // Seek forward, read, then seek backward
594            replay.seek_to(150).unwrap();
595            replay.ensure(50).await.unwrap();
596            assert_eq!(replay.get_u8(), data[150]);
597
598            // Seek back to start
599            replay.seek_to(0).unwrap();
600            replay.ensure(1).await.unwrap();
601            assert_eq!(replay.get_u8(), data[0]);
602
603            // Seek beyond blob size should error
604            assert!(replay.seek_to(data.len() as u64 + 1).is_err());
605
606            // Test that remaining() is correct after seek by reading all data.
607            let seek_offset = 150usize;
608            replay.seek_to(seek_offset as u64).unwrap();
609            let expected_remaining = data.len() - seek_offset;
610            // Read all bytes and verify content
611            let mut collected = Vec::new();
612            loop {
613                // Load more data if needed
614                if !replay.ensure(1).await.unwrap() {
615                    break; // No more data available
616                }
617                let chunk = replay.chunk();
618                if chunk.is_empty() {
619                    break;
620                }
621                collected.extend_from_slice(chunk);
622                let len = chunk.len();
623                replay.advance(len);
624            }
625            assert_eq!(
626                collected.len(),
627                expected_remaining,
628                "After seeking to {}, should read {} bytes but got {}",
629                seek_offset,
630                expected_remaining,
631                collected.len()
632            );
633            assert_eq!(collected, &data[seek_offset..]);
634        });
635    }
636}