1use std::collections::{HashMap, HashSet};
15
16use fsqlite_error::{FrankenError, Result};
17use tracing::{debug, error, info, warn};
18
19use crate::replication_sender::{
20 CHANGESET_HEADER_SIZE, ChangesetId, PageEntry, ReplicationPacket, ReplicationPacketV2Header,
21 SenderConfig, compute_changeset_id, derive_seed_from_changeset_id, encode_changeset,
22};
23use crate::source_block_partition::{K_MAX, SourceBlock, partition_source_blocks};
24
25const BEAD_ID: &str = "bd-1hi.15";
26
27#[derive(Debug, Clone)]
33pub struct BlockResumeState {
34 pub block_id: u32,
36 pub num_received: u32,
38 pub received_isis: HashSet<u32>,
40 pub decoded: bool,
42}
43
44impl BlockResumeState {
45 #[must_use]
47 fn new(block_id: u32) -> Self {
48 Self {
49 block_id,
50 num_received: 0,
51 received_isis: HashSet::new(),
52 decoded: false,
53 }
54 }
55
56 fn record_isi(&mut self, isi: u32) -> bool {
58 if self.received_isis.insert(isi) {
59 self.num_received += 1;
60 true
61 } else {
62 false
63 }
64 }
65
66 #[must_use]
70 pub fn to_bytes(&self) -> Vec<u8> {
71 let n = self.received_isis.len();
72 let mut buf = Vec::with_capacity(13 + n * 4);
73 buf.extend_from_slice(&self.block_id.to_le_bytes());
74 buf.extend_from_slice(&self.num_received.to_le_bytes());
75 buf.push(u8::from(self.decoded));
76 let n_u32 = u32::try_from(n).unwrap_or(u32::MAX);
77 buf.extend_from_slice(&n_u32.to_le_bytes());
78 let mut sorted_isis: Vec<u32> = self.received_isis.iter().copied().collect();
79 sorted_isis.sort_unstable();
80 for isi in sorted_isis {
81 buf.extend_from_slice(&isi.to_le_bytes());
82 }
83 buf
84 }
85
86 pub fn from_bytes(buf: &[u8]) -> Result<(Self, usize)> {
92 if buf.len() < 13 {
93 return Err(FrankenError::DatabaseCorrupt {
94 detail: format!("BlockResumeState too short: {} < 13", buf.len()),
95 });
96 }
97 let block_id = u32::from_le_bytes(buf[0..4].try_into().expect("4 bytes"));
98 let num_received = u32::from_le_bytes(buf[4..8].try_into().expect("4 bytes"));
99 let decoded = buf[8] != 0;
100 let n_isis = u32::from_le_bytes(buf[9..13].try_into().expect("4 bytes"));
101 let n = n_isis as usize;
102 let expected = n
103 .checked_mul(4)
104 .and_then(|v| v.checked_add(13))
105 .ok_or_else(|| FrankenError::DatabaseCorrupt {
106 detail: format!("BlockResumeState n_isis ({n_isis}) causes size overflow"),
107 })?;
108 if buf.len() < expected {
109 return Err(FrankenError::DatabaseCorrupt {
110 detail: format!("BlockResumeState truncated: {} < {expected}", buf.len()),
111 });
112 }
113 let mut received_isis = HashSet::with_capacity(n);
114 for i in 0..n {
115 let offset = 13 + i * 4;
116 let isi = u32::from_le_bytes(buf[offset..offset + 4].try_into().expect("4 bytes"));
117 received_isis.insert(isi);
118 }
119 Ok((
120 Self {
121 block_id,
122 num_received,
123 received_isis,
124 decoded,
125 },
126 expected,
127 ))
128 }
129}
130
131#[derive(Debug, Clone)]
133pub struct ResumeState {
134 pub blocks: Vec<BlockResumeState>,
136 pub total_blocks: u32,
138}
139
140impl ResumeState {
141 #[must_use]
143 pub fn new(total_blocks: u32) -> Self {
144 let blocks = (0..total_blocks).map(BlockResumeState::new).collect();
145 Self {
146 blocks,
147 total_blocks,
148 }
149 }
150
151 #[must_use]
153 pub fn decoded_count(&self) -> u32 {
154 u32::try_from(self.blocks.iter().filter(|b| b.decoded).count()).unwrap_or(u32::MAX)
155 }
156
157 #[must_use]
159 pub fn all_decoded(&self) -> bool {
160 self.blocks.iter().all(|b| b.decoded)
161 }
162
163 #[must_use]
165 pub fn to_bytes(&self) -> Vec<u8> {
166 let mut buf = Vec::new();
167 buf.extend_from_slice(&self.total_blocks.to_le_bytes());
168 for block in &self.blocks {
169 buf.extend_from_slice(&block.to_bytes());
170 }
171 buf
172 }
173
174 pub fn from_bytes(buf: &[u8]) -> Result<Self> {
180 if buf.len() < 4 {
181 return Err(FrankenError::DatabaseCorrupt {
182 detail: format!("ResumeState too short: {} < 4", buf.len()),
183 });
184 }
185 let total_blocks = u32::from_le_bytes(buf[0..4].try_into().expect("4 bytes"));
186 let mut blocks = Vec::with_capacity(total_blocks as usize);
187 let mut offset = 4;
188 for _ in 0..total_blocks {
189 let (block, consumed) = BlockResumeState::from_bytes(&buf[offset..])?;
190 blocks.push(block);
191 offset += consumed;
192 }
193 Ok(Self {
194 blocks,
195 total_blocks,
196 })
197 }
198}
199
200#[derive(Debug)]
206pub struct SnapshotSender {
207 pub source_blocks: Vec<SourceBlock>,
209 pub page_size: u32,
211 current_block: usize,
213 current_isi: u32,
215 block_changeset_ids: Vec<ChangesetId>,
217 block_k_sources: Vec<u32>,
219 block_changesets: Vec<Vec<u8>>,
221 config: SenderConfig,
223 done: bool,
225}
226
227impl SnapshotSender {
228 #[allow(clippy::too_many_lines)]
236 pub fn prepare(
237 page_size: u32,
238 all_pages: &mut [PageEntry],
239 config: SenderConfig,
240 ) -> Result<Self> {
241 if all_pages.is_empty() {
242 return Err(FrankenError::OutOfRange {
243 what: "pages".to_owned(),
244 value: "0".to_owned(),
245 });
246 }
247
248 let total_pages = u32::try_from(all_pages.len()).map_err(|_| FrankenError::OutOfRange {
249 what: "total_pages".to_owned(),
250 value: all_pages.len().to_string(),
251 })?;
252
253 let source_blocks = partition_source_blocks(total_pages)?;
254 info!(
255 bead_id = BEAD_ID,
256 total_pages,
257 n_blocks = source_blocks.len(),
258 page_size,
259 "snapshot partitioned into source blocks"
260 );
261
262 all_pages.sort_by_key(|p| p.page_number);
264
265 let mut block_changeset_ids = Vec::with_capacity(source_blocks.len());
267 let mut block_k_sources = Vec::with_capacity(source_blocks.len());
268 let mut block_changesets = Vec::with_capacity(source_blocks.len());
269
270 let mut page_idx = 0_usize;
271 for block in &source_blocks {
272 let end = page_idx + block.num_pages as usize;
273 if end > all_pages.len() {
274 return Err(FrankenError::Internal(format!(
275 "block {} requires pages up to index {end}, but only {} available",
276 block.index,
277 all_pages.len()
278 )));
279 }
280 let block_pages = &mut all_pages[page_idx..end];
281 let changeset_bytes = encode_changeset(page_size, block_pages)?;
282 let changeset_id = compute_changeset_id(&changeset_bytes);
283
284 let t = u64::from(config.symbol_size);
286 let f = changeset_bytes.len() as u64;
287 let k_source = u32::try_from(f.div_ceil(t)).map_err(|_| FrankenError::OutOfRange {
288 what: "k_source".to_owned(),
289 value: f.div_ceil(t).to_string(),
290 })?;
291
292 debug!(
293 bead_id = BEAD_ID,
294 block_index = block.index,
295 num_pages = block.num_pages,
296 changeset_len = changeset_bytes.len(),
297 k_source,
298 "prepared block changeset"
299 );
300
301 block_changeset_ids.push(changeset_id);
302 block_k_sources.push(k_source);
303 block_changesets.push(changeset_bytes);
304 page_idx = end;
305 }
306
307 Ok(Self {
308 source_blocks,
309 page_size,
310 current_block: 0,
311 current_isi: 0,
312 block_changeset_ids,
313 block_k_sources,
314 block_changesets,
315 config,
316 done: false,
317 })
318 }
319
320 pub fn next_packet(&mut self) -> Option<ReplicationPacket> {
325 if self.done || self.current_block >= self.source_blocks.len() {
326 self.done = true;
327 return None;
328 }
329
330 let k_source = self.block_k_sources[self.current_block];
331 let max_isi = k_source.saturating_mul(self.config.max_isi_multiplier);
332
333 if self.current_isi >= max_isi {
334 self.current_block += 1;
335 self.current_isi = 0;
336 if self.current_block >= self.source_blocks.len() {
337 self.done = true;
338 return None;
339 }
340 }
341
342 let changeset = &self.block_changesets[self.current_block];
343 let changeset_id = self.block_changeset_ids[self.current_block];
344 let k_source = self.block_k_sources[self.current_block];
345 let isi = self.current_isi;
346 let t = usize::from(self.config.symbol_size);
347
348 let symbol_data = if u64::from(isi) < u64::from(k_source) {
350 let start = isi as usize * t;
351 let end = (start + t).min(changeset.len());
352 let mut data = vec![0_u8; t];
353 let available = end.saturating_sub(start);
354 if available > 0 {
355 data[..available].copy_from_slice(&changeset[start..end]);
356 }
357 data
358 } else {
359 let seed = derive_seed_from_changeset_id(&changeset_id);
360 #[cfg(not(target_arch = "wasm32"))]
361 {
362 use asupersync::raptorq::systematic::SystematicEncoder;
366
367 let source_symbols: Vec<Vec<u8>> = (0..k_source as usize)
368 .map(|i| {
369 let start = i * t;
370 let end = (start + t).min(changeset.len());
371 let mut sym = vec![0_u8; t];
372 let available = end.saturating_sub(start);
373 if available > 0 {
374 sym[..available].copy_from_slice(&changeset[start..end]);
375 }
376 sym
377 })
378 .collect();
379
380 match SystematicEncoder::new(&source_symbols, t, seed) {
381 Some(encoder) => {
382 let repair_esi = isi - k_source;
383 encoder.repair_symbol(repair_esi)
384 }
385 None => {
386 crate::replication_sender::generate_deterministic_placeholder(seed, isi, t)
388 }
389 }
390 }
391 #[cfg(target_arch = "wasm32")]
392 {
393 crate::replication_sender::generate_deterministic_placeholder(seed, isi, t)
394 }
395 };
396
397 let seed = derive_seed_from_changeset_id(&changeset_id);
398 let r_repair = max_isi.saturating_sub(k_source);
399 let packet = ReplicationPacket::new_v2(
400 ReplicationPacketV2Header {
401 changeset_id,
402 sbn: 0,
403 esi: isi,
404 k_source,
405 r_repair,
406 symbol_size_t: self.config.symbol_size,
407 seed,
408 },
409 symbol_data,
410 );
411
412 self.current_isi += 1;
413 Some(packet)
414 }
415
416 #[must_use]
418 pub fn num_blocks(&self) -> usize {
419 self.source_blocks.len()
420 }
421
422 #[must_use]
424 pub fn total_source_symbols(&self) -> u64 {
425 self.block_k_sources.iter().map(|&k| u64::from(k)).sum()
426 }
427
428 pub fn restart(&mut self) {
430 self.current_block = 0;
431 self.current_isi = 0;
432 self.done = false;
433 debug!(bead_id = BEAD_ID, "snapshot sender restarted for next pass");
434 }
435}
436
437#[derive(Debug, Clone, Copy, PartialEq, Eq)]
443pub enum SnapshotReceiverState {
444 Waiting,
446 Receiving,
448 Complete,
450}
451
452#[derive(Debug, Clone)]
454pub struct DecodedBlock {
455 pub block_index: u32,
457 pub pages: Vec<DecodedBlockPage>,
459}
460
461#[derive(Debug, Clone, PartialEq, Eq)]
463pub struct DecodedBlockPage {
464 pub page_number: u32,
466 pub page_data: Vec<u8>,
468}
469
470#[derive(Debug)]
472struct BlockDecoder {
473 changeset_id: Option<ChangesetId>,
475 k_source: u32,
477 symbol_size: u32,
479 seed: u64,
481 symbols: HashMap<u32, Vec<u8>>,
483 received_isis: HashSet<u32>,
485 decoded: bool,
487}
488
489impl BlockDecoder {
490 fn new() -> Self {
491 Self {
492 changeset_id: None,
493 k_source: 0,
494 symbol_size: 0,
495 seed: 0,
496 symbols: HashMap::new(),
497 received_isis: HashSet::new(),
498 decoded: false,
499 }
500 }
501
502 fn initialize(&mut self, changeset_id: ChangesetId, k_source: u32, symbol_size: u32) {
503 self.changeset_id = Some(changeset_id);
504 self.k_source = k_source;
505 self.symbol_size = symbol_size;
506 self.seed = derive_seed_from_changeset_id(&changeset_id);
507 }
508
509 fn add_symbol(&mut self, isi: u32, data: Vec<u8>) -> bool {
510 if self.received_isis.insert(isi) {
511 self.symbols.insert(isi, data);
512 true
513 } else {
514 false
515 }
516 }
517
518 fn received_count(&self) -> u32 {
519 u32::try_from(self.received_isis.len()).unwrap_or(u32::MAX)
520 }
521
522 fn ready_to_decode(&self) -> bool {
523 self.received_count() >= self.k_source && self.k_source > 0
524 }
525
526 fn try_decode(&self) -> Option<Vec<u8>> {
527 if !self.ready_to_decode() {
528 return None;
529 }
530 let source_count = self
531 .symbols
532 .keys()
533 .filter(|&&isi| isi < self.k_source)
534 .count();
535 let k = self.k_source as usize;
536 let t = self.symbol_size as usize;
537 if source_count >= k {
538 let padded_len = k * t;
539 let mut padded = vec![0_u8; padded_len];
540 for isi in 0..self.k_source {
541 if let Some(data) = self.symbols.get(&isi) {
542 let start = isi as usize * t;
543 let copy_len = data.len().min(t);
544 padded[start..start + copy_len].copy_from_slice(&data[..copy_len]);
545 }
546 }
547 Some(padded)
548 } else {
549 warn!(
550 bead_id = BEAD_ID,
551 source_count,
552 k_source = self.k_source,
553 "snapshot block decode needs repair symbols (production RaptorQ)"
554 );
555 None
556 }
557 }
558}
559
560#[derive(Debug)]
562pub struct SnapshotReceiver {
563 state: SnapshotReceiverState,
564 changeset_to_block: HashMap<ChangesetId, usize>,
566 block_decoders: Vec<BlockDecoder>,
568 num_blocks: usize,
570 decoded_blocks: Vec<DecodedBlock>,
572 resume: ResumeState,
574 page_size: u32,
576}
577
578impl SnapshotReceiver {
579 #[must_use]
584 pub fn new(num_blocks: usize, page_size: u32) -> Self {
585 let block_decoders = (0..num_blocks).map(|_| BlockDecoder::new()).collect();
586 Self {
587 state: SnapshotReceiverState::Waiting,
588 changeset_to_block: HashMap::new(),
589 block_decoders,
590 num_blocks,
591 decoded_blocks: Vec::new(),
592 resume: ResumeState::new(u32::try_from(num_blocks).unwrap_or(u32::MAX)),
593 page_size,
594 }
595 }
596
597 #[must_use]
599 pub fn from_resume(resume: ResumeState, page_size: u32) -> Self {
600 let num_blocks = resume.total_blocks as usize;
601 let block_decoders = (0..num_blocks).map(|_| BlockDecoder::new()).collect();
602 Self {
603 state: if resume.all_decoded() {
604 SnapshotReceiverState::Complete
605 } else {
606 SnapshotReceiverState::Waiting
607 },
608 changeset_to_block: HashMap::new(),
609 block_decoders,
610 num_blocks,
611 decoded_blocks: Vec::new(),
612 resume,
613 page_size,
614 }
615 }
616
617 #[must_use]
619 pub const fn state(&self) -> SnapshotReceiverState {
620 self.state
621 }
622
623 #[must_use]
625 pub fn blocks_decoded(&self) -> usize {
626 self.decoded_blocks.len()
627 }
628
629 #[must_use]
631 pub fn resume_state(&self) -> &ResumeState {
632 &self.resume
633 }
634
635 pub fn take_decoded_blocks(&mut self) -> Vec<DecodedBlock> {
637 std::mem::take(&mut self.decoded_blocks)
638 }
639
640 #[allow(clippy::too_many_lines)]
649 pub fn process_packet(&mut self, packet: &ReplicationPacket) -> Result<SnapshotPacketResult> {
650 if self.state == SnapshotReceiverState::Complete {
651 return Ok(SnapshotPacketResult::AlreadyComplete);
652 }
653
654 if packet.sbn != 0 {
656 return Err(FrankenError::Internal(format!(
657 "V1: SBN must be 0, got {}",
658 packet.sbn
659 )));
660 }
661 if packet.k_source == 0 || packet.k_source > K_MAX {
662 return Err(FrankenError::OutOfRange {
663 what: "k_source".to_owned(),
664 value: packet.k_source.to_string(),
665 });
666 }
667 let symbol_size =
668 u32::try_from(packet.symbol_data.len()).map_err(|_| FrankenError::OutOfRange {
669 what: "symbol_size".to_owned(),
670 value: packet.symbol_data.len().to_string(),
671 })?;
672 if symbol_size == 0 {
673 return Err(FrankenError::OutOfRange {
674 what: "symbol_size".to_owned(),
675 value: "0".to_owned(),
676 });
677 }
678
679 if self.state == SnapshotReceiverState::Waiting {
680 self.state = SnapshotReceiverState::Receiving;
681 info!(bead_id = BEAD_ID, "snapshot receiving started");
682 }
683
684 let changeset_id = packet.changeset_id;
685
686 let block_idx = if let Some(&idx) = self.changeset_to_block.get(&changeset_id) {
688 idx
689 } else {
690 let next_idx = self
692 .block_decoders
693 .iter()
694 .position(|d| d.changeset_id.is_none() && !d.decoded);
695 if let Some(idx) = next_idx {
696 self.changeset_to_block.insert(changeset_id, idx);
697 self.block_decoders[idx].initialize(changeset_id, packet.k_source, symbol_size);
698 debug!(
699 bead_id = BEAD_ID,
700 block_index = idx,
701 k_source = packet.k_source,
702 "mapped new changeset to block"
703 );
704 idx
705 } else {
706 warn!(
707 bead_id = BEAD_ID,
708 "no available block slot for new changeset_id"
709 );
710 return Ok(SnapshotPacketResult::Rejected);
711 }
712 };
713
714 if block_idx >= self.block_decoders.len() {
715 return Ok(SnapshotPacketResult::Rejected);
716 }
717
718 let decoder = &mut self.block_decoders[block_idx];
719 if decoder.decoded {
720 return Ok(SnapshotPacketResult::BlockAlreadyDecoded);
721 }
722
723 if decoder.k_source != packet.k_source {
725 return Err(FrankenError::DatabaseCorrupt {
726 detail: format!(
727 "k_source mismatch for block {block_idx}: {} vs {}",
728 decoder.k_source, packet.k_source
729 ),
730 });
731 }
732 if decoder.symbol_size != symbol_size {
733 return Err(FrankenError::DatabaseCorrupt {
734 detail: format!(
735 "symbol_size mismatch for block {block_idx}: {} vs {symbol_size}",
736 decoder.symbol_size
737 ),
738 });
739 }
740
741 let accepted = decoder.add_symbol(packet.esi, packet.symbol_data.clone());
743 if !accepted {
744 return Ok(SnapshotPacketResult::Duplicate);
745 }
746
747 if block_idx < self.resume.blocks.len() {
749 self.resume.blocks[block_idx].record_isi(packet.esi);
750 }
751
752 if decoder.ready_to_decode()
754 && !decoder.decoded
755 && let Some(padded) = decoder.try_decode()
756 {
757 match parse_decoded_snapshot_block(&padded, self.page_size) {
758 Ok(pages) => {
759 let block_id = u32::try_from(block_idx).unwrap_or(u32::MAX);
760 decoder.decoded = true;
761 if block_idx < self.resume.blocks.len() {
762 self.resume.blocks[block_idx].decoded = true;
763 }
764 let n_pages = pages.len();
765 self.decoded_blocks.push(DecodedBlock {
766 block_index: block_id,
767 pages,
768 });
769 info!(
770 bead_id = BEAD_ID,
771 block_index = block_idx,
772 n_pages,
773 decoded_so_far = self.decoded_blocks.len(),
774 total_blocks = self.num_blocks,
775 "source block decoded (progressive)"
776 );
777
778 if self.block_decoders.iter().all(|d| d.decoded) {
780 self.state = SnapshotReceiverState::Complete;
781 info!(
782 bead_id = BEAD_ID,
783 total_blocks = self.num_blocks,
784 "snapshot fully received"
785 );
786 }
787 return Ok(SnapshotPacketResult::BlockDecoded(block_id));
788 }
789 Err(e) => {
790 error!(
791 bead_id = BEAD_ID,
792 block_index = block_idx,
793 error = %e,
794 "snapshot block validation failed"
795 );
796 return Err(e);
797 }
798 }
799 }
800
801 Ok(SnapshotPacketResult::Accepted)
802 }
803}
804
805#[derive(Debug, Clone, Copy, PartialEq, Eq)]
807pub enum SnapshotPacketResult {
808 Accepted,
810 Duplicate,
812 BlockDecoded(u32),
814 BlockAlreadyDecoded,
816 Rejected,
818 AlreadyComplete,
820}
821
822fn parse_decoded_snapshot_block(
828 padded_bytes: &[u8],
829 _page_size: u32,
830) -> Result<Vec<DecodedBlockPage>> {
831 use crate::replication_sender::ChangesetHeader;
832
833 if padded_bytes.len() < CHANGESET_HEADER_SIZE {
834 return Err(FrankenError::DatabaseCorrupt {
835 detail: format!(
836 "decoded block too short for header: {} < {CHANGESET_HEADER_SIZE}",
837 padded_bytes.len()
838 ),
839 });
840 }
841
842 let header_bytes: [u8; CHANGESET_HEADER_SIZE] = padded_bytes[..CHANGESET_HEADER_SIZE]
843 .try_into()
844 .expect("checked length");
845 let header = ChangesetHeader::from_bytes(&header_bytes)?;
846
847 let total_len = usize::try_from(header.total_len).map_err(|_| FrankenError::OutOfRange {
848 what: "total_len".to_owned(),
849 value: header.total_len.to_string(),
850 })?;
851 if total_len > padded_bytes.len() {
852 return Err(FrankenError::DatabaseCorrupt {
853 detail: format!(
854 "total_len ({total_len}) exceeds decoded bytes ({})",
855 padded_bytes.len()
856 ),
857 });
858 }
859 let changeset_bytes = &padded_bytes[..total_len];
860
861 let entry_size = 4_usize + 8 + header.page_size as usize;
862 let data_bytes = &changeset_bytes[CHANGESET_HEADER_SIZE..];
863
864 let required_data_len = (header.n_pages as usize)
865 .checked_mul(entry_size)
866 .ok_or_else(|| FrankenError::DatabaseCorrupt {
867 detail: "n_pages causes size overflow".to_owned(),
868 })?;
869
870 if data_bytes.len() < required_data_len {
871 return Err(FrankenError::DatabaseCorrupt {
872 detail: format!(
873 "changeset truncated: expected {} data bytes, got {}",
874 required_data_len,
875 data_bytes.len()
876 ),
877 });
878 }
879
880 let mut pages = Vec::with_capacity(header.n_pages as usize);
881 for i in 0..header.n_pages as usize {
882 let offset = i * entry_size;
883 let page_number =
884 u32::from_le_bytes(data_bytes[offset..offset + 4].try_into().expect("4 bytes"));
885 let page_xxh3 = u64::from_le_bytes(
886 data_bytes[offset + 4..offset + 12]
887 .try_into()
888 .expect("8 bytes"),
889 );
890 let page_data = data_bytes[offset + 12..offset + 12 + header.page_size as usize].to_vec();
891
892 let computed_xxh3 = xxhash_rust::xxh3::xxh3_64(&page_data);
893 if computed_xxh3 != page_xxh3 {
894 error!(
895 bead_id = BEAD_ID,
896 page_number,
897 expected_xxh3 = page_xxh3,
898 computed_xxh3,
899 "snapshot page xxh3 mismatch"
900 );
901 return Err(FrankenError::DatabaseCorrupt {
902 detail: format!(
903 "snapshot page {page_number} xxh3 mismatch: {page_xxh3:#x} vs {computed_xxh3:#x}"
904 ),
905 });
906 }
907
908 pages.push(DecodedBlockPage {
909 page_number,
910 page_data,
911 });
912 }
913
914 Ok(pages)
915}
916
917#[cfg(test)]
918mod tests {
919 use super::*;
920 use crate::replication_sender::PageEntry;
921
922 const TEST_BEAD_ID: &str = "bd-1hi.15";
923
924 #[allow(clippy::cast_possible_truncation)]
925 fn make_pages(page_size: u32, page_numbers: &[u32]) -> Vec<PageEntry> {
926 page_numbers
927 .iter()
928 .map(|&pn| {
929 let mut data = vec![0_u8; page_size as usize];
930 for (i, byte) in data.iter_mut().enumerate() {
931 *byte = ((pn as usize * 251 + i * 31) % 256) as u8;
932 }
933 PageEntry::new(pn, data)
934 })
935 .collect()
936 }
937
938 #[test]
943 fn test_resume_state_persistence() {
944 let mut resume = ResumeState::new(3);
945 resume.blocks[0].record_isi(0);
946 resume.blocks[0].record_isi(5);
947 resume.blocks[0].record_isi(10);
948 resume.blocks[1].decoded = true;
949
950 let bytes = resume.to_bytes();
951 let restored = ResumeState::from_bytes(&bytes).expect("deserialize");
952
953 assert_eq!(
954 restored.total_blocks, 3,
955 "bead_id={TEST_BEAD_ID} case=resume_total_blocks"
956 );
957 assert_eq!(
958 restored.blocks[0].num_received, 3,
959 "bead_id={TEST_BEAD_ID} case=resume_block0_received"
960 );
961 assert!(
962 restored.blocks[0].received_isis.contains(&5),
963 "bead_id={TEST_BEAD_ID} case=resume_block0_isi_5"
964 );
965 assert!(
966 restored.blocks[1].decoded,
967 "bead_id={TEST_BEAD_ID} case=resume_block1_decoded"
968 );
969 assert!(
970 !restored.blocks[2].decoded,
971 "bead_id={TEST_BEAD_ID} case=resume_block2_not_decoded"
972 );
973 }
974
975 #[test]
976 fn test_resume_no_protocol_negotiation() {
977 let mut resume = ResumeState::new(2);
979 resume.blocks[0].record_isi(0);
980 resume.blocks[0].record_isi(1);
981
982 let bytes = resume.to_bytes();
984 let restored = ResumeState::from_bytes(&bytes).expect("deserialize");
985 assert_eq!(
986 restored.blocks[0].num_received, 2,
987 "bead_id={TEST_BEAD_ID} case=resume_no_negotiation"
988 );
989 assert!(!restored.all_decoded());
990 }
991
992 #[test]
997 fn test_snapshot_single_block() {
998 let page_size = 256_u32;
999 let page_numbers: Vec<u32> = (1..=10).collect();
1000 let mut pages = make_pages(page_size, &page_numbers);
1001
1002 let config = SenderConfig {
1003 symbol_size: 256,
1004 max_isi_multiplier: 1,
1005 };
1006 let mut sender = SnapshotSender::prepare(page_size, &mut pages, config).expect("prepare");
1007 assert_eq!(
1008 sender.num_blocks(),
1009 1,
1010 "bead_id={TEST_BEAD_ID} case=single_block"
1011 );
1012
1013 let mut packets = Vec::new();
1015 while let Some(pkt) = sender.next_packet() {
1016 packets.push(pkt);
1017 }
1018 assert!(
1019 !packets.is_empty(),
1020 "bead_id={TEST_BEAD_ID} case=has_packets"
1021 );
1022
1023 let mut receiver = SnapshotReceiver::new(1, page_size);
1025 for pkt in &packets {
1026 let _ = receiver.process_packet(pkt);
1027 }
1028
1029 assert_eq!(
1030 receiver.state(),
1031 SnapshotReceiverState::Complete,
1032 "bead_id={TEST_BEAD_ID} case=single_block_complete"
1033 );
1034
1035 let blocks = receiver.take_decoded_blocks();
1036 assert_eq!(blocks.len(), 1);
1037 assert_eq!(blocks[0].pages.len(), 10);
1038 }
1039
1040 #[test]
1041 fn test_snapshot_multi_block_small() {
1042 let page_size = 64_u32;
1046 let n_pages = 200_u32;
1047 let page_numbers: Vec<u32> = (1..=n_pages).collect();
1048 let mut pages = make_pages(page_size, &page_numbers);
1049
1050 let config = SenderConfig {
1051 symbol_size: 64,
1052 max_isi_multiplier: 1,
1053 };
1054 let mut sender = SnapshotSender::prepare(page_size, &mut pages, config).expect("prepare");
1055
1056 assert_eq!(
1058 sender.num_blocks(),
1059 1,
1060 "bead_id={TEST_BEAD_ID} case=multi_block_small_count"
1061 );
1062
1063 let mut packets = Vec::new();
1064 while let Some(pkt) = sender.next_packet() {
1065 packets.push(pkt);
1066 }
1067
1068 let mut receiver = SnapshotReceiver::new(sender.num_blocks(), page_size);
1069 for pkt in &packets {
1070 let _ = receiver.process_packet(pkt);
1071 }
1072
1073 assert_eq!(
1074 receiver.state(),
1075 SnapshotReceiverState::Complete,
1076 "bead_id={TEST_BEAD_ID} case=multi_block_small_complete"
1077 );
1078
1079 let blocks = receiver.take_decoded_blocks();
1080 let total_pages: usize = blocks.iter().map(|b| b.pages.len()).sum();
1081 assert_eq!(
1082 total_pages, n_pages as usize,
1083 "bead_id={TEST_BEAD_ID} case=multi_block_all_pages"
1084 );
1085 }
1086
1087 #[test]
1088 fn test_duplicate_isi_discarded() {
1089 let page_size = 128_u32;
1090 let mut pages = make_pages(page_size, &[1, 2, 3]);
1091 let config = SenderConfig {
1092 symbol_size: 128,
1093 max_isi_multiplier: 1,
1094 };
1095 let mut sender = SnapshotSender::prepare(page_size, &mut pages, config).expect("prepare");
1096
1097 let mut packets = Vec::new();
1098 while let Some(pkt) = sender.next_packet() {
1099 packets.push(pkt);
1100 }
1101
1102 let mut receiver = SnapshotReceiver::new(1, page_size);
1103
1104 let r1 = receiver.process_packet(&packets[0]).expect("first");
1106 assert_ne!(
1107 r1,
1108 SnapshotPacketResult::Duplicate,
1109 "bead_id={TEST_BEAD_ID} case=first_not_dup"
1110 );
1111 let r2 = receiver.process_packet(&packets[0]).expect("duplicate");
1112 assert_eq!(
1113 r2,
1114 SnapshotPacketResult::Duplicate,
1115 "bead_id={TEST_BEAD_ID} case=dup_discarded"
1116 );
1117 }
1118
1119 #[test]
1120 fn test_snapshot_progressive_receive() {
1121 let page_size = 128_u32;
1125 let mut pages = make_pages(page_size, &[1, 2, 3, 4, 5]);
1126 let config = SenderConfig {
1127 symbol_size: 128,
1128 max_isi_multiplier: 1,
1129 };
1130 let mut sender = SnapshotSender::prepare(page_size, &mut pages, config).expect("prepare");
1131
1132 let mut packets = Vec::new();
1133 while let Some(pkt) = sender.next_packet() {
1134 packets.push(pkt);
1135 }
1136
1137 let mut receiver = SnapshotReceiver::new(1, page_size);
1138 let mut block_decoded_at = None;
1139
1140 for (i, pkt) in packets.iter().enumerate() {
1141 if let Ok(SnapshotPacketResult::BlockDecoded(_)) = receiver.process_packet(pkt) {
1142 block_decoded_at = Some(i);
1143 break;
1144 }
1145 }
1146
1147 assert!(
1148 block_decoded_at.is_some(),
1149 "bead_id={TEST_BEAD_ID} case=progressive_block_decoded"
1150 );
1151
1152 let blocks = receiver.take_decoded_blocks();
1154 assert!(
1155 !blocks.is_empty(),
1156 "bead_id={TEST_BEAD_ID} case=progressive_has_pages"
1157 );
1158 }
1159
1160 #[test]
1165 fn test_e2e_sender_receiver_roundtrip() {
1166 let page_size = 512_u32;
1167 let n_pages = 50_u32;
1168 let page_numbers: Vec<u32> = (1..=n_pages).collect();
1169 let original_pages = make_pages(page_size, &page_numbers);
1170 let mut pages = original_pages.clone();
1171
1172 let config = SenderConfig {
1173 symbol_size: 512,
1174 max_isi_multiplier: 1,
1175 };
1176 let mut sender = SnapshotSender::prepare(page_size, &mut pages, config).expect("prepare");
1177
1178 let mut packets = Vec::new();
1179 while let Some(pkt) = sender.next_packet() {
1180 packets.push(pkt);
1181 }
1182
1183 let mut receiver = SnapshotReceiver::new(sender.num_blocks(), page_size);
1184 for pkt in &packets {
1185 let _ = receiver.process_packet(pkt);
1186 }
1187
1188 assert_eq!(
1189 receiver.state(),
1190 SnapshotReceiverState::Complete,
1191 "bead_id={TEST_BEAD_ID} case=e2e_roundtrip_complete"
1192 );
1193
1194 let blocks = receiver.take_decoded_blocks();
1195 let mut all_decoded_pages: Vec<&DecodedBlockPage> =
1196 blocks.iter().flat_map(|b| b.pages.iter()).collect();
1197 all_decoded_pages.sort_by_key(|p| p.page_number);
1198
1199 assert_eq!(
1200 all_decoded_pages.len(),
1201 original_pages.len(),
1202 "bead_id={TEST_BEAD_ID} case=e2e_page_count"
1203 );
1204
1205 for (decoded, original) in all_decoded_pages.iter().zip(original_pages.iter()) {
1206 assert_eq!(
1207 decoded.page_number, original.page_number,
1208 "bead_id={TEST_BEAD_ID} case=e2e_page_number"
1209 );
1210 assert_eq!(
1211 decoded.page_data, original.page_bytes,
1212 "bead_id={TEST_BEAD_ID} case=e2e_page_data pn={}",
1213 original.page_number
1214 );
1215 }
1216 }
1217
1218 #[test]
1219 fn test_e2e_resume_after_partial() {
1220 let page_size = 128_u32;
1221 let n_pages = 20_u32;
1222 let mut pages = make_pages(page_size, &(1..=n_pages).collect::<Vec<_>>());
1223
1224 let config = SenderConfig {
1225 symbol_size: 128,
1226 max_isi_multiplier: 1,
1227 };
1228 let mut sender = SnapshotSender::prepare(page_size, &mut pages, config).expect("prepare");
1229
1230 let mut packets = Vec::new();
1231 while let Some(pkt) = sender.next_packet() {
1232 packets.push(pkt);
1233 }
1234
1235 let half = packets.len() / 2;
1237 let mut receiver1 = SnapshotReceiver::new(sender.num_blocks(), page_size);
1238 for pkt in &packets[..half] {
1239 let _ = receiver1.process_packet(pkt);
1240 }
1241
1242 let resume_bytes = receiver1.resume_state().to_bytes();
1244
1245 let resume = ResumeState::from_bytes(&resume_bytes).expect("restore");
1247 let mut receiver2 = SnapshotReceiver::from_resume(resume, page_size);
1248
1249 for pkt in &packets {
1251 let _ = receiver2.process_packet(pkt);
1252 }
1253
1254 assert_eq!(
1256 receiver2.state(),
1257 SnapshotReceiverState::Complete,
1258 "bead_id={TEST_BEAD_ID} case=e2e_resume_complete"
1259 );
1260 }
1261
1262 #[test]
1263 fn test_e2e_bd_1hi_15_compliance() {
1264 let page_size = 256_u32;
1266 let n_pages = 30_u32;
1267 let original_pages = make_pages(page_size, &(1..=n_pages).collect::<Vec<_>>());
1268 let mut pages = original_pages;
1269
1270 let config = SenderConfig {
1271 symbol_size: 256,
1272 max_isi_multiplier: 1,
1273 };
1274 let mut sender = SnapshotSender::prepare(page_size, &mut pages, config).expect("prepare");
1275
1276 assert!(
1278 sender.num_blocks() >= 1,
1279 "bead_id={TEST_BEAD_ID} case=compliance_has_blocks"
1280 );
1281 assert!(
1282 sender.total_source_symbols() > 0,
1283 "bead_id={TEST_BEAD_ID} case=compliance_has_symbols"
1284 );
1285
1286 let mut packets = Vec::new();
1287 while let Some(pkt) = sender.next_packet() {
1288 packets.push(pkt);
1289 }
1290
1291 let mut receiver = SnapshotReceiver::new(sender.num_blocks(), page_size);
1292 assert_eq!(receiver.state(), SnapshotReceiverState::Waiting);
1293
1294 for pkt in &packets {
1295 let _ = receiver.process_packet(pkt);
1296 }
1297 assert_eq!(receiver.state(), SnapshotReceiverState::Complete);
1298
1299 let blocks = receiver.take_decoded_blocks();
1300 let total_decoded: usize = blocks.iter().map(|b| b.pages.len()).sum();
1301 assert_eq!(
1302 total_decoded, n_pages as usize,
1303 "bead_id={TEST_BEAD_ID} case=compliance_all_pages_decoded"
1304 );
1305
1306 assert!(
1308 receiver.resume_state().all_decoded(),
1309 "bead_id={TEST_BEAD_ID} case=compliance_resume_all_decoded"
1310 );
1311 }
1312
1313 #[test]
1318 fn prop_partition_covers_all_pages() {
1319 for p in [1_u32, 10, 100, 1000, 56_403, 56_404, 100_000] {
1320 let blocks = partition_source_blocks(p).expect("partition");
1321 let total: u32 = blocks.iter().map(|b| b.num_pages).sum();
1322 assert_eq!(
1323 total, p,
1324 "bead_id={TEST_BEAD_ID} case=prop_partition_covers p={p}"
1325 );
1326 }
1327 }
1328
1329 #[test]
1330 fn prop_partition_block_sizes_valid() {
1331 for p in [1_u32, 56_403, 56_404, 200_000] {
1332 let blocks = partition_source_blocks(p).expect("partition");
1333 for block in &blocks {
1334 assert!(
1335 block.num_pages <= K_MAX,
1336 "bead_id={TEST_BEAD_ID} case=prop_block_size p={p} block={} num_pages={}",
1337 block.index,
1338 block.num_pages
1339 );
1340 }
1341 }
1342 }
1343
1344 #[test]
1349 fn test_bd_1hi_15_unit_compliance_gate() {
1350 let _ = SnapshotReceiverState::Waiting;
1352 let _ = SnapshotReceiverState::Receiving;
1353 let _ = SnapshotReceiverState::Complete;
1354
1355 let _ = SnapshotPacketResult::Accepted;
1356 let _ = SnapshotPacketResult::Duplicate;
1357 let _ = SnapshotPacketResult::Rejected;
1358 let _ = SnapshotPacketResult::AlreadyComplete;
1359
1360 let resume = ResumeState::new(3);
1361 assert_eq!(resume.total_blocks, 3);
1362 assert!(!resume.all_decoded());
1363 assert_eq!(resume.decoded_count(), 0);
1364
1365 let block = BlockResumeState::new(0);
1367 let bytes = block.to_bytes();
1368 let (restored, _) = BlockResumeState::from_bytes(&bytes).expect("deser");
1369 assert_eq!(restored.block_id, 0);
1370 }
1371
1372 #[test]
1373 fn prop_bd_1hi_15_structure_compliance() {
1374 let page_size = 128_u32;
1376 let mut pages = make_pages(page_size, &[1, 2]);
1377 let config = SenderConfig {
1378 symbol_size: 128,
1379 max_isi_multiplier: 1,
1380 };
1381 let mut sender = SnapshotSender::prepare(page_size, &mut pages, config).expect("prepare");
1382 assert!(sender.num_blocks() >= 1);
1383
1384 let mut packets = Vec::new();
1385 while let Some(pkt) = sender.next_packet() {
1386 packets.push(pkt);
1387 }
1388
1389 let mut receiver = SnapshotReceiver::new(sender.num_blocks(), page_size);
1390 for pkt in &packets {
1391 let _ = receiver.process_packet(pkt);
1392 }
1393 assert_eq!(receiver.state(), SnapshotReceiverState::Complete);
1394 }
1395}