Skip to main content

fsqlite_core/
replication_receiver.rs

1//! §3.4.2 Fountain-Coded Replication Receiver (bd-1hi.14).
2//!
3//! Implements the receiver-side state machine for fountain-coded database
4//! replication. Listens for UDP packets, collects symbols per changeset,
5//! decodes when sufficient, validates and applies recovered pages.
6//!
7//! State machine: LISTENING → COLLECTING → DECODING → APPLYING → COMPLETE
8
9use 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// ---------------------------------------------------------------------------
27// Receiver State Machine
28// ---------------------------------------------------------------------------
29
30/// Receiver state (§3.4.2).
31#[derive(Debug, Clone, Copy, PartialEq, Eq)]
32pub enum ReceiverState {
33    /// Ready to accept replication data.
34    Listening,
35    /// At least one packet received; collecting symbols.
36    Collecting,
37    /// Sufficient symbols collected; decoding in progress.
38    Decoding,
39    /// Pages decoded; applying to local database.
40    Applying,
41    /// All pages applied; ready for next changeset.
42    Complete,
43}
44
45/// Per-changeset decoder state, created on first packet.
46#[derive(Debug)]
47pub struct DecoderState {
48    /// Number of source symbols expected.
49    pub k_source: u32,
50    /// Symbol size in bytes (inferred from first packet).
51    pub symbol_size: u32,
52    /// Deterministic seed derived from changeset_id.
53    pub seed: u64,
54    /// Collected symbols indexed by ISI.
55    symbols: HashMap<u32, Vec<u8>>,
56    /// Set of received ISIs for O(1) deduplication.
57    received_isis: HashSet<u32>,
58}
59
60impl DecoderState {
61    /// Create a new decoder state for a changeset.
62    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    /// Number of unique symbols received.
73    #[must_use]
74    pub fn received_count(&self) -> u32 {
75        u32::try_from(self.received_isis.len()).unwrap_or(u32::MAX)
76    }
77
78    /// Whether enough symbols have been collected to attempt decode.
79    #[must_use]
80    pub fn ready_to_decode(&self) -> bool {
81        self.received_count() >= self.k_source
82    }
83
84    /// Number of collected source symbols (`isi < k_source`).
85    #[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    /// Whether any collected symbol is a repair symbol (`isi >= k_source`).
96    #[must_use]
97    pub fn has_repair_symbols(&self) -> bool {
98        self.symbols.keys().any(|&isi| isi >= self.k_source)
99    }
100
101    /// Sorted unique ISIs of all collected symbols.
102    #[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    /// Add a symbol. Returns `true` if the symbol was new (accepted).
111    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    /// Attempt to decode the collected symbols into changeset bytes.
131    ///
132    /// For source symbols (ISI < k_source), this reconstructs the padded
133    /// changeset by placing each symbol at offset `ISI * symbol_size`.
134    /// Repair symbols would require RaptorQ decoding in production;
135    /// this implementation handles the source-symbol-only case.
136    ///
137    /// Returns `None` if insufficient symbols or decode fails.
138    fn try_decode(&self) -> Option<Vec<u8>> {
139        if !self.ready_to_decode() {
140            return None;
141        }
142
143        // Count source symbols available.
144        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            // All source symbols available — reconstruct directly.
151            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            // Need repair symbols + RaptorQ decoder (production path via asupersync).
163            // For now, return None to stay in COLLECTING.
164            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/// A decoded and validated page ready for application.
177#[derive(Debug, Clone, PartialEq, Eq)]
178pub struct DecodedPage {
179    /// Page number in the database.
180    pub page_number: u32,
181    /// Validated page data.
182    pub page_data: Vec<u8>,
183}
184
185/// Result of a successful decode operation.
186#[derive(Debug)]
187pub struct DecodeResult {
188    /// The changeset identifier that was decoded.
189    pub changeset_id: ChangesetId,
190    /// Decoded and validated pages, sorted by page number.
191    pub pages: Vec<DecodedPage>,
192    /// Number of symbols used for decoding.
193    pub symbols_used: u32,
194    /// Optional decode proof emitted under policy control.
195    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/// Replication receiver state machine.
211#[derive(Debug)]
212pub struct ReplicationReceiver {
213    config: ReceiverConfig,
214    state: ReceiverState,
215    /// Per-changeset decoder states.
216    decoders: HashMap<ChangesetId, DecoderState>,
217    /// Received symbol counts per changeset.
218    received_counts: HashMap<ChangesetId, u32>,
219    /// Total bytes currently buffered across all decoder symbol sets.
220    buffered_symbol_bytes: usize,
221    /// Decoded results waiting for application.
222    pending_results: Vec<DecodeResult>,
223    /// Applied results (for metrics/ACK).
224    applied_count: u64,
225    /// Decode-proof audit entries emitted by this receiver.
226    decode_audit: Vec<DecodeAuditEntry>,
227    /// Monotonic audit sequence.
228    decode_audit_seq: u64,
229}
230
231/// Receiver policy knobs for packet integrity/auth enforcement.
232#[derive(Debug, Clone, Copy, PartialEq, Eq)]
233pub struct DecodeProofEmissionPolicy {
234    /// Emit proofs on decode failure (durability-critical requirement).
235    pub emit_on_decode_failure: bool,
236    /// Emit proofs on successful decode that included repair symbols.
237    pub emit_on_repair_success: bool,
238}
239
240impl DecodeProofEmissionPolicy {
241    /// Default production posture: disabled.
242    #[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    /// Durability-critical posture for replication apply paths.
251    #[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/// Receiver policy knobs for packet integrity/auth enforcement.
267#[derive(Debug, Clone)]
268pub struct ReceiverConfig {
269    /// Optional auth key for validating packet auth tags.
270    pub auth_key: Option<[u8; 32]>,
271    /// Decode proof emission hooks.
272    pub decode_proof_policy: DecodeProofEmissionPolicy,
273    /// Maximum number of concurrent in-flight changeset decoders.
274    pub max_inflight_decoders: usize,
275    /// Maximum total bytes buffered across all decoder symbol maps.
276    pub max_buffered_symbol_bytes: usize,
277}
278
279impl ReceiverConfig {
280    /// Build a receiver config with authenticated transport enabled.
281    #[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    /// Create a new receiver with explicit configuration.
314    #[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    /// Create a new receiver in LISTENING state.
330    #[must_use]
331    pub fn new() -> Self {
332        Self::with_config(ReceiverConfig::default())
333    }
334
335    /// Current state.
336    #[must_use]
337    pub const fn state(&self) -> ReceiverState {
338        self.state
339    }
340
341    /// Number of changesets successfully applied.
342    #[must_use]
343    pub const fn applied_count(&self) -> u64 {
344        self.applied_count
345    }
346
347    /// Number of active decoder sessions.
348    #[must_use]
349    pub fn active_decoders(&self) -> usize {
350        self.decoders.len()
351    }
352
353    /// View decode-proof audit entries emitted so far.
354    #[must_use]
355    pub fn decode_audit_entries(&self) -> &[DecodeAuditEntry] {
356        &self.decode_audit
357    }
358
359    /// Drain decode-proof audit entries.
360    pub fn take_decode_audit_entries(&mut self) -> Vec<DecodeAuditEntry> {
361        std::mem::take(&mut self.decode_audit)
362    }
363
364    /// Process a raw packet from the wire.
365    ///
366    /// # Errors
367    ///
368    /// Returns error if:
369    /// - Packet is malformed (too short, symbol_size = 0)
370    /// - V1 rule violated (SBN != 0)
371    /// - K_source out of range
372    /// - K_source or symbol_size mismatch for existing decoder
373    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    /// Process a parsed packet.
391    ///
392    /// # Errors
393    ///
394    /// See `process_packet`.
395    #[allow(clippy::too_many_lines)]
396    pub fn process_parsed_packet(&mut self, packet: &ReplicationPacket) -> Result<PacketResult> {
397        // V1 rule: reject multi-block packets.
398        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        // Validate K_source range.
411        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        // Compute symbol_size from packet header and validate payload consistency.
425        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        // Transition LISTENING → COLLECTING on first packet.
443        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        // Get or create decoder state.
452        if let Some(decoder) = self.decoders.get(&changeset_id) {
453            // Validate consistency with existing decoder.
454            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            // Create new decoder state.
503            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        // Enforce global buffered-symbol bound before accepting a new symbol.
532        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        // Add symbol to decoder (with ISI deduplication) and capture decode context.
559        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                // Decode succeeded: truncate to total_len and parse pages.
637                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                        // Clean up decoder for this changeset.
651                        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                        // Clean up failed decoder.
661                        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            // Decode failed (need more symbols).
687            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    /// Parse and validate decoded changeset bytes.
701    #[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        // Parse header.
717        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        // Truncate to total_len.
723        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        // Parse page entries.
754        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            })?; // page_number + xxh3 + data
772        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            // Validate page xxh3.
820            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        // Pages should already be sorted (sender sorts them).
843        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    /// Apply pending decoded results. Returns applied page counts.
883    ///
884    /// In production, this writes pages to the local database. Here we
885    /// validate and return the results for the caller to apply.
886    ///
887    /// # Errors
888    ///
889    /// Returns error if not in APPLYING state.
890    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        // Transition to COMPLETE.
910        self.state = ReceiverState::Complete;
911        Ok(results)
912    }
913
914    /// Transition from COMPLETE back to LISTENING for the next changeset.
915    ///
916    /// # Errors
917    ///
918    /// Returns error if not in COMPLETE state.
919    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    /// Force reset to LISTENING from any state (e.g., on error recovery).
932    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/// Result of processing a single packet.
957#[derive(Debug, Clone, Copy, PartialEq, Eq)]
958pub enum PacketResult {
959    /// Symbol accepted, need more for decode.
960    Accepted,
961    /// Integrity/auth invalid; packet ignored as erasure.
962    Erasure,
963    /// Duplicate ISI, silently ignored.
964    Duplicate,
965    /// Enough symbols collected, decode succeeded and ready to apply.
966    DecodeReady,
967    /// Had enough symbols but decode failed, need more.
968    NeedMore,
969}
970
971// ---------------------------------------------------------------------------
972// Changeset parsing utility (used by tests and receiver)
973// ---------------------------------------------------------------------------
974
975/// Parse changeset bytes into page entries (for validation/testing).
976///
977/// # Errors
978///
979/// Returns error if the changeset is malformed.
980pub 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    /// Helper: generate sender packets for a set of pages.
1119    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    // -----------------------------------------------------------------------
1287    // State transition tests
1288    // -----------------------------------------------------------------------
1289
1290    #[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        // Should have created exactly one decoder.
1318        // Note: if decode triggers, the decoder may be cleaned up,
1319        // so just check that processing succeeded.
1320        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        // New changeset would exceed budget and should be rejected/cleaned up.
1378        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        // Verify seed = xxh3_64(changeset_id_bytes) matches sender.
1393        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, // V1 violation
1409            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        // K_source = 0 → rejected.
1426        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        // K_source = K_MAX + 1 → rejected.
1440        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        // ESI only has 24 bits, K_source > K_MAX might not fit in packet format
1448        // but we test the validation path directly.
1449        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        // K_source = K_MAX → accepted.
1456        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        // Symbol size should be inferred as 1024.
1485        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        // Zero-length symbol data → rejected.
1495        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        // Same changeset_id, different K_source.
1514        let p2 = make_packet(id, 0, 1, 200, vec![0x42; 512]); // mismatch
1515        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        // Same changeset_id, different symbol_size.
1532        let p2 = make_packet(id, 0, 1, 100, vec![0x42; 1024]); // different size
1533        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        // Same ISI again → duplicate.
1554        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        // Count should still be 1.
1562        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        // Use the sender to generate proper packets, then feed to receiver.
1634        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        // Two repair-only symbols at K=2: ready_to_decode => true, but decode fails.
1664        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        // Build K=2 source symbols from encoded bytes.
1702        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        // Interleave source+repair so K is reached with at least one repair symbol present.
1713        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        // Apply and check that pages are correctly truncated.
1759        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        // Tamper with a page byte in the changeset (after header + page_number + xxh3).
1784        let mut tampered = changeset_bytes.clone();
1785        let tamper_offset = CHANGESET_HEADER_SIZE + 4 + 8 + 10; // into page data
1786        tampered[tamper_offset] ^= 0xFF;
1787
1788        // Now create a "decoded" changeset and try to parse it.
1789        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    // -----------------------------------------------------------------------
1906    // Property tests
1907    // -----------------------------------------------------------------------
1908
1909    #[test]
1910    fn prop_any_k_symbols_decode() {
1911        // With only source symbols and k_source = actual source count,
1912        // providing all k source symbols always decodes.
1913        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        // Use a large K_source so we can feed duplicates before decode triggers.
1936        let mut receiver = ReplicationReceiver::new();
1937        let id = ChangesetId::from_bytes([0x77; 16]);
1938
1939        // Feed the same ISI multiple times within a single decoder session.
1940        let p1 = make_packet(id, 0, 0, 100, vec![0x42; 512]); // large enough that one symbol won't trigger decode
1941
1942        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        // Count should still be 1.
1959        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    // -----------------------------------------------------------------------
1964    // E2E tests
1965    // -----------------------------------------------------------------------
1966
1967    #[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        // Sender encodes pages. Receiver collects and decodes. Byte-identical.
1978        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        // Complete the cycle.
2022        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        // Two changesets streaming simultaneously.
2033        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        // Interleave packets from two different changesets.
2039        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                // Apply immediately and reset if needed.
2055                if receiver.state() == ReceiverState::Applying {
2056                    let _ = receiver.apply_pending();
2057                    // If more decoders remain, go back to collecting.
2058                    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        // Full end-to-end compliance test.
2074        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        // Encode via sender.
2079        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        // Collect all packets.
2087        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        // Feed to receiver.
2093        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        // Verify decode happened.
2108        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        // Apply.
2116        let results = receiver.apply_pending().expect("apply");
2117        assert_eq!(receiver.state(), ReceiverState::Complete);
2118        assert_eq!(results.len(), 1);
2119
2120        // Verify byte-identical pages.
2121        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        // Reset and verify.
2129        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    // -----------------------------------------------------------------------
2313    // Compliance gate tests
2314    // -----------------------------------------------------------------------
2315
2316    #[test]
2317    fn test_bd_1hi_14_unit_compliance_gate() {
2318        // Verify all required types and functions exist.
2319        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        // Verify REPLICATION_HEADER_SIZE is correct.
2337        assert_eq!(REPLICATION_HEADER_SIZE, 72);
2338    }
2339
2340    #[test]
2341    fn prop_bd_1hi_14_structure_compliance() {
2342        // Full state machine cycle.
2343        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        // Should have transitioned through the state machine.
2353        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}