1use std::collections::{HashMap, HashSet};
10
11use fsqlite_error::{FrankenError, Result};
12use fsqlite_types::ObjectId;
13use tracing::{debug, error, info, warn};
14
15use crate::decode_proofs::{DecodeAuditEntry, EcsDecodeProof};
16use crate::replication_sender::{
17 CHANGESET_HEADER_SIZE, ChangesetHeader, ChangesetId, DEFAULT_RPC_MESSAGE_CAP_BYTES, PageEntry,
18 ReplicationPacket, ReplicationWireVersion, compute_changeset_id,
19};
20use crate::source_block_partition::K_MAX;
21
22const BEAD_ID: &str = "bd-1hi.14";
23const DEFAULT_MAX_INFLIGHT_DECODERS: usize = 128;
24const DEFAULT_MAX_BUFFERED_SYMBOL_BYTES: usize = 64 * 1024 * 1024;
25
26#[derive(Debug, Clone, Copy, PartialEq, Eq)]
32pub enum ReceiverState {
33 Listening,
35 Collecting,
37 Decoding,
39 Applying,
41 Complete,
43}
44
45#[derive(Debug)]
47pub struct DecoderState {
48 pub k_source: u32,
50 pub symbol_size: u32,
52 pub seed: u64,
54 symbols: HashMap<u32, Vec<u8>>,
56 received_isis: HashSet<u32>,
58}
59
60impl DecoderState {
61 fn new(k_source: u32, symbol_size: u32, seed: u64) -> Self {
63 Self {
64 k_source,
65 symbol_size,
66 seed,
67 symbols: HashMap::with_capacity(k_source as usize),
68 received_isis: HashSet::with_capacity(k_source as usize),
69 }
70 }
71
72 #[must_use]
74 pub fn received_count(&self) -> u32 {
75 u32::try_from(self.received_isis.len()).unwrap_or(u32::MAX)
76 }
77
78 #[must_use]
80 pub fn ready_to_decode(&self) -> bool {
81 self.received_count() >= self.k_source
82 }
83
84 #[must_use]
86 pub fn source_symbol_count(&self) -> u32 {
87 let count = self
88 .symbols
89 .keys()
90 .filter(|&&isi| isi < self.k_source)
91 .count();
92 u32::try_from(count).unwrap_or(u32::MAX)
93 }
94
95 #[must_use]
97 pub fn has_repair_symbols(&self) -> bool {
98 self.symbols.keys().any(|&isi| isi >= self.k_source)
99 }
100
101 #[must_use]
103 pub fn sorted_isis(&self) -> Vec<u32> {
104 let mut isis: Vec<u32> = self.symbols.keys().copied().collect();
105 isis.sort_unstable();
106 isis.dedup();
107 isis
108 }
109
110 fn add_symbol(&mut self, isi: u32, data: Vec<u8>) -> bool {
112 if self.received_isis.contains(&isi) {
113 return false;
114 }
115 self.received_isis.insert(isi);
116 self.symbols.insert(isi, data);
117 true
118 }
119
120 #[must_use]
121 fn has_symbol(&self, isi: u32) -> bool {
122 self.received_isis.contains(&isi)
123 }
124
125 #[must_use]
126 fn buffered_bytes(&self) -> usize {
127 self.symbols.values().map(Vec::len).sum()
128 }
129
130 fn try_decode(&self) -> Option<Vec<u8>> {
139 if !self.ready_to_decode() {
140 return None;
141 }
142
143 let source_count = usize::try_from(self.source_symbol_count()).unwrap_or(usize::MAX);
145
146 let k = self.k_source as usize;
147 let t = self.symbol_size as usize;
148
149 if source_count >= k {
150 let padded_len = k * t;
152 let mut padded = vec![0_u8; padded_len];
153 for isi in 0..self.k_source {
154 if let Some(data) = self.symbols.get(&isi) {
155 let start = isi as usize * t;
156 let copy_len = data.len().min(t);
157 padded[start..start + copy_len].copy_from_slice(&data[..copy_len]);
158 }
159 }
160 Some(padded)
161 } else {
162 warn!(
165 bead_id = BEAD_ID,
166 source_count,
167 k_source = self.k_source,
168 total_received = self.received_count(),
169 "decode requires repair symbols (production uses RaptorQ decoder)"
170 );
171 None
172 }
173 }
174}
175
176#[derive(Debug, Clone, PartialEq, Eq)]
178pub struct DecodedPage {
179 pub page_number: u32,
181 pub page_data: Vec<u8>,
183}
184
185#[derive(Debug)]
187pub struct DecodeResult {
188 pub changeset_id: ChangesetId,
190 pub pages: Vec<DecodedPage>,
192 pub symbols_used: u32,
194 pub decode_proof: Option<EcsDecodeProof>,
196}
197
198#[derive(Debug, Clone, Copy)]
199struct DecodeProofBuildInput<'a> {
200 changeset_id: ChangesetId,
201 k_source: u32,
202 symbol_size: u32,
203 seed: u64,
204 received_isis: &'a [u32],
205 decode_success: bool,
206 intermediate_rank: Option<u32>,
207 symbols_used: u32,
208}
209
210#[derive(Debug)]
212pub struct ReplicationReceiver {
213 config: ReceiverConfig,
214 state: ReceiverState,
215 decoders: HashMap<ChangesetId, DecoderState>,
217 received_counts: HashMap<ChangesetId, u32>,
219 buffered_symbol_bytes: usize,
221 pending_results: Vec<DecodeResult>,
223 applied_count: u64,
225 decode_audit: Vec<DecodeAuditEntry>,
227 decode_audit_seq: u64,
229}
230
231#[derive(Debug, Clone, Copy, PartialEq, Eq)]
233pub struct DecodeProofEmissionPolicy {
234 pub emit_on_decode_failure: bool,
236 pub emit_on_repair_success: bool,
238}
239
240impl DecodeProofEmissionPolicy {
241 #[must_use]
243 pub const fn disabled() -> Self {
244 Self {
245 emit_on_decode_failure: false,
246 emit_on_repair_success: false,
247 }
248 }
249
250 #[must_use]
252 pub const fn durability_critical() -> Self {
253 Self {
254 emit_on_decode_failure: true,
255 emit_on_repair_success: true,
256 }
257 }
258}
259
260impl Default for DecodeProofEmissionPolicy {
261 fn default() -> Self {
262 Self::disabled()
263 }
264}
265
266#[derive(Debug, Clone)]
268pub struct ReceiverConfig {
269 pub auth_key: Option<[u8; 32]>,
271 pub decode_proof_policy: DecodeProofEmissionPolicy,
273 pub max_inflight_decoders: usize,
275 pub max_buffered_symbol_bytes: usize,
277}
278
279impl ReceiverConfig {
280 #[must_use]
282 pub const fn with_auth_key(auth_key: [u8; 32]) -> Self {
283 Self {
284 auth_key: Some(auth_key),
285 decode_proof_policy: DecodeProofEmissionPolicy::disabled(),
286 max_inflight_decoders: DEFAULT_MAX_INFLIGHT_DECODERS,
287 max_buffered_symbol_bytes: DEFAULT_MAX_BUFFERED_SYMBOL_BYTES,
288 }
289 }
290}
291
292impl Default for ReceiverConfig {
293 fn default() -> Self {
294 Self {
295 auth_key: None,
296 decode_proof_policy: DecodeProofEmissionPolicy::disabled(),
297 max_inflight_decoders: DEFAULT_MAX_INFLIGHT_DECODERS,
298 max_buffered_symbol_bytes: DEFAULT_MAX_BUFFERED_SYMBOL_BYTES,
299 }
300 }
301}
302
303impl ReplicationReceiver {
304 fn remove_decoder(&mut self, changeset_id: ChangesetId) {
305 if let Some(decoder) = self.decoders.remove(&changeset_id) {
306 self.buffered_symbol_bytes = self
307 .buffered_symbol_bytes
308 .saturating_sub(decoder.buffered_bytes());
309 }
310 self.received_counts.remove(&changeset_id);
311 }
312
313 #[must_use]
315 pub fn with_config(config: ReceiverConfig) -> Self {
316 Self {
317 config,
318 state: ReceiverState::Listening,
319 decoders: HashMap::new(),
320 received_counts: HashMap::new(),
321 buffered_symbol_bytes: 0,
322 pending_results: Vec::new(),
323 applied_count: 0,
324 decode_audit: Vec::new(),
325 decode_audit_seq: 0,
326 }
327 }
328
329 #[must_use]
331 pub fn new() -> Self {
332 Self::with_config(ReceiverConfig::default())
333 }
334
335 #[must_use]
337 pub const fn state(&self) -> ReceiverState {
338 self.state
339 }
340
341 #[must_use]
343 pub const fn applied_count(&self) -> u64 {
344 self.applied_count
345 }
346
347 #[must_use]
349 pub fn active_decoders(&self) -> usize {
350 self.decoders.len()
351 }
352
353 #[must_use]
355 pub fn decode_audit_entries(&self) -> &[DecodeAuditEntry] {
356 &self.decode_audit
357 }
358
359 pub fn take_decode_audit_entries(&mut self) -> Vec<DecodeAuditEntry> {
361 std::mem::take(&mut self.decode_audit)
362 }
363
364 pub fn process_packet(&mut self, packet_bytes: &[u8]) -> Result<PacketResult> {
374 if packet_bytes.len() > DEFAULT_RPC_MESSAGE_CAP_BYTES {
375 return Err(FrankenError::TooBig);
376 }
377 let packet = ReplicationPacket::from_bytes(packet_bytes)?;
378 if !packet.verify_integrity(self.config.auth_key.as_ref()) {
379 warn!(
380 bead_id = BEAD_ID,
381 wire_version = ?packet.wire_version,
382 has_auth = packet.auth_tag.is_some(),
383 "packet integrity/auth verification failed; treating as erasure"
384 );
385 return Ok(PacketResult::Erasure);
386 }
387 self.process_parsed_packet(&packet)
388 }
389
390 #[allow(clippy::too_many_lines)]
396 pub fn process_parsed_packet(&mut self, packet: &ReplicationPacket) -> Result<PacketResult> {
397 if packet.sbn != 0 {
399 error!(
400 bead_id = BEAD_ID,
401 sbn = packet.sbn,
402 "V1 rule: SBN must be 0"
403 );
404 return Err(FrankenError::Internal(format!(
405 "V1 replication: source_block must be 0, got {}",
406 packet.sbn
407 )));
408 }
409
410 if packet.k_source == 0 || packet.k_source > K_MAX {
412 error!(
413 bead_id = BEAD_ID,
414 k_source = packet.k_source,
415 k_max = K_MAX,
416 "K_source out of valid range"
417 );
418 return Err(FrankenError::OutOfRange {
419 what: "k_source".to_owned(),
420 value: packet.k_source.to_string(),
421 });
422 }
423
424 if usize::from(packet.symbol_size_t) != packet.symbol_data.len() {
426 return Err(FrankenError::DatabaseCorrupt {
427 detail: format!(
428 "symbol_size_t mismatch: header={}, payload={}",
429 packet.symbol_size_t,
430 packet.symbol_data.len()
431 ),
432 });
433 }
434 let symbol_size = u32::from(packet.symbol_size_t);
435 if symbol_size == 0 {
436 return Err(FrankenError::OutOfRange {
437 what: "symbol_size".to_owned(),
438 value: "0".to_owned(),
439 });
440 }
441
442 if self.state == ReceiverState::Listening {
444 self.state = ReceiverState::Collecting;
445 info!(bead_id = BEAD_ID, "first packet received, now COLLECTING");
446 }
447
448 let changeset_id = packet.changeset_id;
449 let mut created_decoder = false;
450
451 if let Some(decoder) = self.decoders.get(&changeset_id) {
453 if decoder.k_source != packet.k_source {
455 error!(
456 bead_id = BEAD_ID,
457 expected_k = decoder.k_source,
458 got_k = packet.k_source,
459 "K_source mismatch for existing changeset"
460 );
461 return Err(FrankenError::DatabaseCorrupt {
462 detail: format!(
463 "K_source mismatch: expected {}, got {}",
464 decoder.k_source, packet.k_source
465 ),
466 });
467 }
468 if decoder.symbol_size != symbol_size {
469 error!(
470 bead_id = BEAD_ID,
471 expected_t = decoder.symbol_size,
472 got_t = symbol_size,
473 "symbol_size mismatch for existing changeset"
474 );
475 return Err(FrankenError::DatabaseCorrupt {
476 detail: format!(
477 "symbol_size mismatch: expected {}, got {}",
478 decoder.symbol_size, symbol_size
479 ),
480 });
481 }
482 if packet.wire_version == ReplicationWireVersion::FramedV2
483 && decoder.seed != packet.seed
484 {
485 return Err(FrankenError::DatabaseCorrupt {
486 detail: format!(
487 "seed mismatch: expected {}, got {}",
488 decoder.seed, packet.seed
489 ),
490 });
491 }
492 } else {
493 if self.decoders.len() >= self.config.max_inflight_decoders {
494 warn!(
495 bead_id = BEAD_ID,
496 active_decoders = self.decoders.len(),
497 max_inflight_decoders = self.config.max_inflight_decoders,
498 "decoder cap reached; rejecting new changeset"
499 );
500 return Err(FrankenError::Busy);
501 }
502 let expected_seed =
504 crate::replication_sender::derive_seed_from_changeset_id(&changeset_id);
505 if packet.wire_version == ReplicationWireVersion::FramedV2
506 && packet.seed != expected_seed
507 {
508 return Err(FrankenError::DatabaseCorrupt {
509 detail: format!(
510 "seed does not match deterministic derivation for changeset: expected {expected_seed}, got {}",
511 packet.seed
512 ),
513 });
514 }
515 let seed = expected_seed;
516 debug!(
517 bead_id = BEAD_ID,
518 k_source = packet.k_source,
519 symbol_size,
520 seed,
521 "created decoder for new changeset"
522 );
523 self.decoders.insert(
524 changeset_id,
525 DecoderState::new(packet.k_source, symbol_size, seed),
526 );
527 self.received_counts.insert(changeset_id, 0);
528 created_decoder = true;
529 }
530
531 if let Some(decoder) = self.decoders.get(&changeset_id)
533 && !decoder.has_symbol(packet.esi)
534 {
535 let next_total = self
536 .buffered_symbol_bytes
537 .saturating_add(packet.symbol_data.len());
538 if next_total > self.config.max_buffered_symbol_bytes {
539 warn!(
540 bead_id = BEAD_ID,
541 buffered_symbol_bytes = self.buffered_symbol_bytes,
542 incoming_symbol_bytes = packet.symbol_data.len(),
543 max_buffered_symbol_bytes = self.config.max_buffered_symbol_bytes,
544 "buffered symbol budget exceeded"
545 );
546 if created_decoder {
547 self.remove_decoder(changeset_id);
548 self.state = if self.decoders.is_empty() {
549 ReceiverState::Listening
550 } else {
551 ReceiverState::Collecting
552 };
553 }
554 return Err(FrankenError::TooBig);
555 }
556 }
557
558 let (
560 ready_to_decode,
561 k_source_ctx,
562 symbol_size_ctx,
563 seed_ctx,
564 received_isis_ctx,
565 received_count_ctx,
566 source_count_ctx,
567 has_repair_ctx,
568 decoded_padded,
569 ) = {
570 let decoder = self.decoders.get_mut(&changeset_id).expect("just inserted");
571 let accepted = decoder.add_symbol(packet.esi, packet.symbol_data.clone());
572
573 if !accepted {
574 debug!(
575 bead_id = BEAD_ID,
576 isi = packet.esi,
577 "duplicate ISI, symbol ignored"
578 );
579 return Ok(PacketResult::Duplicate);
580 }
581
582 self.buffered_symbol_bytes = self
583 .buffered_symbol_bytes
584 .saturating_add(packet.symbol_data.len());
585 let count = self.received_counts.entry(changeset_id).or_insert(0);
586 *count += 1;
587 debug!(
588 bead_id = BEAD_ID,
589 isi = packet.esi,
590 received = *count,
591 k_source = packet.k_source,
592 "symbol accepted"
593 );
594
595 let ready = decoder.ready_to_decode();
596 let padded = if ready { decoder.try_decode() } else { None };
597 (
598 ready,
599 decoder.k_source,
600 decoder.symbol_size,
601 decoder.seed,
602 decoder.sorted_isis(),
603 decoder.received_count(),
604 decoder.source_symbol_count(),
605 decoder.has_repair_symbols(),
606 padded,
607 )
608 };
609
610 if ready_to_decode {
611 info!(
612 bead_id = BEAD_ID,
613 received = received_count_ctx,
614 k_source = k_source_ctx,
615 "attempting decode"
616 );
617 self.state = ReceiverState::Decoding;
618
619 if let Some(padded_bytes) = decoded_padded {
620 let success_proof =
621 if self.config.decode_proof_policy.emit_on_repair_success && has_repair_ctx {
622 Some(Self::build_decode_proof(DecodeProofBuildInput {
623 changeset_id,
624 k_source: k_source_ctx,
625 symbol_size: symbol_size_ctx,
626 seed: seed_ctx,
627 received_isis: &received_isis_ctx,
628 decode_success: true,
629 intermediate_rank: Some(k_source_ctx),
630 symbols_used: received_count_ctx,
631 }))
632 } else {
633 None
634 };
635
636 match self.parse_and_validate_changeset(changeset_id, &padded_bytes) {
638 Ok(mut result) => {
639 let n_pages = result.pages.len();
640 if let Some(proof) = success_proof {
641 self.record_decode_proof(proof.clone());
642 result.decode_proof = Some(proof);
643 }
644 self.pending_results.push(result);
645 self.state = ReceiverState::Applying;
646 info!(
647 bead_id = BEAD_ID,
648 n_pages, "decode succeeded, ready to apply"
649 );
650 self.remove_decoder(changeset_id);
652 return Ok(PacketResult::DecodeReady);
653 }
654 Err(e) => {
655 error!(
656 bead_id = BEAD_ID,
657 error = %e,
658 "changeset validation failed after decode"
659 );
660 self.remove_decoder(changeset_id);
662 self.state = if self.decoders.is_empty() {
663 ReceiverState::Listening
664 } else {
665 ReceiverState::Collecting
666 };
667 return Err(e);
668 }
669 }
670 }
671
672 if self.config.decode_proof_policy.emit_on_decode_failure {
673 let failure_proof = Self::build_decode_proof(DecodeProofBuildInput {
674 changeset_id,
675 k_source: k_source_ctx,
676 symbol_size: symbol_size_ctx,
677 seed: seed_ctx,
678 received_isis: &received_isis_ctx,
679 decode_success: false,
680 intermediate_rank: Some(source_count_ctx),
681 symbols_used: received_count_ctx,
682 });
683 self.record_decode_proof(failure_proof);
684 }
685
686 warn!(
688 bead_id = BEAD_ID,
689 source_count = source_count_ctx,
690 k_source = k_source_ctx,
691 "decode failed at K_source, continuing collection"
692 );
693 self.state = ReceiverState::Collecting;
694 return Ok(PacketResult::NeedMore);
695 }
696
697 Ok(PacketResult::Accepted)
698 }
699
700 #[allow(clippy::too_many_lines)]
702 fn parse_and_validate_changeset(
703 &self,
704 changeset_id: ChangesetId,
705 padded_bytes: &[u8],
706 ) -> Result<DecodeResult> {
707 if padded_bytes.len() < CHANGESET_HEADER_SIZE {
708 return Err(FrankenError::DatabaseCorrupt {
709 detail: format!(
710 "decoded bytes too short for header: {} < {CHANGESET_HEADER_SIZE}",
711 padded_bytes.len()
712 ),
713 });
714 }
715
716 let header_bytes: [u8; CHANGESET_HEADER_SIZE] = padded_bytes[..CHANGESET_HEADER_SIZE]
718 .try_into()
719 .expect("checked length");
720 let header = ChangesetHeader::from_bytes(&header_bytes)?;
721
722 let total_len =
724 usize::try_from(header.total_len).map_err(|_| FrankenError::OutOfRange {
725 what: "total_len".to_owned(),
726 value: header.total_len.to_string(),
727 })?;
728 if total_len < CHANGESET_HEADER_SIZE {
729 return Err(FrankenError::DatabaseCorrupt {
730 detail: format!(
731 "total_len ({total_len}) smaller than changeset header size ({CHANGESET_HEADER_SIZE})"
732 ),
733 });
734 }
735 if total_len > padded_bytes.len() {
736 return Err(FrankenError::DatabaseCorrupt {
737 detail: format!(
738 "total_len ({total_len}) exceeds decoded bytes ({})",
739 padded_bytes.len()
740 ),
741 });
742 }
743 let changeset_bytes = &padded_bytes[..total_len];
744 let computed_id = compute_changeset_id(changeset_bytes);
745 if computed_id != changeset_id {
746 return Err(FrankenError::DatabaseCorrupt {
747 detail: format!(
748 "changeset id mismatch: expected {changeset_id:?}, computed {computed_id:?}"
749 ),
750 });
751 }
752
753 let page_size =
755 usize::try_from(header.page_size).map_err(|_| FrankenError::OutOfRange {
756 what: "page_size".to_owned(),
757 value: header.page_size.to_string(),
758 })?;
759 if page_size == 0 {
760 return Err(FrankenError::OutOfRange {
761 what: "page_size".to_owned(),
762 value: "0".to_owned(),
763 });
764 }
765 let entry_size = 4_usize
766 .checked_add(8)
767 .and_then(|value| value.checked_add(page_size))
768 .ok_or_else(|| FrankenError::OutOfRange {
769 what: "entry_size".to_owned(),
770 value: format!("page_size={}", header.page_size),
771 })?; let n_pages = usize::try_from(header.n_pages).map_err(|_| FrankenError::OutOfRange {
773 what: "n_pages".to_owned(),
774 value: header.n_pages.to_string(),
775 })?;
776 let data_start = CHANGESET_HEADER_SIZE;
777 let data_bytes = &changeset_bytes[data_start..];
778 let required_bytes =
779 entry_size
780 .checked_mul(n_pages)
781 .ok_or_else(|| FrankenError::OutOfRange {
782 what: "changeset payload size".to_owned(),
783 value: format!("entry_size={entry_size}, n_pages={}", header.n_pages),
784 })?;
785
786 if data_bytes.len() != required_bytes {
787 return Err(FrankenError::DatabaseCorrupt {
788 detail: format!(
789 "changeset payload length mismatch for {} pages: {} != {}",
790 header.n_pages,
791 data_bytes.len(),
792 required_bytes,
793 ),
794 });
795 }
796
797 let mut pages = Vec::with_capacity(n_pages);
798 let decoder_state_symbols = self
799 .decoders
800 .get(&changeset_id)
801 .map_or(0, DecoderState::received_count);
802
803 for i in 0..n_pages {
804 let offset = i
805 .checked_mul(entry_size)
806 .ok_or_else(|| FrankenError::OutOfRange {
807 what: "page entry offset".to_owned(),
808 value: format!("index={i}, entry_size={entry_size}"),
809 })?;
810 let page_number =
811 u32::from_le_bytes(data_bytes[offset..offset + 4].try_into().expect("4 bytes"));
812 let page_xxh3 = u64::from_le_bytes(
813 data_bytes[offset + 4..offset + 12]
814 .try_into()
815 .expect("8 bytes"),
816 );
817 let page_data = data_bytes[offset + 12..offset + 12 + page_size].to_vec();
818
819 let computed_xxh3 = xxhash_rust::xxh3::xxh3_64(&page_data);
821 if computed_xxh3 != page_xxh3 {
822 error!(
823 bead_id = BEAD_ID,
824 page_number,
825 expected_xxh3 = page_xxh3,
826 computed_xxh3,
827 "page xxh3 validation failed"
828 );
829 return Err(FrankenError::DatabaseCorrupt {
830 detail: format!(
831 "page {page_number} xxh3 mismatch: expected {page_xxh3:#x}, got {computed_xxh3:#x}"
832 ),
833 });
834 }
835
836 pages.push(DecodedPage {
837 page_number,
838 page_data,
839 });
840 }
841
842 debug_assert!(
844 pages
845 .windows(2)
846 .all(|w| w[0].page_number <= w[1].page_number)
847 );
848
849 Ok(DecodeResult {
850 changeset_id,
851 pages,
852 symbols_used: decoder_state_symbols,
853 decode_proof: None,
854 })
855 }
856
857 fn build_decode_proof(input: DecodeProofBuildInput<'_>) -> EcsDecodeProof {
858 let object_id = ObjectId::from_bytes(*input.changeset_id.as_bytes());
859 let timing_ns =
860 deterministic_timing_ns(input.k_source, input.symbol_size, input.symbols_used);
861 EcsDecodeProof::from_esis(
862 object_id,
863 input.k_source,
864 input.received_isis,
865 input.decode_success,
866 input.intermediate_rank,
867 timing_ns,
868 input.seed,
869 )
870 .with_changeset_id(*input.changeset_id.as_bytes())
871 }
872
873 fn record_decode_proof(&mut self, proof: EcsDecodeProof) {
874 self.decode_audit_seq = self.decode_audit_seq.saturating_add(1);
875 self.decode_audit.push(DecodeAuditEntry {
876 proof,
877 seq: self.decode_audit_seq,
878 lab_mode: false,
879 });
880 }
881
882 pub fn apply_pending(&mut self) -> Result<Vec<DecodeResult>> {
891 if self.state != ReceiverState::Applying {
892 return Err(FrankenError::Internal(format!(
893 "receiver must be APPLYING to apply, current state: {:?}",
894 self.state
895 )));
896 }
897
898 let results = std::mem::take(&mut self.pending_results);
899 let n = results.len();
900 self.applied_count += u64::try_from(n).unwrap_or(u64::MAX);
901
902 info!(
903 bead_id = BEAD_ID,
904 applied = n,
905 total_applied = self.applied_count,
906 "applied pending changesets"
907 );
908
909 self.state = ReceiverState::Complete;
911 Ok(results)
912 }
913
914 pub fn reset_to_listening(&mut self) -> Result<()> {
920 if self.state != ReceiverState::Complete {
921 return Err(FrankenError::Internal(format!(
922 "receiver must be COMPLETE to reset, current state: {:?}",
923 self.state
924 )));
925 }
926 self.state = ReceiverState::Listening;
927 debug!(bead_id = BEAD_ID, "receiver reset to LISTENING");
928 Ok(())
929 }
930
931 pub fn force_reset(&mut self) {
933 self.decoders.clear();
934 self.received_counts.clear();
935 self.buffered_symbol_bytes = 0;
936 self.pending_results.clear();
937 self.state = ReceiverState::Listening;
938 warn!(bead_id = BEAD_ID, "receiver force-reset to LISTENING");
939 }
940}
941
942impl Default for ReplicationReceiver {
943 fn default() -> Self {
944 Self::new()
945 }
946}
947
948fn deterministic_timing_ns(k_source: u32, symbol_size: u32, symbols_used: u32) -> u64 {
949 let mut material = [0_u8; 12];
950 material[..4].copy_from_slice(&k_source.to_le_bytes());
951 material[4..8].copy_from_slice(&symbol_size.to_le_bytes());
952 material[8..12].copy_from_slice(&symbols_used.to_le_bytes());
953 xxhash_rust::xxh3::xxh3_64(&material)
954}
955
956#[derive(Debug, Clone, Copy, PartialEq, Eq)]
958pub enum PacketResult {
959 Accepted,
961 Erasure,
963 Duplicate,
965 DecodeReady,
967 NeedMore,
969}
970
971pub fn parse_changeset_pages(changeset_bytes: &[u8]) -> Result<(ChangesetHeader, Vec<PageEntry>)> {
981 if changeset_bytes.len() < CHANGESET_HEADER_SIZE {
982 return Err(FrankenError::DatabaseCorrupt {
983 detail: format!(
984 "changeset too short: {} < {CHANGESET_HEADER_SIZE}",
985 changeset_bytes.len()
986 ),
987 });
988 }
989
990 let header_bytes: [u8; CHANGESET_HEADER_SIZE] = changeset_bytes[..CHANGESET_HEADER_SIZE]
991 .try_into()
992 .expect("checked length");
993 let header = ChangesetHeader::from_bytes(&header_bytes)?;
994
995 let total_len = usize::try_from(header.total_len).map_err(|_| FrankenError::OutOfRange {
996 what: "total_len".to_owned(),
997 value: header.total_len.to_string(),
998 })?;
999 if total_len < CHANGESET_HEADER_SIZE {
1000 return Err(FrankenError::DatabaseCorrupt {
1001 detail: format!(
1002 "total_len ({total_len}) smaller than changeset header size ({CHANGESET_HEADER_SIZE})"
1003 ),
1004 });
1005 }
1006 if total_len > changeset_bytes.len() {
1007 return Err(FrankenError::DatabaseCorrupt {
1008 detail: format!(
1009 "total_len ({total_len}) exceeds available bytes ({})",
1010 changeset_bytes.len()
1011 ),
1012 });
1013 }
1014 let changeset_bytes = &changeset_bytes[..total_len];
1015
1016 let page_size = usize::try_from(header.page_size).map_err(|_| FrankenError::OutOfRange {
1017 what: "page_size".to_owned(),
1018 value: header.page_size.to_string(),
1019 })?;
1020 if page_size == 0 {
1021 return Err(FrankenError::OutOfRange {
1022 what: "page_size".to_owned(),
1023 value: "0".to_owned(),
1024 });
1025 }
1026 let entry_size = 4_usize
1027 .checked_add(8)
1028 .and_then(|value| value.checked_add(page_size))
1029 .ok_or_else(|| FrankenError::OutOfRange {
1030 what: "entry_size".to_owned(),
1031 value: format!("page_size={}", header.page_size),
1032 })?;
1033 let n_pages = usize::try_from(header.n_pages).map_err(|_| FrankenError::OutOfRange {
1034 what: "n_pages".to_owned(),
1035 value: header.n_pages.to_string(),
1036 })?;
1037 let data_start = CHANGESET_HEADER_SIZE;
1038 let data_bytes = &changeset_bytes[data_start..];
1039 let required_bytes =
1040 entry_size
1041 .checked_mul(n_pages)
1042 .ok_or_else(|| FrankenError::OutOfRange {
1043 what: "changeset payload size".to_owned(),
1044 value: format!("entry_size={entry_size}, n_pages={}", header.n_pages),
1045 })?;
1046 if data_bytes.len() != required_bytes {
1047 return Err(FrankenError::DatabaseCorrupt {
1048 detail: format!(
1049 "changeset payload length mismatch for {} pages: {} != {}",
1050 header.n_pages,
1051 data_bytes.len(),
1052 required_bytes
1053 ),
1054 });
1055 }
1056
1057 let mut pages = Vec::with_capacity(n_pages);
1058 for i in 0..n_pages {
1059 let offset = i
1060 .checked_mul(entry_size)
1061 .ok_or_else(|| FrankenError::OutOfRange {
1062 what: "page entry offset".to_owned(),
1063 value: format!("index={i}, entry_size={entry_size}"),
1064 })?;
1065 let page_number =
1066 u32::from_le_bytes(data_bytes[offset..offset + 4].try_into().expect("4 bytes"));
1067 let page_xxh3 = u64::from_le_bytes(
1068 data_bytes[offset + 4..offset + 12]
1069 .try_into()
1070 .expect("8 bytes"),
1071 );
1072 let page_bytes = data_bytes[offset + 12..offset + 12 + page_size].to_vec();
1073
1074 pages.push(PageEntry {
1075 page_number,
1076 page_xxh3,
1077 page_bytes,
1078 });
1079 }
1080
1081 Ok((header, pages))
1082}
1083
1084#[cfg(test)]
1085mod tests {
1086 use asupersync::runtime::RuntimeBuilder;
1087 use asupersync::security::authenticated::AuthenticatedSymbol;
1088 use asupersync::security::tag::AuthenticationTag;
1089 use asupersync::transport::{
1090 SimNetwork, SimTransportConfig, SymbolSinkExt as _, SymbolStreamExt as _,
1091 };
1092 use asupersync::types::{Symbol, SymbolId, SymbolKind};
1093 use std::collections::HashSet;
1094
1095 use super::*;
1096 use crate::replication_sender::{
1097 CHANGESET_HEADER_SIZE, ChangesetId, PageEntry, REPLICATION_HEADER_SIZE, ReplicationPacket,
1098 ReplicationPacketV2Header, ReplicationSender, ReplicationWireVersion, SenderConfig,
1099 compute_changeset_id, derive_seed_from_changeset_id, encode_changeset,
1100 };
1101
1102 const TEST_BEAD_ID: &str = "bd-1hi.14";
1103
1104 #[allow(clippy::cast_possible_truncation)]
1105 fn make_pages(page_size: u32, page_numbers: &[u32]) -> Vec<PageEntry> {
1106 page_numbers
1107 .iter()
1108 .map(|&pn| {
1109 let mut data = vec![0_u8; page_size as usize];
1110 for (i, byte) in data.iter_mut().enumerate() {
1111 *byte = ((pn as usize * 251 + i * 31) % 256) as u8;
1112 }
1113 PageEntry::new(pn, data)
1114 })
1115 .collect()
1116 }
1117
1118 fn generate_sender_packets(
1120 page_size: u32,
1121 page_numbers: &[u32],
1122 symbol_size: u16,
1123 ) -> Vec<Vec<u8>> {
1124 generate_sender_packets_with_multiplier(page_size, page_numbers, symbol_size, 1)
1125 }
1126
1127 fn generate_sender_packets_with_multiplier(
1128 page_size: u32,
1129 page_numbers: &[u32],
1130 symbol_size: u16,
1131 max_isi_multiplier: u32,
1132 ) -> Vec<Vec<u8>> {
1133 let mut sender = ReplicationSender::new();
1134 let mut pages = make_pages(page_size, page_numbers);
1135 let config = SenderConfig {
1136 symbol_size,
1137 max_isi_multiplier,
1138 };
1139 sender
1140 .prepare(page_size, &mut pages, config)
1141 .expect("prepare");
1142 sender.start_streaming().expect("start");
1143
1144 let mut packets = Vec::new();
1145 while let Some(packet) = sender.next_packet().expect("next") {
1146 packets.push(packet.to_bytes().expect("encode"));
1147 }
1148 packets
1149 }
1150
1151 #[derive(Debug)]
1152 struct SimNetworkDelivery {
1153 sent_count: usize,
1154 delivered: Vec<(u32, Vec<u8>)>,
1155 }
1156
1157 fn packet_symbol(esi: u32, wire_bytes: Vec<u8>) -> AuthenticatedSymbol {
1158 let symbol_id = SymbolId::new_for_test(0xBEEF, 0, esi);
1159 let symbol = Symbol::new(symbol_id, wire_bytes, SymbolKind::Source);
1160 AuthenticatedSymbol::from_parts(symbol, AuthenticationTag::zero())
1161 }
1162
1163 fn transmit_packets_simnetwork(
1164 config: SimTransportConfig,
1165 packet_bytes: &[Vec<u8>],
1166 ) -> SimNetworkDelivery {
1167 let network = SimNetwork::fully_connected(2, config);
1168 let (mut sink, mut stream) = network.transport(0, 1);
1169 let runtime = RuntimeBuilder::current_thread()
1170 .build()
1171 .expect("runtime build");
1172
1173 runtime.block_on(async {
1174 for (index, bytes) in packet_bytes.iter().enumerate() {
1175 let esi = u32::try_from(index).expect("test packet index fits u32");
1176 sink.send(packet_symbol(esi, bytes.clone()))
1177 .await
1178 .expect("send simulated symbol");
1179 }
1180 sink.close().await.expect("close simulated sink");
1181
1182 let mut delivered = Vec::new();
1183 while let Some(item) = stream.next().await {
1184 let auth = item.expect("sim stream item");
1185 delivered.push((auth.symbol().id().esi(), auth.symbol().data().to_vec()));
1186 }
1187
1188 SimNetworkDelivery {
1189 sent_count: packet_bytes.len(),
1190 delivered,
1191 }
1192 })
1193 }
1194
1195 fn has_duplicate_esies(delivery: &SimNetworkDelivery) -> bool {
1196 let mut seen = HashSet::new();
1197 delivery.delivered.iter().any(|(esi, _)| !seen.insert(*esi))
1198 }
1199
1200 fn has_reordered_esies(delivery: &SimNetworkDelivery) -> bool {
1201 delivery
1202 .delivered
1203 .windows(2)
1204 .any(|window| window[0].0 > window[1].0)
1205 }
1206
1207 fn has_corrupted_wire_bytes(delivery: &SimNetworkDelivery, original: &[Vec<u8>]) -> bool {
1208 delivery.delivered.iter().any(|(esi, bytes)| {
1209 usize::try_from(*esi)
1210 .ok()
1211 .and_then(|index| original.get(index))
1212 .is_some_and(|expected| expected.as_slice() != bytes.as_slice())
1213 })
1214 }
1215
1216 fn decode_from_wire_packets(
1217 delivered: &[(u32, Vec<u8>)],
1218 ) -> (Option<Vec<DecodedPage>>, usize, usize) {
1219 let mut receiver = ReplicationReceiver::new();
1220 let mut erasures = 0_usize;
1221 let mut parse_errors = 0_usize;
1222
1223 for (_, wire) in delivered {
1224 match receiver.process_packet(wire) {
1225 Ok(PacketResult::DecodeReady) => {
1226 let mut applied = receiver.apply_pending().expect("apply decoded changeset");
1227 let pages = applied.pop().expect("decode result pages").pages;
1228 return (Some(pages), erasures, parse_errors);
1229 }
1230 Ok(PacketResult::Erasure) => erasures += 1,
1231 Ok(PacketResult::Accepted | PacketResult::Duplicate | PacketResult::NeedMore) => {}
1232 Err(_) => parse_errors += 1,
1233 }
1234 }
1235
1236 (None, erasures, parse_errors)
1237 }
1238
1239 fn decoded_matches_original(decoded: &[DecodedPage], original: &[PageEntry]) -> bool {
1240 if decoded.len() != original.len() {
1241 return false;
1242 }
1243 for (decoded, original) in decoded.iter().zip(original.iter()) {
1244 if decoded.page_number != original.page_number {
1245 return false;
1246 }
1247 if decoded.page_data != original.page_bytes {
1248 return false;
1249 }
1250 }
1251 true
1252 }
1253
1254 fn make_packet(
1255 changeset_id: ChangesetId,
1256 sbn: u8,
1257 esi: u32,
1258 k_source: u32,
1259 symbol_data: Vec<u8>,
1260 ) -> ReplicationPacket {
1261 let symbol_size_t =
1262 u16::try_from(symbol_data.len()).expect("test symbol payload must fit u16");
1263 let seed = derive_seed_from_changeset_id(&changeset_id);
1264 ReplicationPacket::new_v2(
1265 ReplicationPacketV2Header {
1266 changeset_id,
1267 sbn,
1268 esi,
1269 k_source,
1270 r_repair: 0,
1271 symbol_size_t,
1272 seed,
1273 },
1274 symbol_data,
1275 )
1276 }
1277
1278 fn receiver_with_decode_proofs() -> ReplicationReceiver {
1279 ReplicationReceiver::with_config(ReceiverConfig {
1280 auth_key: None,
1281 decode_proof_policy: DecodeProofEmissionPolicy::durability_critical(),
1282 ..ReceiverConfig::default()
1283 })
1284 }
1285
1286 #[test]
1291 fn test_receiver_listening_to_collecting() {
1292 let mut receiver = ReplicationReceiver::new();
1293 assert_eq!(
1294 receiver.state(),
1295 ReceiverState::Listening,
1296 "bead_id={TEST_BEAD_ID} case=initial_state"
1297 );
1298
1299 let packets = generate_sender_packets(512, &[1], 512);
1300 assert!(!packets.is_empty());
1301
1302 receiver.process_packet(&packets[0]).expect("first packet");
1303 assert_ne!(
1304 receiver.state(),
1305 ReceiverState::Listening,
1306 "bead_id={TEST_BEAD_ID} case=transition_on_first_packet"
1307 );
1308 }
1309
1310 #[test]
1311 fn test_receiver_decoder_creation() {
1312 let mut receiver = ReplicationReceiver::new();
1313 let packets = generate_sender_packets(512, &[1, 2], 512);
1314 assert_eq!(receiver.active_decoders(), 0);
1315
1316 receiver.process_packet(&packets[0]).expect("first packet");
1317 assert_ne!(
1321 receiver.state(),
1322 ReceiverState::Listening,
1323 "bead_id={TEST_BEAD_ID} case=decoder_created"
1324 );
1325 }
1326
1327 #[test]
1328 fn test_receiver_rejects_new_changeset_when_decoder_limit_hit() {
1329 let mut receiver = ReplicationReceiver::with_config(ReceiverConfig {
1330 max_inflight_decoders: 1,
1331 ..ReceiverConfig::default()
1332 });
1333
1334 let first = make_packet(
1335 ChangesetId::from_bytes([0x31; 16]),
1336 0,
1337 0,
1338 100,
1339 vec![0x11; 256],
1340 );
1341 receiver
1342 .process_parsed_packet(&first)
1343 .expect("first decoder");
1344 assert_eq!(receiver.active_decoders(), 1);
1345
1346 let second = make_packet(
1347 ChangesetId::from_bytes([0x32; 16]),
1348 0,
1349 0,
1350 100,
1351 vec![0x22; 256],
1352 );
1353 let err = receiver.process_parsed_packet(&second).unwrap_err();
1354 assert!(matches!(err, FrankenError::Busy));
1355 assert_eq!(receiver.active_decoders(), 1);
1356 }
1357
1358 #[test]
1359 fn test_receiver_enforces_buffered_symbol_budget() {
1360 let mut receiver = ReplicationReceiver::with_config(ReceiverConfig {
1361 max_buffered_symbol_bytes: 512,
1362 ..ReceiverConfig::default()
1363 });
1364
1365 let first = make_packet(
1366 ChangesetId::from_bytes([0x41; 16]),
1367 0,
1368 0,
1369 100,
1370 vec![0x55; 400],
1371 );
1372 receiver
1373 .process_parsed_packet(&first)
1374 .expect("first packet");
1375 assert_eq!(receiver.active_decoders(), 1);
1376
1377 let second = make_packet(
1379 ChangesetId::from_bytes([0x42; 16]),
1380 0,
1381 0,
1382 100,
1383 vec![0x77; 200],
1384 );
1385 let err = receiver.process_parsed_packet(&second).unwrap_err();
1386 assert!(matches!(err, FrankenError::TooBig));
1387 assert_eq!(receiver.active_decoders(), 1);
1388 }
1389
1390 #[test]
1391 fn test_receiver_seed_derivation() {
1392 let id = ChangesetId::from_bytes([1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15, 16]);
1394 let seed = derive_seed_from_changeset_id(&id);
1395
1396 let expected = xxhash_rust::xxh3::xxh3_64(id.as_bytes());
1397 assert_eq!(
1398 seed, expected,
1399 "bead_id={TEST_BEAD_ID} case=seed_matches_sender"
1400 );
1401 }
1402
1403 #[test]
1404 fn test_receiver_v1_reject_sbn_nonzero() {
1405 let mut receiver = ReplicationReceiver::new();
1406 let packet = make_packet(
1407 ChangesetId::from_bytes([0xAA; 16]),
1408 1, 0,
1410 10,
1411 vec![0x55; 512],
1412 );
1413 let wire = packet.to_bytes().expect("encode");
1414 let result = receiver.process_packet(&wire);
1415 assert!(
1416 result.is_err(),
1417 "bead_id={TEST_BEAD_ID} case=v1_sbn_rejected"
1418 );
1419 }
1420
1421 #[test]
1422 fn test_receiver_k_source_validation() {
1423 let mut receiver = ReplicationReceiver::new();
1424
1425 let packet_zero = make_packet(
1427 ChangesetId::from_bytes([0xBB; 16]),
1428 0,
1429 0,
1430 0,
1431 vec![0x55; 512],
1432 );
1433 let wire_zero = packet_zero.to_bytes().expect("encode");
1434 assert!(
1435 receiver.process_packet(&wire_zero).is_err(),
1436 "bead_id={TEST_BEAD_ID} case=k_source_zero_rejected"
1437 );
1438
1439 let packet_over = make_packet(
1441 ChangesetId::from_bytes([0xCC; 16]),
1442 0,
1443 0,
1444 K_MAX + 1,
1445 vec![0x55; 512],
1446 );
1447 let result = receiver.process_parsed_packet(&packet_over);
1450 assert!(
1451 result.is_err(),
1452 "bead_id={TEST_BEAD_ID} case=k_source_over_max_rejected"
1453 );
1454
1455 let packet_max = make_packet(
1457 ChangesetId::from_bytes([0xDD; 16]),
1458 0,
1459 0,
1460 K_MAX,
1461 vec![0x55; 512],
1462 );
1463 let result = receiver.process_parsed_packet(&packet_max);
1464 assert!(
1465 result.is_ok(),
1466 "bead_id={TEST_BEAD_ID} case=k_source_at_max_accepted"
1467 );
1468 }
1469
1470 #[test]
1471 fn test_receiver_symbol_size_inference() {
1472 let mut receiver = ReplicationReceiver::new();
1473 let packet = make_packet(
1474 ChangesetId::from_bytes([0xEE; 16]),
1475 0,
1476 0,
1477 100,
1478 vec![0x42; 1024],
1479 );
1480 receiver
1481 .process_parsed_packet(&packet)
1482 .expect("accept packet");
1483
1484 let decoder = receiver
1486 .decoders
1487 .get(&packet.changeset_id)
1488 .expect("decoder exists");
1489 assert_eq!(
1490 decoder.symbol_size, 1024,
1491 "bead_id={TEST_BEAD_ID} case=symbol_size_inferred"
1492 );
1493
1494 let mut receiver2 = ReplicationReceiver::new();
1496 let empty_packet = make_packet(ChangesetId::from_bytes([0xFF; 16]), 0, 0, 10, vec![]);
1497 assert!(
1498 receiver2.process_parsed_packet(&empty_packet).is_err(),
1499 "bead_id={TEST_BEAD_ID} case=zero_symbol_size_rejected"
1500 );
1501 }
1502
1503 #[test]
1504 fn test_receiver_k_source_mismatch_rejected() {
1505 let mut receiver = ReplicationReceiver::new();
1506 let id = ChangesetId::from_bytes([0x11; 16]);
1507
1508 let p1 = make_packet(id, 0, 0, 100, vec![0x42; 512]);
1509 receiver
1510 .process_parsed_packet(&p1)
1511 .expect("first packet ok");
1512
1513 let p2 = make_packet(id, 0, 1, 200, vec![0x42; 512]); assert!(
1516 receiver.process_parsed_packet(&p2).is_err(),
1517 "bead_id={TEST_BEAD_ID} case=k_source_mismatch_rejected"
1518 );
1519 }
1520
1521 #[test]
1522 fn test_receiver_symbol_size_mismatch_rejected() {
1523 let mut receiver = ReplicationReceiver::new();
1524 let id = ChangesetId::from_bytes([0x22; 16]);
1525
1526 let p1 = make_packet(id, 0, 0, 100, vec![0x42; 512]);
1527 receiver
1528 .process_parsed_packet(&p1)
1529 .expect("first packet ok");
1530
1531 let p2 = make_packet(id, 0, 1, 100, vec![0x42; 1024]); assert!(
1534 receiver.process_parsed_packet(&p2).is_err(),
1535 "bead_id={TEST_BEAD_ID} case=symbol_size_mismatch_rejected"
1536 );
1537 }
1538
1539 #[test]
1540 fn test_receiver_isi_deduplication() {
1541 let mut receiver = ReplicationReceiver::new();
1542 let id = ChangesetId::from_bytes([0x33; 16]);
1543
1544 let p1 = make_packet(id, 0, 0, 100, vec![0x42; 512]);
1545
1546 let r1 = receiver.process_parsed_packet(&p1).expect("first");
1547 assert_eq!(
1548 r1,
1549 PacketResult::Accepted,
1550 "bead_id={TEST_BEAD_ID} case=first_accepted"
1551 );
1552
1553 let r2 = receiver.process_parsed_packet(&p1).expect("duplicate");
1555 assert_eq!(
1556 r2,
1557 PacketResult::Duplicate,
1558 "bead_id={TEST_BEAD_ID} case=isi_dedup"
1559 );
1560
1561 let count = receiver.received_counts.get(&id).copied().unwrap_or(0);
1563 assert_eq!(
1564 count, 1,
1565 "bead_id={TEST_BEAD_ID} case=dedup_count_unchanged"
1566 );
1567 }
1568
1569 #[test]
1570 fn test_receiver_treats_payload_hash_mismatch_as_erasure() {
1571 let mut receiver = ReplicationReceiver::new();
1572 let packet = make_packet(
1573 ChangesetId::from_bytes([0x44; 16]),
1574 0,
1575 0,
1576 100,
1577 vec![0x42; 512],
1578 );
1579 let mut wire = packet.to_bytes().expect("encode packet");
1580 wire[48] ^= 0xFF;
1581 let result = receiver.process_packet(&wire).expect("process packet");
1582 assert_eq!(result, PacketResult::Erasure);
1583 }
1584
1585 #[test]
1586 fn test_receiver_treats_invalid_auth_tag_as_erasure() {
1587 let receiver_key = [0x11_u8; 32];
1588 let sender_key = [0x22_u8; 32];
1589 let mut receiver =
1590 ReplicationReceiver::with_config(ReceiverConfig::with_auth_key(receiver_key));
1591 let mut packet = make_packet(
1592 ChangesetId::from_bytes([0x45; 16]),
1593 0,
1594 0,
1595 100,
1596 vec![0x24; 512],
1597 );
1598 packet.attach_auth_tag(&sender_key);
1599 let wire = packet.to_bytes().expect("encode auth packet");
1600 let result = receiver.process_packet(&wire).expect("process packet");
1601 assert_eq!(result, PacketResult::Erasure);
1602 }
1603
1604 #[test]
1605 fn test_receiver_accepts_legacy_v1_packets() {
1606 let mut receiver = ReplicationReceiver::new();
1607 let id = ChangesetId::from_bytes([0x46; 16]);
1608 let symbol_data = vec![0x5A; 512];
1609 let legacy = ReplicationPacket {
1610 wire_version: ReplicationWireVersion::LegacyV1,
1611 changeset_id: id,
1612 sbn: 0,
1613 esi: 0,
1614 k_source: 100,
1615 r_repair: 0,
1616 symbol_size_t: 512,
1617 seed: derive_seed_from_changeset_id(&id),
1618 payload_xxh3: ReplicationPacket::compute_payload_xxh3(&symbol_data),
1619 auth_tag: None,
1620 symbol_data,
1621 };
1622 let wire = legacy.to_bytes().expect("encode legacy packet");
1623 let parsed = ReplicationPacket::from_bytes(&wire).expect("decode legacy packet");
1624 assert_eq!(parsed.wire_version, ReplicationWireVersion::LegacyV1);
1625 let result = receiver
1626 .process_packet(&wire)
1627 .expect("process legacy packet");
1628 assert_eq!(result, PacketResult::Accepted);
1629 }
1630
1631 #[test]
1632 fn test_receiver_decode_at_k_source() {
1633 let page_size = 512_u32;
1635 let mut receiver = ReplicationReceiver::new();
1636 let packets = generate_sender_packets(page_size, &[1, 2, 3], 512);
1637
1638 let mut last_result = PacketResult::Accepted;
1639 for pkt in &packets {
1640 let result = receiver
1641 .process_packet(pkt)
1642 .expect("bead_id={TEST_BEAD_ID} case=decode_at_k unexpected error");
1643 last_result = result;
1644 }
1645
1646 assert_eq!(
1647 last_result,
1648 PacketResult::DecodeReady,
1649 "bead_id={TEST_BEAD_ID} case=decode_triggers_at_k_source"
1650 );
1651 assert_eq!(
1652 receiver.state(),
1653 ReceiverState::Applying,
1654 "bead_id={TEST_BEAD_ID} case=state_applying_after_decode"
1655 );
1656 }
1657
1658 #[test]
1659 fn test_receiver_decode_failure_emits_proof_when_enabled() {
1660 let mut receiver = receiver_with_decode_proofs();
1661 let changeset_id = ChangesetId::from_bytes([0x5A; 16]);
1662
1663 let p1 = make_packet(changeset_id, 0, 2, 2, vec![0xA1; 64]);
1665 let p2 = make_packet(changeset_id, 0, 3, 2, vec![0xA2; 64]);
1666
1667 let r1 = receiver.process_parsed_packet(&p1).expect("first packet");
1668 assert_eq!(r1, PacketResult::Accepted);
1669 let r2 = receiver.process_parsed_packet(&p2).expect("second packet");
1670 assert_eq!(r2, PacketResult::NeedMore);
1671
1672 let audit = receiver.take_decode_audit_entries();
1673 assert_eq!(audit.len(), 1, "bead_id=bd-faz4 case=failure_proof_emitted");
1674 let proof = &audit[0].proof;
1675 assert!(
1676 !proof.decode_success,
1677 "bead_id=bd-faz4 case=failure_proof_decode_success_false"
1678 );
1679 assert_eq!(proof.changeset_id, Some(*changeset_id.as_bytes()));
1680 assert!(
1681 proof.is_consistent(),
1682 "bead_id=bd-faz4 case=failure_proof_consistent"
1683 );
1684 }
1685
1686 #[test]
1687 fn test_receiver_decode_success_with_repair_emits_proof_when_enabled() {
1688 let mut receiver = ReplicationReceiver::with_config(ReceiverConfig {
1689 auth_key: None,
1690 decode_proof_policy: DecodeProofEmissionPolicy {
1691 emit_on_decode_failure: false,
1692 emit_on_repair_success: true,
1693 },
1694 ..ReceiverConfig::default()
1695 });
1696 let page_size = 64_u32;
1697 let mut pages = make_pages(page_size, &[7]);
1698 let changeset_bytes = encode_changeset(page_size, &mut pages).expect("encode changeset");
1699 let changeset_id = compute_changeset_id(&changeset_bytes);
1700
1701 let symbol_size = 64_usize;
1703 let mut s0 = vec![0_u8; symbol_size];
1704 let mut s1 = vec![0_u8; symbol_size];
1705 let split = changeset_bytes.len().min(symbol_size);
1706 s0[..split].copy_from_slice(&changeset_bytes[..split]);
1707 if changeset_bytes.len() > symbol_size {
1708 let rem = changeset_bytes.len() - symbol_size;
1709 s1[..rem].copy_from_slice(&changeset_bytes[symbol_size..]);
1710 }
1711
1712 let p0 = make_packet(changeset_id, 0, 0, 2, s0);
1714 let p_repair = make_packet(changeset_id, 0, 2, 2, vec![0xCC; symbol_size]);
1715 let p1 = make_packet(changeset_id, 0, 1, 2, s1);
1716
1717 assert_eq!(
1718 receiver.process_parsed_packet(&p0).expect("p0"),
1719 PacketResult::Accepted
1720 );
1721 assert_eq!(
1722 receiver.process_parsed_packet(&p_repair).expect("repair"),
1723 PacketResult::NeedMore
1724 );
1725 assert_eq!(
1726 receiver.process_parsed_packet(&p1).expect("p1"),
1727 PacketResult::DecodeReady
1728 );
1729 assert_eq!(receiver.state(), ReceiverState::Applying);
1730
1731 let results = receiver.apply_pending().expect("apply");
1732 assert_eq!(results.len(), 1);
1733 let decode_proof = results[0]
1734 .decode_proof
1735 .as_ref()
1736 .expect("bead_id=bd-faz4 case=success_proof_attached_to_result");
1737 assert!(decode_proof.decode_success);
1738 assert!(decode_proof.is_repair());
1739 assert!(
1740 decode_proof.is_consistent(),
1741 "bead_id=bd-faz4 case=success_proof_consistent"
1742 );
1743
1744 let audit = receiver.take_decode_audit_entries();
1745 assert_eq!(audit.len(), 1, "bead_id=bd-faz4 case=success_proof_emitted");
1746 }
1747
1748 #[test]
1749 fn test_receiver_decode_success_truncation() {
1750 let page_size = 128_u32;
1751 let mut receiver = ReplicationReceiver::new();
1752 let packets = generate_sender_packets(page_size, &[1], 128);
1753
1754 for pkt in &packets {
1755 let _ = receiver.process_packet(pkt);
1756 }
1757
1758 if receiver.state() == ReceiverState::Applying {
1760 let results = receiver.apply_pending().expect("apply");
1761 assert!(
1762 !results.is_empty(),
1763 "bead_id={TEST_BEAD_ID} case=has_results"
1764 );
1765 for result in &results {
1766 for page in &result.pages {
1767 assert_eq!(
1768 page.page_data.len(),
1769 page_size as usize,
1770 "bead_id={TEST_BEAD_ID} case=page_data_correct_size"
1771 );
1772 }
1773 }
1774 }
1775 }
1776
1777 #[test]
1778 fn test_receiver_page_xxh3_validation() {
1779 let page_size = 256_u32;
1780 let mut pages = make_pages(page_size, &[1]);
1781 let changeset_bytes = encode_changeset(page_size, &mut pages).expect("encode");
1782
1783 let mut tampered = changeset_bytes.clone();
1785 let tamper_offset = CHANGESET_HEADER_SIZE + 4 + 8 + 10; tampered[tamper_offset] ^= 0xFF;
1787
1788 let receiver = ReplicationReceiver::new();
1790 let changeset_id = compute_changeset_id(&changeset_bytes);
1791 let result = receiver.parse_and_validate_changeset(changeset_id, &tampered);
1792 assert!(
1793 result.is_err(),
1794 "bead_id={TEST_BEAD_ID} case=xxh3_validation_catches_corruption"
1795 );
1796 }
1797
1798 #[test]
1799 fn test_parse_and_validate_rejects_changeset_id_mismatch() {
1800 let page_size = 128_u32;
1801 let mut pages = make_pages(page_size, &[1]);
1802 let changeset_bytes = encode_changeset(page_size, &mut pages).expect("encode");
1803 let wrong_changeset_id = ChangesetId::from_bytes([0x42; 16]);
1804
1805 let receiver = ReplicationReceiver::new();
1806 let result = receiver.parse_and_validate_changeset(wrong_changeset_id, &changeset_bytes);
1807 assert!(
1808 matches!(result, Err(FrankenError::DatabaseCorrupt { .. })),
1809 "bead_id={TEST_BEAD_ID} case=changeset_id_mismatch_rejected"
1810 );
1811 }
1812
1813 #[test]
1814 fn test_parse_and_validate_rejects_total_len_smaller_than_header() {
1815 let receiver = ReplicationReceiver::new();
1816 let changeset_id = ChangesetId::from_bytes([0xA5; 16]);
1817
1818 let mut malformed = vec![0_u8; CHANGESET_HEADER_SIZE];
1819 malformed[0..4].copy_from_slice(b"FSRP");
1820 malformed[4..6].copy_from_slice(&1_u16.to_le_bytes());
1821 malformed[6..10].copy_from_slice(&4096_u32.to_le_bytes());
1822 malformed[10..14].copy_from_slice(&1_u32.to_le_bytes());
1823 malformed[14..22].copy_from_slice(&1_u64.to_le_bytes());
1824
1825 let result = receiver.parse_and_validate_changeset(changeset_id, &malformed);
1826 assert!(matches!(result, Err(FrankenError::DatabaseCorrupt { .. })));
1827 }
1828
1829 #[test]
1830 fn test_parse_and_validate_rejects_trailing_payload_bytes() {
1831 let page_size = 128_u32;
1832 let mut pages = make_pages(page_size, &[1]);
1833 let mut malformed = encode_changeset(page_size, &mut pages).expect("encode");
1834 malformed.push(0x99);
1835 let total_len = u64::try_from(malformed.len()).expect("test total_len fits u64");
1836 malformed[14..22].copy_from_slice(&total_len.to_le_bytes());
1837 let changeset_id = compute_changeset_id(&malformed);
1838
1839 let receiver = ReplicationReceiver::new();
1840 let result = receiver.parse_and_validate_changeset(changeset_id, &malformed);
1841 assert!(
1842 matches!(result, Err(FrankenError::DatabaseCorrupt { .. })),
1843 "bead_id={TEST_BEAD_ID} case=parse_rejects_trailing_payload"
1844 );
1845 }
1846
1847 #[test]
1848 fn test_parse_changeset_pages_rejects_truncated_payload() {
1849 let total_len = CHANGESET_HEADER_SIZE + 8;
1850 let mut malformed = vec![0_u8; total_len];
1851 malformed[0..4].copy_from_slice(b"FSRP");
1852 malformed[4..6].copy_from_slice(&1_u16.to_le_bytes());
1853 malformed[6..10].copy_from_slice(&4096_u32.to_le_bytes());
1854 malformed[10..14].copy_from_slice(&1_u32.to_le_bytes());
1855 malformed[14..22].copy_from_slice(
1856 &u64::try_from(total_len)
1857 .expect("test total_len fits into u64")
1858 .to_le_bytes(),
1859 );
1860
1861 let result = parse_changeset_pages(&malformed);
1862 assert!(matches!(result, Err(FrankenError::DatabaseCorrupt { .. })));
1863 }
1864
1865 #[test]
1866 fn test_parse_changeset_pages_rejects_trailing_payload() {
1867 let page_size = 128_u32;
1868 let mut pages = make_pages(page_size, &[1]);
1869 let mut malformed = encode_changeset(page_size, &mut pages).expect("encode");
1870 malformed.push(0xA5);
1871 let total_len = u64::try_from(malformed.len()).expect("test total_len fits u64");
1872 malformed[14..22].copy_from_slice(&total_len.to_le_bytes());
1873
1874 let result = parse_changeset_pages(&malformed);
1875 assert!(
1876 matches!(result, Err(FrankenError::DatabaseCorrupt { .. })),
1877 "bead_id={TEST_BEAD_ID} case=parse_pages_rejects_trailing_payload"
1878 );
1879 }
1880
1881 #[test]
1882 fn test_receiver_pages_applied_in_order() {
1883 let page_size = 256_u32;
1884 let mut receiver = ReplicationReceiver::new();
1885 let packets = generate_sender_packets(page_size, &[5, 1, 3, 2, 4], 256);
1886
1887 for pkt in &packets {
1888 let _ = receiver.process_packet(pkt);
1889 }
1890
1891 if receiver.state() == ReceiverState::Applying {
1892 let results = receiver.apply_pending().expect("apply");
1893 let pages = &results[0].pages;
1894 for w in pages.windows(2) {
1895 assert!(
1896 w[0].page_number <= w[1].page_number,
1897 "bead_id={TEST_BEAD_ID} case=pages_sorted pn0={} pn1={}",
1898 w[0].page_number,
1899 w[1].page_number
1900 );
1901 }
1902 }
1903 }
1904
1905 #[test]
1910 fn prop_any_k_symbols_decode() {
1911 for n_pages in [1_u32, 3, 5, 10] {
1914 let page_size = 256_u32;
1915 let mut receiver = ReplicationReceiver::new();
1916 let packets =
1917 generate_sender_packets(page_size, &(1..=n_pages).collect::<Vec<_>>(), 256);
1918
1919 let mut decode_ready = false;
1920 for pkt in &packets {
1921 if matches!(receiver.process_packet(pkt), Ok(PacketResult::DecodeReady)) {
1922 decode_ready = true;
1923 break;
1924 }
1925 }
1926 assert!(
1927 decode_ready,
1928 "bead_id={TEST_BEAD_ID} case=prop_any_k_decode n_pages={n_pages}"
1929 );
1930 }
1931 }
1932
1933 #[test]
1934 fn prop_dedup_idempotent() {
1935 let mut receiver = ReplicationReceiver::new();
1937 let id = ChangesetId::from_bytes([0x77; 16]);
1938
1939 let p1 = make_packet(id, 0, 0, 100, vec![0x42; 512]); let r1 = receiver.process_parsed_packet(&p1).expect("first");
1943 assert_eq!(
1944 r1,
1945 PacketResult::Accepted,
1946 "bead_id={TEST_BEAD_ID} case=dedup_first_accepted"
1947 );
1948
1949 for _ in 0..5 {
1950 let r = receiver.process_parsed_packet(&p1).expect("duplicate");
1951 assert_eq!(
1952 r,
1953 PacketResult::Duplicate,
1954 "bead_id={TEST_BEAD_ID} case=dedup_subsequent_always_duplicate"
1955 );
1956 }
1957
1958 let count = receiver.received_counts.get(&id).copied().unwrap_or(0);
1960 assert_eq!(count, 1, "bead_id={TEST_BEAD_ID} case=dedup_count_stable");
1961 }
1962
1963 #[test]
1968 fn test_packet_reject_over_message_cap() {
1969 let mut receiver = ReplicationReceiver::new();
1970 let oversized = vec![0_u8; DEFAULT_RPC_MESSAGE_CAP_BYTES + 1];
1971 let err = receiver.process_packet(&oversized).unwrap_err();
1972 assert!(matches!(err, FrankenError::TooBig));
1973 }
1974
1975 #[test]
1976 fn test_e2e_sender_receiver_roundtrip() {
1977 let page_size = 512_u32;
1979 let page_numbers: Vec<u32> = (1..=20).collect();
1980 let original_pages = make_pages(page_size, &page_numbers);
1981
1982 let mut receiver = ReplicationReceiver::new();
1983 let packets = generate_sender_packets(page_size, &page_numbers, 512);
1984
1985 for pkt in &packets {
1986 let _ = receiver.process_packet(pkt);
1987 }
1988
1989 assert_eq!(
1990 receiver.state(),
1991 ReceiverState::Applying,
1992 "bead_id={TEST_BEAD_ID} case=e2e_roundtrip_applying"
1993 );
1994
1995 let results = receiver.apply_pending().expect("apply");
1996 assert_eq!(
1997 results.len(),
1998 1,
1999 "bead_id={TEST_BEAD_ID} case=e2e_one_changeset"
2000 );
2001
2002 let decoded_pages = &results[0].pages;
2003 assert_eq!(
2004 decoded_pages.len(),
2005 original_pages.len(),
2006 "bead_id={TEST_BEAD_ID} case=e2e_page_count"
2007 );
2008
2009 for (decoded, original) in decoded_pages.iter().zip(original_pages.iter()) {
2010 assert_eq!(
2011 decoded.page_number, original.page_number,
2012 "bead_id={TEST_BEAD_ID} case=e2e_page_number_match"
2013 );
2014 assert_eq!(
2015 decoded.page_data, original.page_bytes,
2016 "bead_id={TEST_BEAD_ID} case=e2e_page_data_identical pn={}",
2017 original.page_number
2018 );
2019 }
2020
2021 receiver.reset_to_listening().expect("reset");
2023 assert_eq!(
2024 receiver.state(),
2025 ReceiverState::Listening,
2026 "bead_id={TEST_BEAD_ID} case=e2e_back_to_listening"
2027 );
2028 }
2029
2030 #[test]
2031 fn test_e2e_concurrent_changesets() {
2032 let mut receiver = ReplicationReceiver::new();
2034
2035 let packets_a = generate_sender_packets(256, &[1, 2, 3], 256);
2036 let packets_b = generate_sender_packets(256, &[10, 20, 30], 256);
2037
2038 let mut all_packets = Vec::new();
2040 let max_len = packets_a.len().max(packets_b.len());
2041 for i in 0..max_len {
2042 if i < packets_a.len() {
2043 all_packets.push(packets_a[i].clone());
2044 }
2045 if i < packets_b.len() {
2046 all_packets.push(packets_b[i].clone());
2047 }
2048 }
2049
2050 let mut decode_count = 0_u32;
2051 for pkt in &all_packets {
2052 if matches!(receiver.process_packet(pkt), Ok(PacketResult::DecodeReady)) {
2053 decode_count += 1;
2054 if receiver.state() == ReceiverState::Applying {
2056 let _ = receiver.apply_pending();
2057 if !receiver.decoders.is_empty() {
2059 receiver.state = ReceiverState::Collecting;
2060 }
2061 }
2062 }
2063 }
2064
2065 assert!(
2066 decode_count >= 1,
2067 "bead_id={TEST_BEAD_ID} case=e2e_concurrent_at_least_one_decoded count={decode_count}"
2068 );
2069 }
2070
2071 #[test]
2072 fn test_e2e_bd_1hi_14_compliance() {
2073 let page_size = 1024_u32;
2075 let page_numbers: Vec<u32> = (1..=10).collect();
2076 let original_pages = make_pages(page_size, &page_numbers);
2077
2078 let mut sender = ReplicationSender::new();
2080 let mut pages = make_pages(page_size, &page_numbers);
2081 sender
2082 .prepare(page_size, &mut pages, SenderConfig::default())
2083 .expect("prepare");
2084 sender.start_streaming().expect("start");
2085
2086 let mut wire_packets = Vec::new();
2088 while let Some(packet) = sender.next_packet().expect("next") {
2089 wire_packets.push(packet.to_bytes().expect("encode"));
2090 }
2091
2092 let mut receiver = ReplicationReceiver::new();
2094 assert_eq!(receiver.state(), ReceiverState::Listening);
2095
2096 let mut last_result = PacketResult::Accepted;
2097 for pkt in &wire_packets {
2098 let result = receiver
2099 .process_packet(pkt)
2100 .expect("bead_id={TEST_BEAD_ID} case=e2e_compliance unexpected error");
2101 last_result = result;
2102 if result == PacketResult::DecodeReady {
2103 break;
2104 }
2105 }
2106
2107 assert_eq!(
2109 last_result,
2110 PacketResult::DecodeReady,
2111 "bead_id={TEST_BEAD_ID} case=e2e_compliance_decoded"
2112 );
2113 assert_eq!(receiver.state(), ReceiverState::Applying);
2114
2115 let results = receiver.apply_pending().expect("apply");
2117 assert_eq!(receiver.state(), ReceiverState::Complete);
2118 assert_eq!(results.len(), 1);
2119
2120 let decoded = &results[0].pages;
2122 assert_eq!(decoded.len(), original_pages.len());
2123 for (d, o) in decoded.iter().zip(original_pages.iter()) {
2124 assert_eq!(d.page_number, o.page_number);
2125 assert_eq!(d.page_data, o.page_bytes);
2126 }
2127
2128 receiver.reset_to_listening().expect("reset");
2130 assert_eq!(
2131 receiver.state(),
2132 ReceiverState::Listening,
2133 "bead_id={TEST_BEAD_ID} case=e2e_compliance_reset"
2134 );
2135 assert_eq!(receiver.applied_count(), 1);
2136 }
2137
2138 #[test]
2139 fn test_simnetwork_loss_profiles_converge_with_repair_symbols() {
2140 let page_size = 128_u32;
2141 let page_numbers = [1_u32, 2];
2142 let original_pages = make_pages(page_size, &page_numbers);
2143 let packets = generate_sender_packets_with_multiplier(page_size, &page_numbers, 128, 2);
2144 let loss_packets: Vec<Vec<u8>> = packets
2145 .iter()
2146 .flat_map(|packet| [packet.clone(), packet.clone()])
2147 .collect();
2148
2149 for (loss_rate, require_observed_drop) in [(0.05_f64, false), (0.30_f64, true)] {
2150 let mut found_seed = None;
2151 for seed in 1_u64..=20_000 {
2152 let mut config = SimTransportConfig::deterministic(seed);
2153 config.loss_rate = loss_rate;
2154 config.preserve_order = true;
2155
2156 let delivery = transmit_packets_simnetwork(config, &loss_packets);
2157 let observed_drop = delivery.delivered.len() < delivery.sent_count;
2158 if require_observed_drop && !observed_drop {
2159 continue;
2160 }
2161 let saw_repair_symbol = delivery.delivered.iter().any(|(_, wire)| {
2162 ReplicationPacket::from_bytes(wire)
2163 .is_ok_and(|packet| !packet.is_source_symbol())
2164 });
2165 if !saw_repair_symbol {
2166 continue;
2167 }
2168
2169 let (decoded, _erasures, _parse_errors) =
2170 decode_from_wire_packets(&delivery.delivered);
2171 if decoded
2172 .as_ref()
2173 .is_some_and(|pages| decoded_matches_original(pages, &original_pages))
2174 {
2175 found_seed = Some(seed);
2176 break;
2177 }
2178 }
2179
2180 assert!(
2181 found_seed.is_some(),
2182 "bead_id=bd-xgoe case=loss_profile_convergence loss_rate={loss_rate} require_drop={require_observed_drop} did not find deterministic convergent seed"
2183 );
2184 }
2185 }
2186
2187 #[test]
2188 fn test_simnetwork_reorder_and_dup_converge() {
2189 let page_size = 128_u32;
2190 let page_numbers = [7_u32, 11];
2191 let original_pages = make_pages(page_size, &page_numbers);
2192 let packets = generate_sender_packets_with_multiplier(page_size, &page_numbers, 128, 2);
2193
2194 let mut found_seed = None;
2195 for seed in 1_u64..=2_000 {
2196 let mut config = SimTransportConfig::deterministic(seed);
2197 config.preserve_order = false;
2198 config.duplication_rate = 0.35;
2199
2200 let delivery = transmit_packets_simnetwork(config, &packets);
2201 if !has_duplicate_esies(&delivery) || !has_reordered_esies(&delivery) {
2202 continue;
2203 }
2204
2205 let (decoded, _erasures, _parse_errors) = decode_from_wire_packets(&delivery.delivered);
2206 if decoded
2207 .as_ref()
2208 .is_some_and(|pages| decoded_matches_original(pages, &original_pages))
2209 {
2210 found_seed = Some(seed);
2211 break;
2212 }
2213 }
2214
2215 assert!(
2216 found_seed.is_some(),
2217 "bead_id=bd-xgoe case=reorder_dup_convergence no deterministic seed achieved reorder+dup convergence"
2218 );
2219 }
2220
2221 #[test]
2222 fn test_simnetwork_corruption_is_rejected_and_recovered() {
2223 let page_size = 128_u32;
2224 let page_numbers = [21_u32, 34];
2225 let original_pages = make_pages(page_size, &page_numbers);
2226 let packets = generate_sender_packets_with_multiplier(page_size, &page_numbers, 128, 2);
2227
2228 let mut found_seed = None;
2229 for seed in 1_u64..=20_000 {
2230 let mut config = SimTransportConfig::deterministic(seed);
2231 config.corruption_rate = 0.20;
2232 config.preserve_order = false;
2233
2234 let delivery = transmit_packets_simnetwork(config, &packets);
2235 if !has_corrupted_wire_bytes(&delivery, &packets) {
2236 continue;
2237 }
2238
2239 let (decoded, erasures, parse_errors) = decode_from_wire_packets(&delivery.delivered);
2240 if erasures + parse_errors == 0 {
2241 continue;
2242 }
2243 if decoded
2244 .as_ref()
2245 .is_some_and(|pages| decoded_matches_original(pages, &original_pages))
2246 {
2247 found_seed = Some(seed);
2248 break;
2249 }
2250 }
2251
2252 assert!(
2253 found_seed.is_some(),
2254 "bead_id=bd-xgoe case=corruption_recovery no deterministic seed achieved corruption rejection + convergence"
2255 );
2256 }
2257
2258 #[test]
2259 fn test_simnetwork_stop_early_reduces_traffic() {
2260 let page_size = 256_u32;
2261 let page_numbers = [1_u32, 2, 3];
2262 let packets = generate_sender_packets_with_multiplier(page_size, &page_numbers, 256, 2);
2263
2264 let full_delivery = transmit_packets_simnetwork(SimTransportConfig::reliable(), &packets);
2265 let full_sent = full_delivery.sent_count;
2266
2267 let network = SimNetwork::fully_connected(2, SimTransportConfig::reliable());
2268 let (mut sink, mut stream) = network.transport(0, 1);
2269 let runtime = RuntimeBuilder::current_thread()
2270 .build()
2271 .expect("runtime build");
2272
2273 let mut receiver = ReplicationReceiver::new();
2274 let mut stop_early_sent = 0_usize;
2275 let mut decoded = false;
2276
2277 runtime.block_on(async {
2278 for (index, bytes) in packets.iter().enumerate() {
2279 let esi = u32::try_from(index).expect("test packet index fits u32");
2280 sink.send(packet_symbol(esi, bytes.clone()))
2281 .await
2282 .expect("send simulated symbol");
2283 stop_early_sent += 1;
2284
2285 let delivered = stream
2286 .next()
2287 .await
2288 .expect("delivered packet")
2289 .expect("stream item");
2290 let wire = delivered.symbol().data().to_vec();
2291 if matches!(
2292 receiver.process_packet(&wire).expect("receiver process"),
2293 PacketResult::DecodeReady
2294 ) {
2295 decoded = true;
2296 break;
2297 }
2298 }
2299 sink.close().await.expect("close simulated sink");
2300 });
2301
2302 assert!(
2303 decoded,
2304 "bead_id=bd-xgoe case=stop_early_decode_not_reached"
2305 );
2306 assert!(
2307 stop_early_sent < full_sent,
2308 "bead_id=bd-xgoe case=stop_early_not_reduced stop_early_sent={stop_early_sent} full_sent={full_sent}"
2309 );
2310 }
2311
2312 #[test]
2317 fn test_bd_1hi_14_unit_compliance_gate() {
2318 let _ = ReceiverState::Listening;
2320 let _ = ReceiverState::Collecting;
2321 let _ = ReceiverState::Decoding;
2322 let _ = ReceiverState::Applying;
2323 let _ = ReceiverState::Complete;
2324
2325 let _ = PacketResult::Accepted;
2326 let _ = PacketResult::Erasure;
2327 let _ = PacketResult::Duplicate;
2328 let _ = PacketResult::DecodeReady;
2329 let _ = PacketResult::NeedMore;
2330
2331 let receiver = ReplicationReceiver::new();
2332 assert_eq!(receiver.state(), ReceiverState::Listening);
2333 assert_eq!(receiver.applied_count(), 0);
2334 assert_eq!(receiver.active_decoders(), 0);
2335
2336 assert_eq!(REPLICATION_HEADER_SIZE, 72);
2338 }
2339
2340 #[test]
2341 fn prop_bd_1hi_14_structure_compliance() {
2342 let page_size = 256_u32;
2344 let mut receiver = ReplicationReceiver::new();
2345 assert_eq!(receiver.state(), ReceiverState::Listening);
2346
2347 let packets = generate_sender_packets(page_size, &[1, 2], 256);
2348 for pkt in &packets {
2349 let _ = receiver.process_packet(pkt);
2350 }
2351
2352 assert!(
2354 receiver.state() == ReceiverState::Applying
2355 || receiver.state() == ReceiverState::Collecting,
2356 "bead_id={TEST_BEAD_ID} case=prop_state_machine state={:?}",
2357 receiver.state()
2358 );
2359
2360 if receiver.state() == ReceiverState::Applying {
2361 let results = receiver.apply_pending().expect("apply");
2362 assert!(!results.is_empty());
2363 assert_eq!(receiver.state(), ReceiverState::Complete);
2364 receiver.reset_to_listening().expect("reset");
2365 assert_eq!(receiver.state(), ReceiverState::Listening);
2366 }
2367 }
2368}