commonware_runtime/utils/buffer/paged/
read.rs1use super::Checksum;
2use crate::{Blob, Buf, Error, IoBuf, ReadOptions};
3use commonware_codec::FixedSize;
4use std::{collections::VecDeque, num::NonZeroU16};
5use tracing::error;
6
7pub(super) struct BufferState {
13 buffer: IoBuf,
15 num_pages: usize,
17 last_page_len: usize,
19}
20
21pub(super) struct PageReader<B: Blob> {
26 blob: B,
28 physical_page_size: usize,
30 page_size: usize,
32 physical_blob_size: u64,
34 logical_blob_size: u64,
36 blob_page: u64,
38 prefetch_count: usize,
40 read_options: ReadOptions,
42}
43
44impl<B: Blob> PageReader<B> {
45 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 pub(super) const fn blob_size(&self) -> u64 {
86 self.logical_blob_size
87 }
88
89 pub(super) const fn physical_page_size(&self) -> usize {
91 self.physical_page_size
92 }
93
94 pub(super) const fn page_size(&self) -> usize {
96 self.page_size
97 }
98
99 pub(super) async fn fill(&mut self) -> Result<Option<(BufferState, usize)>, Error> {
104 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); }
112
113 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 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 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 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
178struct ReplayBuf {
184 physical_page_size: usize,
186 page_size: usize,
188 buffers: VecDeque<BufferState>,
190 current_page: usize,
192 offset_in_page: usize,
194 remaining: usize,
196}
197
198impl ReplayBuf {
199 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 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 fn push(&mut self, state: BufferState, logical_bytes: usize) {
221 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 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 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 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
291pub struct Replay<B: Blob> {
296 reader: PageReader<B>,
298 buffer: ReplayBuf,
300 exhausted: bool,
302}
303
304impl<B: Blob> Replay<B> {
305 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 pub const fn blob_size(&self) -> u64 {
318 self.reader.blob_size()
319 }
320
321 pub const fn is_exhausted(&self) -> bool {
326 self.exhausted
327 }
328
329 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 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 let data: Vec<u8> = (0u8..=255).cycle().take(300).collect();
411 append.append(&data).await.unwrap();
412 append.sync().await.unwrap();
413
414 let mut replay = append
416 .replay(NZUsize!(BUFFER_PAGES), ReadOptions::default())
417 .await
418 .unwrap();
419
420 replay.ensure(300).await.unwrap();
422
423 assert_eq!(replay.remaining(), 300);
425
426 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 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 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 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 let data: Vec<u8> = (0u8..=255).cycle().take(400).collect();
484 append.append(&data).await.unwrap();
485 append.sync().await.unwrap();
486
487 let mut replay = append
491 .replay(NZUsize!(115), ReadOptions::default())
492 .await
493 .unwrap();
494
495 assert!(replay.ensure(400).await.unwrap());
498 assert_eq!(replay.remaining(), 400);
499
500 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 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 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 assert_eq!(append.size(), 0);
543
544 let mut replay = append
546 .replay(NZUsize!(BUFFER_PAGES), ReadOptions::default())
547 .await
548 .unwrap();
549
550 assert_eq!(replay.remaining(), 0);
553
554 assert!(replay.ensure(0).await.unwrap());
556
557 assert!(!replay.ensure(1).await.unwrap());
559
560 assert!(replay.is_exhausted());
562
563 assert!(replay.chunk().is_empty());
565
566 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 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 replay.seek_to(150).unwrap();
595 replay.ensure(50).await.unwrap();
596 assert_eq!(replay.get_u8(), data[150]);
597
598 replay.seek_to(0).unwrap();
600 replay.ensure(1).await.unwrap();
601 assert_eq!(replay.get_u8(), data[0]);
602
603 assert!(replay.seek_to(data.len() as u64 + 1).is_err());
605
606 let seek_offset = 150usize;
608 replay.seek_to(seek_offset as u64).unwrap();
609 let expected_remaining = data.len() - seek_offset;
610 let mut collected = Vec::new();
612 loop {
613 if !replay.ensure(1).await.unwrap() {
615 break; }
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}