Skip to main content

fsqlite_core/
snapshot_shipping.rs

1//! §3.4.3 Fountain-Coded Snapshot Shipping (bd-1hi.15).
2//!
3//! Implements snapshot transfer for initializing new replicas using
4//! fountain coding. The entire database is partitioned into source blocks
5//! and streamed as rateless-coded symbols over UDP.
6//!
7//! Key advantages:
8//! - No handshake or acknowledgment needed
9//! - Receiver can start receiving from any point in the stream
10//! - Inherently resumable with zero protocol overhead
11//! - Natural multicast: initialize many replicas simultaneously
12//! - Progressive receive: partial queries after first block decoded
13
14use std::collections::{HashMap, HashSet};
15
16use fsqlite_error::{FrankenError, Result};
17use tracing::{debug, error, info, warn};
18
19use crate::replication_sender::{
20    CHANGESET_HEADER_SIZE, ChangesetId, PageEntry, ReplicationPacket, ReplicationPacketV2Header,
21    SenderConfig, compute_changeset_id, derive_seed_from_changeset_id, encode_changeset,
22};
23use crate::source_block_partition::{K_MAX, SourceBlock, partition_source_blocks};
24
25const BEAD_ID: &str = "bd-1hi.15";
26
27// ---------------------------------------------------------------------------
28// Resume State (persistent across connection losses)
29// ---------------------------------------------------------------------------
30
31/// Per-block resume state: tracks which ISIs have been received.
32#[derive(Debug, Clone)]
33pub struct BlockResumeState {
34    /// Source block index (SBN).
35    pub block_id: u32,
36    /// Number of unique symbols received.
37    pub num_received: u32,
38    /// Set of received ISIs (for O(1) dedup).
39    pub received_isis: HashSet<u32>,
40    /// Whether this block has been fully decoded.
41    pub decoded: bool,
42}
43
44impl BlockResumeState {
45    /// Create a new empty resume state for a block.
46    #[must_use]
47    fn new(block_id: u32) -> Self {
48        Self {
49            block_id,
50            num_received: 0,
51            received_isis: HashSet::new(),
52            decoded: false,
53        }
54    }
55
56    /// Record a received ISI. Returns true if new (accepted).
57    fn record_isi(&mut self, isi: u32) -> bool {
58        if self.received_isis.insert(isi) {
59            self.num_received += 1;
60            true
61        } else {
62            false
63        }
64    }
65
66    /// Serialize to a compact binary format for persistence.
67    ///
68    /// Format: `block_id(4 LE) | num_received(4 LE) | decoded(1) | n_isis(4 LE) | isis(4 LE each)`
69    #[must_use]
70    pub fn to_bytes(&self) -> Vec<u8> {
71        let n = self.received_isis.len();
72        let mut buf = Vec::with_capacity(13 + n * 4);
73        buf.extend_from_slice(&self.block_id.to_le_bytes());
74        buf.extend_from_slice(&self.num_received.to_le_bytes());
75        buf.push(u8::from(self.decoded));
76        let n_u32 = u32::try_from(n).unwrap_or(u32::MAX);
77        buf.extend_from_slice(&n_u32.to_le_bytes());
78        let mut sorted_isis: Vec<u32> = self.received_isis.iter().copied().collect();
79        sorted_isis.sort_unstable();
80        for isi in sorted_isis {
81            buf.extend_from_slice(&isi.to_le_bytes());
82        }
83        buf
84    }
85
86    /// Deserialize from bytes.
87    ///
88    /// # Errors
89    ///
90    /// Returns error if buffer is too short or malformed.
91    pub fn from_bytes(buf: &[u8]) -> Result<(Self, usize)> {
92        if buf.len() < 13 {
93            return Err(FrankenError::DatabaseCorrupt {
94                detail: format!("BlockResumeState too short: {} < 13", buf.len()),
95            });
96        }
97        let block_id = u32::from_le_bytes(buf[0..4].try_into().expect("4 bytes"));
98        let num_received = u32::from_le_bytes(buf[4..8].try_into().expect("4 bytes"));
99        let decoded = buf[8] != 0;
100        let n_isis = u32::from_le_bytes(buf[9..13].try_into().expect("4 bytes"));
101        let n = n_isis as usize;
102        let expected = n
103            .checked_mul(4)
104            .and_then(|v| v.checked_add(13))
105            .ok_or_else(|| FrankenError::DatabaseCorrupt {
106                detail: format!("BlockResumeState n_isis ({n_isis}) causes size overflow"),
107            })?;
108        if buf.len() < expected {
109            return Err(FrankenError::DatabaseCorrupt {
110                detail: format!("BlockResumeState truncated: {} < {expected}", buf.len()),
111            });
112        }
113        let mut received_isis = HashSet::with_capacity(n);
114        for i in 0..n {
115            let offset = 13 + i * 4;
116            let isi = u32::from_le_bytes(buf[offset..offset + 4].try_into().expect("4 bytes"));
117            received_isis.insert(isi);
118        }
119        Ok((
120            Self {
121                block_id,
122                num_received,
123                received_isis,
124                decoded,
125            },
126            expected,
127        ))
128    }
129}
130
131/// Full resume state for a snapshot transfer.
132#[derive(Debug, Clone)]
133pub struct ResumeState {
134    /// Per-block resume states.
135    pub blocks: Vec<BlockResumeState>,
136    /// Total number of source blocks expected.
137    pub total_blocks: u32,
138}
139
140impl ResumeState {
141    /// Create a new resume state for a snapshot with `total_blocks` blocks.
142    #[must_use]
143    pub fn new(total_blocks: u32) -> Self {
144        let blocks = (0..total_blocks).map(BlockResumeState::new).collect();
145        Self {
146            blocks,
147            total_blocks,
148        }
149    }
150
151    /// Number of blocks fully decoded.
152    #[must_use]
153    pub fn decoded_count(&self) -> u32 {
154        u32::try_from(self.blocks.iter().filter(|b| b.decoded).count()).unwrap_or(u32::MAX)
155    }
156
157    /// Whether all blocks are decoded.
158    #[must_use]
159    pub fn all_decoded(&self) -> bool {
160        self.blocks.iter().all(|b| b.decoded)
161    }
162
163    /// Serialize to bytes.
164    #[must_use]
165    pub fn to_bytes(&self) -> Vec<u8> {
166        let mut buf = Vec::new();
167        buf.extend_from_slice(&self.total_blocks.to_le_bytes());
168        for block in &self.blocks {
169            buf.extend_from_slice(&block.to_bytes());
170        }
171        buf
172    }
173
174    /// Deserialize from bytes.
175    ///
176    /// # Errors
177    ///
178    /// Returns error if buffer is malformed.
179    pub fn from_bytes(buf: &[u8]) -> Result<Self> {
180        if buf.len() < 4 {
181            return Err(FrankenError::DatabaseCorrupt {
182                detail: format!("ResumeState too short: {} < 4", buf.len()),
183            });
184        }
185        let total_blocks = u32::from_le_bytes(buf[0..4].try_into().expect("4 bytes"));
186        let mut blocks = Vec::with_capacity(total_blocks as usize);
187        let mut offset = 4;
188        for _ in 0..total_blocks {
189            let (block, consumed) = BlockResumeState::from_bytes(&buf[offset..])?;
190            blocks.push(block);
191            offset += consumed;
192        }
193        Ok(Self {
194            blocks,
195            total_blocks,
196        })
197    }
198}
199
200// ---------------------------------------------------------------------------
201// Snapshot Sender
202// ---------------------------------------------------------------------------
203
204/// Snapshot sender: partitions a database into source blocks and streams symbols.
205#[derive(Debug)]
206pub struct SnapshotSender {
207    /// Source blocks from the partition algorithm.
208    pub source_blocks: Vec<SourceBlock>,
209    /// Page size of the database.
210    pub page_size: u32,
211    /// Current block being streamed.
212    current_block: usize,
213    /// Current ISI within the current block.
214    current_isi: u32,
215    /// Per-block changeset IDs (computed during prepare).
216    block_changeset_ids: Vec<ChangesetId>,
217    /// Per-block K_source values.
218    block_k_sources: Vec<u32>,
219    /// Per-block changeset bytes.
220    block_changesets: Vec<Vec<u8>>,
221    /// Sender config.
222    config: SenderConfig,
223    /// Whether we're done.
224    done: bool,
225}
226
227impl SnapshotSender {
228    /// Prepare a snapshot sender for the given database pages.
229    ///
230    /// `all_pages` must be sorted by page number and cover the entire database.
231    ///
232    /// # Errors
233    ///
234    /// Returns error if partitioning fails or pages are empty.
235    #[allow(clippy::too_many_lines)]
236    pub fn prepare(
237        page_size: u32,
238        all_pages: &mut [PageEntry],
239        config: SenderConfig,
240    ) -> Result<Self> {
241        if all_pages.is_empty() {
242            return Err(FrankenError::OutOfRange {
243                what: "pages".to_owned(),
244                value: "0".to_owned(),
245            });
246        }
247
248        let total_pages = u32::try_from(all_pages.len()).map_err(|_| FrankenError::OutOfRange {
249            what: "total_pages".to_owned(),
250            value: all_pages.len().to_string(),
251        })?;
252
253        let source_blocks = partition_source_blocks(total_pages)?;
254        info!(
255            bead_id = BEAD_ID,
256            total_pages,
257            n_blocks = source_blocks.len(),
258            page_size,
259            "snapshot partitioned into source blocks"
260        );
261
262        // Sort all pages by page_number.
263        all_pages.sort_by_key(|p| p.page_number);
264
265        // Build per-block changesets.
266        let mut block_changeset_ids = Vec::with_capacity(source_blocks.len());
267        let mut block_k_sources = Vec::with_capacity(source_blocks.len());
268        let mut block_changesets = Vec::with_capacity(source_blocks.len());
269
270        let mut page_idx = 0_usize;
271        for block in &source_blocks {
272            let end = page_idx + block.num_pages as usize;
273            if end > all_pages.len() {
274                return Err(FrankenError::Internal(format!(
275                    "block {} requires pages up to index {end}, but only {} available",
276                    block.index,
277                    all_pages.len()
278                )));
279            }
280            let block_pages = &mut all_pages[page_idx..end];
281            let changeset_bytes = encode_changeset(page_size, block_pages)?;
282            let changeset_id = compute_changeset_id(&changeset_bytes);
283
284            // Compute K_source from changeset + symbol_size.
285            let t = u64::from(config.symbol_size);
286            let f = changeset_bytes.len() as u64;
287            let k_source = u32::try_from(f.div_ceil(t)).map_err(|_| FrankenError::OutOfRange {
288                what: "k_source".to_owned(),
289                value: f.div_ceil(t).to_string(),
290            })?;
291
292            debug!(
293                bead_id = BEAD_ID,
294                block_index = block.index,
295                num_pages = block.num_pages,
296                changeset_len = changeset_bytes.len(),
297                k_source,
298                "prepared block changeset"
299            );
300
301            block_changeset_ids.push(changeset_id);
302            block_k_sources.push(k_source);
303            block_changesets.push(changeset_bytes);
304            page_idx = end;
305        }
306
307        Ok(Self {
308            source_blocks,
309            page_size,
310            current_block: 0,
311            current_isi: 0,
312            block_changeset_ids,
313            block_k_sources,
314            block_changesets,
315            config,
316            done: false,
317        })
318    }
319
320    /// Generate the next snapshot packet.
321    ///
322    /// Returns `None` when the current streaming pass is complete.
323    /// Caller can restart from block 0 for continuous streaming.
324    pub fn next_packet(&mut self) -> Option<ReplicationPacket> {
325        if self.done || self.current_block >= self.source_blocks.len() {
326            self.done = true;
327            return None;
328        }
329
330        let k_source = self.block_k_sources[self.current_block];
331        let max_isi = k_source.saturating_mul(self.config.max_isi_multiplier);
332
333        if self.current_isi >= max_isi {
334            self.current_block += 1;
335            self.current_isi = 0;
336            if self.current_block >= self.source_blocks.len() {
337                self.done = true;
338                return None;
339            }
340        }
341
342        let changeset = &self.block_changesets[self.current_block];
343        let changeset_id = self.block_changeset_ids[self.current_block];
344        let k_source = self.block_k_sources[self.current_block];
345        let isi = self.current_isi;
346        let t = usize::from(self.config.symbol_size);
347
348        // Extract or generate symbol data.
349        let symbol_data = if u64::from(isi) < u64::from(k_source) {
350            let start = isi as usize * t;
351            let end = (start + t).min(changeset.len());
352            let mut data = vec![0_u8; t];
353            let available = end.saturating_sub(start);
354            if available > 0 {
355                data[..available].copy_from_slice(&changeset[start..end]);
356            }
357            data
358        } else {
359            let seed = derive_seed_from_changeset_id(&changeset_id);
360            #[cfg(not(target_arch = "wasm32"))]
361            {
362                // Repair symbol: use RaptorQ SystematicEncoder from asupersync.
363                // NOTE: Encoder is rebuilt per-symbol (expensive for large K).
364                // See replication_sender.rs for the same pattern and caching note.
365                use asupersync::raptorq::systematic::SystematicEncoder;
366
367                let source_symbols: Vec<Vec<u8>> = (0..k_source as usize)
368                    .map(|i| {
369                        let start = i * t;
370                        let end = (start + t).min(changeset.len());
371                        let mut sym = vec![0_u8; t];
372                        let available = end.saturating_sub(start);
373                        if available > 0 {
374                            sym[..available].copy_from_slice(&changeset[start..end]);
375                        }
376                        sym
377                    })
378                    .collect();
379
380                match SystematicEncoder::new(&source_symbols, t, seed) {
381                    Some(encoder) => {
382                        let repair_esi = isi - k_source;
383                        encoder.repair_symbol(repair_esi)
384                    }
385                    None => {
386                        // Fallback for degenerate parameters.
387                        crate::replication_sender::generate_deterministic_placeholder(seed, isi, t)
388                    }
389                }
390            }
391            #[cfg(target_arch = "wasm32")]
392            {
393                crate::replication_sender::generate_deterministic_placeholder(seed, isi, t)
394            }
395        };
396
397        let seed = derive_seed_from_changeset_id(&changeset_id);
398        let r_repair = max_isi.saturating_sub(k_source);
399        let packet = ReplicationPacket::new_v2(
400            ReplicationPacketV2Header {
401                changeset_id,
402                sbn: 0,
403                esi: isi,
404                k_source,
405                r_repair,
406                symbol_size_t: self.config.symbol_size,
407                seed,
408            },
409            symbol_data,
410        );
411
412        self.current_isi += 1;
413        Some(packet)
414    }
415
416    /// Number of source blocks.
417    #[must_use]
418    pub fn num_blocks(&self) -> usize {
419        self.source_blocks.len()
420    }
421
422    /// Total source symbols across all blocks.
423    #[must_use]
424    pub fn total_source_symbols(&self) -> u64 {
425        self.block_k_sources.iter().map(|&k| u64::from(k)).sum()
426    }
427
428    /// Reset to re-stream from the beginning (for continuous multicast).
429    pub fn restart(&mut self) {
430        self.current_block = 0;
431        self.current_isi = 0;
432        self.done = false;
433        debug!(bead_id = BEAD_ID, "snapshot sender restarted for next pass");
434    }
435}
436
437// ---------------------------------------------------------------------------
438// Snapshot Receiver
439// ---------------------------------------------------------------------------
440
441/// Snapshot receiver state.
442#[derive(Debug, Clone, Copy, PartialEq, Eq)]
443pub enum SnapshotReceiverState {
444    /// Waiting for first packet.
445    Waiting,
446    /// Actively collecting symbols.
447    Receiving,
448    /// All blocks decoded, snapshot complete.
449    Complete,
450}
451
452/// A decoded source block's pages.
453#[derive(Debug, Clone)]
454pub struct DecodedBlock {
455    /// Block index.
456    pub block_index: u32,
457    /// Decoded pages sorted by page number.
458    pub pages: Vec<DecodedBlockPage>,
459}
460
461/// A single page from a decoded block.
462#[derive(Debug, Clone, PartialEq, Eq)]
463pub struct DecodedBlockPage {
464    /// Page number.
465    pub page_number: u32,
466    /// Page data.
467    pub page_data: Vec<u8>,
468}
469
470/// Per-block decoder used by the snapshot receiver.
471#[derive(Debug)]
472struct BlockDecoder {
473    /// The changeset_id for this block (determined from first packet).
474    changeset_id: Option<ChangesetId>,
475    /// K_source for this block.
476    k_source: u32,
477    /// Symbol size.
478    symbol_size: u32,
479    /// Seed for RaptorQ.
480    seed: u64,
481    /// Symbols collected by ISI.
482    symbols: HashMap<u32, Vec<u8>>,
483    /// ISI dedup set.
484    received_isis: HashSet<u32>,
485    /// Whether decoded.
486    decoded: bool,
487}
488
489impl BlockDecoder {
490    fn new() -> Self {
491        Self {
492            changeset_id: None,
493            k_source: 0,
494            symbol_size: 0,
495            seed: 0,
496            symbols: HashMap::new(),
497            received_isis: HashSet::new(),
498            decoded: false,
499        }
500    }
501
502    fn initialize(&mut self, changeset_id: ChangesetId, k_source: u32, symbol_size: u32) {
503        self.changeset_id = Some(changeset_id);
504        self.k_source = k_source;
505        self.symbol_size = symbol_size;
506        self.seed = derive_seed_from_changeset_id(&changeset_id);
507    }
508
509    fn add_symbol(&mut self, isi: u32, data: Vec<u8>) -> bool {
510        if self.received_isis.insert(isi) {
511            self.symbols.insert(isi, data);
512            true
513        } else {
514            false
515        }
516    }
517
518    fn received_count(&self) -> u32 {
519        u32::try_from(self.received_isis.len()).unwrap_or(u32::MAX)
520    }
521
522    fn ready_to_decode(&self) -> bool {
523        self.received_count() >= self.k_source && self.k_source > 0
524    }
525
526    fn try_decode(&self) -> Option<Vec<u8>> {
527        if !self.ready_to_decode() {
528            return None;
529        }
530        let source_count = self
531            .symbols
532            .keys()
533            .filter(|&&isi| isi < self.k_source)
534            .count();
535        let k = self.k_source as usize;
536        let t = self.symbol_size as usize;
537        if source_count >= k {
538            let padded_len = k * t;
539            let mut padded = vec![0_u8; padded_len];
540            for isi in 0..self.k_source {
541                if let Some(data) = self.symbols.get(&isi) {
542                    let start = isi as usize * t;
543                    let copy_len = data.len().min(t);
544                    padded[start..start + copy_len].copy_from_slice(&data[..copy_len]);
545                }
546            }
547            Some(padded)
548        } else {
549            warn!(
550                bead_id = BEAD_ID,
551                source_count,
552                k_source = self.k_source,
553                "snapshot block decode needs repair symbols (production RaptorQ)"
554            );
555            None
556        }
557    }
558}
559
560/// Snapshot receiver: collects symbols per source block, decodes progressively.
561#[derive(Debug)]
562pub struct SnapshotReceiver {
563    state: SnapshotReceiverState,
564    /// Per-changeset_id → block index mapping.
565    changeset_to_block: HashMap<ChangesetId, usize>,
566    /// Per-block decoders.
567    block_decoders: Vec<BlockDecoder>,
568    /// Number of blocks expected (set after first packet or from resume state).
569    num_blocks: usize,
570    /// Decoded blocks ready for application.
571    decoded_blocks: Vec<DecodedBlock>,
572    /// Resume state.
573    resume: ResumeState,
574    /// Page size.
575    page_size: u32,
576}
577
578impl SnapshotReceiver {
579    /// Create a new snapshot receiver.
580    ///
581    /// `num_blocks` is the expected number of source blocks (from partitioning).
582    /// `page_size` is the database page size.
583    #[must_use]
584    pub fn new(num_blocks: usize, page_size: u32) -> Self {
585        let block_decoders = (0..num_blocks).map(|_| BlockDecoder::new()).collect();
586        Self {
587            state: SnapshotReceiverState::Waiting,
588            changeset_to_block: HashMap::new(),
589            block_decoders,
590            num_blocks,
591            decoded_blocks: Vec::new(),
592            resume: ResumeState::new(u32::try_from(num_blocks).unwrap_or(u32::MAX)),
593            page_size,
594        }
595    }
596
597    /// Create from a resume state (after crash/reconnect).
598    #[must_use]
599    pub fn from_resume(resume: ResumeState, page_size: u32) -> Self {
600        let num_blocks = resume.total_blocks as usize;
601        let block_decoders = (0..num_blocks).map(|_| BlockDecoder::new()).collect();
602        Self {
603            state: if resume.all_decoded() {
604                SnapshotReceiverState::Complete
605            } else {
606                SnapshotReceiverState::Waiting
607            },
608            changeset_to_block: HashMap::new(),
609            block_decoders,
610            num_blocks,
611            decoded_blocks: Vec::new(),
612            resume,
613            page_size,
614        }
615    }
616
617    /// Current state.
618    #[must_use]
619    pub const fn state(&self) -> SnapshotReceiverState {
620        self.state
621    }
622
623    /// Number of blocks decoded so far.
624    #[must_use]
625    pub fn blocks_decoded(&self) -> usize {
626        self.decoded_blocks.len()
627    }
628
629    /// Get the resume state for persistence.
630    #[must_use]
631    pub fn resume_state(&self) -> &ResumeState {
632        &self.resume
633    }
634
635    /// Take decoded blocks (for application to local database).
636    pub fn take_decoded_blocks(&mut self) -> Vec<DecodedBlock> {
637        std::mem::take(&mut self.decoded_blocks)
638    }
639
640    /// Process a snapshot packet.
641    ///
642    /// The receiver maps packets to blocks by changeset_id. The first packet
643    /// for a new changeset_id establishes the mapping to the next unmapped block.
644    ///
645    /// # Errors
646    ///
647    /// Returns error if the packet is malformed or validation fails.
648    #[allow(clippy::too_many_lines)]
649    pub fn process_packet(&mut self, packet: &ReplicationPacket) -> Result<SnapshotPacketResult> {
650        if self.state == SnapshotReceiverState::Complete {
651            return Ok(SnapshotPacketResult::AlreadyComplete);
652        }
653
654        // V1 rule.
655        if packet.sbn != 0 {
656            return Err(FrankenError::Internal(format!(
657                "V1: SBN must be 0, got {}",
658                packet.sbn
659            )));
660        }
661        if packet.k_source == 0 || packet.k_source > K_MAX {
662            return Err(FrankenError::OutOfRange {
663                what: "k_source".to_owned(),
664                value: packet.k_source.to_string(),
665            });
666        }
667        let symbol_size =
668            u32::try_from(packet.symbol_data.len()).map_err(|_| FrankenError::OutOfRange {
669                what: "symbol_size".to_owned(),
670                value: packet.symbol_data.len().to_string(),
671            })?;
672        if symbol_size == 0 {
673            return Err(FrankenError::OutOfRange {
674                what: "symbol_size".to_owned(),
675                value: "0".to_owned(),
676            });
677        }
678
679        if self.state == SnapshotReceiverState::Waiting {
680            self.state = SnapshotReceiverState::Receiving;
681            info!(bead_id = BEAD_ID, "snapshot receiving started");
682        }
683
684        let changeset_id = packet.changeset_id;
685
686        // Map changeset_id to block index.
687        let block_idx = if let Some(&idx) = self.changeset_to_block.get(&changeset_id) {
688            idx
689        } else {
690            // Find the next unmapped, undecoded block.
691            let next_idx = self
692                .block_decoders
693                .iter()
694                .position(|d| d.changeset_id.is_none() && !d.decoded);
695            if let Some(idx) = next_idx {
696                self.changeset_to_block.insert(changeset_id, idx);
697                self.block_decoders[idx].initialize(changeset_id, packet.k_source, symbol_size);
698                debug!(
699                    bead_id = BEAD_ID,
700                    block_index = idx,
701                    k_source = packet.k_source,
702                    "mapped new changeset to block"
703                );
704                idx
705            } else {
706                warn!(
707                    bead_id = BEAD_ID,
708                    "no available block slot for new changeset_id"
709                );
710                return Ok(SnapshotPacketResult::Rejected);
711            }
712        };
713
714        if block_idx >= self.block_decoders.len() {
715            return Ok(SnapshotPacketResult::Rejected);
716        }
717
718        let decoder = &mut self.block_decoders[block_idx];
719        if decoder.decoded {
720            return Ok(SnapshotPacketResult::BlockAlreadyDecoded);
721        }
722
723        // Validate consistency.
724        if decoder.k_source != packet.k_source {
725            return Err(FrankenError::DatabaseCorrupt {
726                detail: format!(
727                    "k_source mismatch for block {block_idx}: {} vs {}",
728                    decoder.k_source, packet.k_source
729                ),
730            });
731        }
732        if decoder.symbol_size != symbol_size {
733            return Err(FrankenError::DatabaseCorrupt {
734                detail: format!(
735                    "symbol_size mismatch for block {block_idx}: {} vs {symbol_size}",
736                    decoder.symbol_size
737                ),
738            });
739        }
740
741        // Add symbol.
742        let accepted = decoder.add_symbol(packet.esi, packet.symbol_data.clone());
743        if !accepted {
744            return Ok(SnapshotPacketResult::Duplicate);
745        }
746
747        // Update resume state.
748        if block_idx < self.resume.blocks.len() {
749            self.resume.blocks[block_idx].record_isi(packet.esi);
750        }
751
752        // Check if ready to decode this block.
753        if decoder.ready_to_decode()
754            && !decoder.decoded
755            && let Some(padded) = decoder.try_decode()
756        {
757            match parse_decoded_snapshot_block(&padded, self.page_size) {
758                Ok(pages) => {
759                    let block_id = u32::try_from(block_idx).unwrap_or(u32::MAX);
760                    decoder.decoded = true;
761                    if block_idx < self.resume.blocks.len() {
762                        self.resume.blocks[block_idx].decoded = true;
763                    }
764                    let n_pages = pages.len();
765                    self.decoded_blocks.push(DecodedBlock {
766                        block_index: block_id,
767                        pages,
768                    });
769                    info!(
770                        bead_id = BEAD_ID,
771                        block_index = block_idx,
772                        n_pages,
773                        decoded_so_far = self.decoded_blocks.len(),
774                        total_blocks = self.num_blocks,
775                        "source block decoded (progressive)"
776                    );
777
778                    // Check if all blocks are done.
779                    if self.block_decoders.iter().all(|d| d.decoded) {
780                        self.state = SnapshotReceiverState::Complete;
781                        info!(
782                            bead_id = BEAD_ID,
783                            total_blocks = self.num_blocks,
784                            "snapshot fully received"
785                        );
786                    }
787                    return Ok(SnapshotPacketResult::BlockDecoded(block_id));
788                }
789                Err(e) => {
790                    error!(
791                        bead_id = BEAD_ID,
792                        block_index = block_idx,
793                        error = %e,
794                        "snapshot block validation failed"
795                    );
796                    return Err(e);
797                }
798            }
799        }
800
801        Ok(SnapshotPacketResult::Accepted)
802    }
803}
804
805/// Result of processing a snapshot packet.
806#[derive(Debug, Clone, Copy, PartialEq, Eq)]
807pub enum SnapshotPacketResult {
808    /// Symbol accepted, need more.
809    Accepted,
810    /// Duplicate ISI, ignored.
811    Duplicate,
812    /// A source block was fully decoded (progressive).
813    BlockDecoded(u32),
814    /// This block was already decoded.
815    BlockAlreadyDecoded,
816    /// Packet rejected (no available block slot or already complete).
817    Rejected,
818    /// Snapshot already complete.
819    AlreadyComplete,
820}
821
822// ---------------------------------------------------------------------------
823// Helpers
824// ---------------------------------------------------------------------------
825
826/// Parse decoded snapshot block bytes into pages with xxh3 validation.
827fn parse_decoded_snapshot_block(
828    padded_bytes: &[u8],
829    _page_size: u32,
830) -> Result<Vec<DecodedBlockPage>> {
831    use crate::replication_sender::ChangesetHeader;
832
833    if padded_bytes.len() < CHANGESET_HEADER_SIZE {
834        return Err(FrankenError::DatabaseCorrupt {
835            detail: format!(
836                "decoded block too short for header: {} < {CHANGESET_HEADER_SIZE}",
837                padded_bytes.len()
838            ),
839        });
840    }
841
842    let header_bytes: [u8; CHANGESET_HEADER_SIZE] = padded_bytes[..CHANGESET_HEADER_SIZE]
843        .try_into()
844        .expect("checked length");
845    let header = ChangesetHeader::from_bytes(&header_bytes)?;
846
847    let total_len = usize::try_from(header.total_len).map_err(|_| FrankenError::OutOfRange {
848        what: "total_len".to_owned(),
849        value: header.total_len.to_string(),
850    })?;
851    if total_len > padded_bytes.len() {
852        return Err(FrankenError::DatabaseCorrupt {
853            detail: format!(
854                "total_len ({total_len}) exceeds decoded bytes ({})",
855                padded_bytes.len()
856            ),
857        });
858    }
859    let changeset_bytes = &padded_bytes[..total_len];
860
861    let entry_size = 4_usize + 8 + header.page_size as usize;
862    let data_bytes = &changeset_bytes[CHANGESET_HEADER_SIZE..];
863
864    let required_data_len = (header.n_pages as usize)
865        .checked_mul(entry_size)
866        .ok_or_else(|| FrankenError::DatabaseCorrupt {
867            detail: "n_pages causes size overflow".to_owned(),
868        })?;
869
870    if data_bytes.len() < required_data_len {
871        return Err(FrankenError::DatabaseCorrupt {
872            detail: format!(
873                "changeset truncated: expected {} data bytes, got {}",
874                required_data_len,
875                data_bytes.len()
876            ),
877        });
878    }
879
880    let mut pages = Vec::with_capacity(header.n_pages as usize);
881    for i in 0..header.n_pages as usize {
882        let offset = i * entry_size;
883        let page_number =
884            u32::from_le_bytes(data_bytes[offset..offset + 4].try_into().expect("4 bytes"));
885        let page_xxh3 = u64::from_le_bytes(
886            data_bytes[offset + 4..offset + 12]
887                .try_into()
888                .expect("8 bytes"),
889        );
890        let page_data = data_bytes[offset + 12..offset + 12 + header.page_size as usize].to_vec();
891
892        let computed_xxh3 = xxhash_rust::xxh3::xxh3_64(&page_data);
893        if computed_xxh3 != page_xxh3 {
894            error!(
895                bead_id = BEAD_ID,
896                page_number,
897                expected_xxh3 = page_xxh3,
898                computed_xxh3,
899                "snapshot page xxh3 mismatch"
900            );
901            return Err(FrankenError::DatabaseCorrupt {
902                detail: format!(
903                    "snapshot page {page_number} xxh3 mismatch: {page_xxh3:#x} vs {computed_xxh3:#x}"
904                ),
905            });
906        }
907
908        pages.push(DecodedBlockPage {
909            page_number,
910            page_data,
911        });
912    }
913
914    Ok(pages)
915}
916
917#[cfg(test)]
918mod tests {
919    use super::*;
920    use crate::replication_sender::PageEntry;
921
922    const TEST_BEAD_ID: &str = "bd-1hi.15";
923
924    #[allow(clippy::cast_possible_truncation)]
925    fn make_pages(page_size: u32, page_numbers: &[u32]) -> Vec<PageEntry> {
926        page_numbers
927            .iter()
928            .map(|&pn| {
929                let mut data = vec![0_u8; page_size as usize];
930                for (i, byte) in data.iter_mut().enumerate() {
931                    *byte = ((pn as usize * 251 + i * 31) % 256) as u8;
932                }
933                PageEntry::new(pn, data)
934            })
935            .collect()
936    }
937
938    // -----------------------------------------------------------------------
939    // Resume state tests
940    // -----------------------------------------------------------------------
941
942    #[test]
943    fn test_resume_state_persistence() {
944        let mut resume = ResumeState::new(3);
945        resume.blocks[0].record_isi(0);
946        resume.blocks[0].record_isi(5);
947        resume.blocks[0].record_isi(10);
948        resume.blocks[1].decoded = true;
949
950        let bytes = resume.to_bytes();
951        let restored = ResumeState::from_bytes(&bytes).expect("deserialize");
952
953        assert_eq!(
954            restored.total_blocks, 3,
955            "bead_id={TEST_BEAD_ID} case=resume_total_blocks"
956        );
957        assert_eq!(
958            restored.blocks[0].num_received, 3,
959            "bead_id={TEST_BEAD_ID} case=resume_block0_received"
960        );
961        assert!(
962            restored.blocks[0].received_isis.contains(&5),
963            "bead_id={TEST_BEAD_ID} case=resume_block0_isi_5"
964        );
965        assert!(
966            restored.blocks[1].decoded,
967            "bead_id={TEST_BEAD_ID} case=resume_block1_decoded"
968        );
969        assert!(
970            !restored.blocks[2].decoded,
971            "bead_id={TEST_BEAD_ID} case=resume_block2_not_decoded"
972        );
973    }
974
975    #[test]
976    fn test_resume_no_protocol_negotiation() {
977        // Resume state works without any sender-side coordination.
978        let mut resume = ResumeState::new(2);
979        resume.blocks[0].record_isi(0);
980        resume.blocks[0].record_isi(1);
981
982        // Persist and restore.
983        let bytes = resume.to_bytes();
984        let restored = ResumeState::from_bytes(&bytes).expect("deserialize");
985        assert_eq!(
986            restored.blocks[0].num_received, 2,
987            "bead_id={TEST_BEAD_ID} case=resume_no_negotiation"
988        );
989        assert!(!restored.all_decoded());
990    }
991
992    // -----------------------------------------------------------------------
993    // Snapshot sender/receiver integration
994    // -----------------------------------------------------------------------
995
996    #[test]
997    fn test_snapshot_single_block() {
998        let page_size = 256_u32;
999        let page_numbers: Vec<u32> = (1..=10).collect();
1000        let mut pages = make_pages(page_size, &page_numbers);
1001
1002        let config = SenderConfig {
1003            symbol_size: 256,
1004            max_isi_multiplier: 1,
1005        };
1006        let mut sender = SnapshotSender::prepare(page_size, &mut pages, config).expect("prepare");
1007        assert_eq!(
1008            sender.num_blocks(),
1009            1,
1010            "bead_id={TEST_BEAD_ID} case=single_block"
1011        );
1012
1013        // Collect all packets.
1014        let mut packets = Vec::new();
1015        while let Some(pkt) = sender.next_packet() {
1016            packets.push(pkt);
1017        }
1018        assert!(
1019            !packets.is_empty(),
1020            "bead_id={TEST_BEAD_ID} case=has_packets"
1021        );
1022
1023        // Feed to receiver.
1024        let mut receiver = SnapshotReceiver::new(1, page_size);
1025        for pkt in &packets {
1026            let _ = receiver.process_packet(pkt);
1027        }
1028
1029        assert_eq!(
1030            receiver.state(),
1031            SnapshotReceiverState::Complete,
1032            "bead_id={TEST_BEAD_ID} case=single_block_complete"
1033        );
1034
1035        let blocks = receiver.take_decoded_blocks();
1036        assert_eq!(blocks.len(), 1);
1037        assert_eq!(blocks[0].pages.len(), 10);
1038    }
1039
1040    #[test]
1041    fn test_snapshot_multi_block_small() {
1042        // Force multi-block by using many pages.
1043        // Use smaller page count that still creates multiple blocks
1044        // by using the sender's internal sharding mechanism.
1045        let page_size = 64_u32;
1046        let n_pages = 200_u32;
1047        let page_numbers: Vec<u32> = (1..=n_pages).collect();
1048        let mut pages = make_pages(page_size, &page_numbers);
1049
1050        let config = SenderConfig {
1051            symbol_size: 64,
1052            max_isi_multiplier: 1,
1053        };
1054        let mut sender = SnapshotSender::prepare(page_size, &mut pages, config).expect("prepare");
1055
1056        // Should be 1 block (200 < K_MAX).
1057        assert_eq!(
1058            sender.num_blocks(),
1059            1,
1060            "bead_id={TEST_BEAD_ID} case=multi_block_small_count"
1061        );
1062
1063        let mut packets = Vec::new();
1064        while let Some(pkt) = sender.next_packet() {
1065            packets.push(pkt);
1066        }
1067
1068        let mut receiver = SnapshotReceiver::new(sender.num_blocks(), page_size);
1069        for pkt in &packets {
1070            let _ = receiver.process_packet(pkt);
1071        }
1072
1073        assert_eq!(
1074            receiver.state(),
1075            SnapshotReceiverState::Complete,
1076            "bead_id={TEST_BEAD_ID} case=multi_block_small_complete"
1077        );
1078
1079        let blocks = receiver.take_decoded_blocks();
1080        let total_pages: usize = blocks.iter().map(|b| b.pages.len()).sum();
1081        assert_eq!(
1082            total_pages, n_pages as usize,
1083            "bead_id={TEST_BEAD_ID} case=multi_block_all_pages"
1084        );
1085    }
1086
1087    #[test]
1088    fn test_duplicate_isi_discarded() {
1089        let page_size = 128_u32;
1090        let mut pages = make_pages(page_size, &[1, 2, 3]);
1091        let config = SenderConfig {
1092            symbol_size: 128,
1093            max_isi_multiplier: 1,
1094        };
1095        let mut sender = SnapshotSender::prepare(page_size, &mut pages, config).expect("prepare");
1096
1097        let mut packets = Vec::new();
1098        while let Some(pkt) = sender.next_packet() {
1099            packets.push(pkt);
1100        }
1101
1102        let mut receiver = SnapshotReceiver::new(1, page_size);
1103
1104        // Feed first packet twice.
1105        let r1 = receiver.process_packet(&packets[0]).expect("first");
1106        assert_ne!(
1107            r1,
1108            SnapshotPacketResult::Duplicate,
1109            "bead_id={TEST_BEAD_ID} case=first_not_dup"
1110        );
1111        let r2 = receiver.process_packet(&packets[0]).expect("duplicate");
1112        assert_eq!(
1113            r2,
1114            SnapshotPacketResult::Duplicate,
1115            "bead_id={TEST_BEAD_ID} case=dup_discarded"
1116        );
1117    }
1118
1119    #[test]
1120    fn test_snapshot_progressive_receive() {
1121        // With a single block, after decode the receiver is complete.
1122        // Progressive receive means we can query pages from decoded blocks
1123        // while other blocks are still being received.
1124        let page_size = 128_u32;
1125        let mut pages = make_pages(page_size, &[1, 2, 3, 4, 5]);
1126        let config = SenderConfig {
1127            symbol_size: 128,
1128            max_isi_multiplier: 1,
1129        };
1130        let mut sender = SnapshotSender::prepare(page_size, &mut pages, config).expect("prepare");
1131
1132        let mut packets = Vec::new();
1133        while let Some(pkt) = sender.next_packet() {
1134            packets.push(pkt);
1135        }
1136
1137        let mut receiver = SnapshotReceiver::new(1, page_size);
1138        let mut block_decoded_at = None;
1139
1140        for (i, pkt) in packets.iter().enumerate() {
1141            if let Ok(SnapshotPacketResult::BlockDecoded(_)) = receiver.process_packet(pkt) {
1142                block_decoded_at = Some(i);
1143                break;
1144            }
1145        }
1146
1147        assert!(
1148            block_decoded_at.is_some(),
1149            "bead_id={TEST_BEAD_ID} case=progressive_block_decoded"
1150        );
1151
1152        // After decoding, pages are available.
1153        let blocks = receiver.take_decoded_blocks();
1154        assert!(
1155            !blocks.is_empty(),
1156            "bead_id={TEST_BEAD_ID} case=progressive_has_pages"
1157        );
1158    }
1159
1160    // -----------------------------------------------------------------------
1161    // E2E tests
1162    // -----------------------------------------------------------------------
1163
1164    #[test]
1165    fn test_e2e_sender_receiver_roundtrip() {
1166        let page_size = 512_u32;
1167        let n_pages = 50_u32;
1168        let page_numbers: Vec<u32> = (1..=n_pages).collect();
1169        let original_pages = make_pages(page_size, &page_numbers);
1170        let mut pages = original_pages.clone();
1171
1172        let config = SenderConfig {
1173            symbol_size: 512,
1174            max_isi_multiplier: 1,
1175        };
1176        let mut sender = SnapshotSender::prepare(page_size, &mut pages, config).expect("prepare");
1177
1178        let mut packets = Vec::new();
1179        while let Some(pkt) = sender.next_packet() {
1180            packets.push(pkt);
1181        }
1182
1183        let mut receiver = SnapshotReceiver::new(sender.num_blocks(), page_size);
1184        for pkt in &packets {
1185            let _ = receiver.process_packet(pkt);
1186        }
1187
1188        assert_eq!(
1189            receiver.state(),
1190            SnapshotReceiverState::Complete,
1191            "bead_id={TEST_BEAD_ID} case=e2e_roundtrip_complete"
1192        );
1193
1194        let blocks = receiver.take_decoded_blocks();
1195        let mut all_decoded_pages: Vec<&DecodedBlockPage> =
1196            blocks.iter().flat_map(|b| b.pages.iter()).collect();
1197        all_decoded_pages.sort_by_key(|p| p.page_number);
1198
1199        assert_eq!(
1200            all_decoded_pages.len(),
1201            original_pages.len(),
1202            "bead_id={TEST_BEAD_ID} case=e2e_page_count"
1203        );
1204
1205        for (decoded, original) in all_decoded_pages.iter().zip(original_pages.iter()) {
1206            assert_eq!(
1207                decoded.page_number, original.page_number,
1208                "bead_id={TEST_BEAD_ID} case=e2e_page_number"
1209            );
1210            assert_eq!(
1211                decoded.page_data, original.page_bytes,
1212                "bead_id={TEST_BEAD_ID} case=e2e_page_data pn={}",
1213                original.page_number
1214            );
1215        }
1216    }
1217
1218    #[test]
1219    fn test_e2e_resume_after_partial() {
1220        let page_size = 128_u32;
1221        let n_pages = 20_u32;
1222        let mut pages = make_pages(page_size, &(1..=n_pages).collect::<Vec<_>>());
1223
1224        let config = SenderConfig {
1225            symbol_size: 128,
1226            max_isi_multiplier: 1,
1227        };
1228        let mut sender = SnapshotSender::prepare(page_size, &mut pages, config).expect("prepare");
1229
1230        let mut packets = Vec::new();
1231        while let Some(pkt) = sender.next_packet() {
1232            packets.push(pkt);
1233        }
1234
1235        // First receiver: receive only half the packets.
1236        let half = packets.len() / 2;
1237        let mut receiver1 = SnapshotReceiver::new(sender.num_blocks(), page_size);
1238        for pkt in &packets[..half] {
1239            let _ = receiver1.process_packet(pkt);
1240        }
1241
1242        // Persist resume state.
1243        let resume_bytes = receiver1.resume_state().to_bytes();
1244
1245        // "Crash" — create new receiver from resume state.
1246        let resume = ResumeState::from_bytes(&resume_bytes).expect("restore");
1247        let mut receiver2 = SnapshotReceiver::from_resume(resume, page_size);
1248
1249        // Continue with remaining packets (and possibly some overlap).
1250        for pkt in &packets {
1251            let _ = receiver2.process_packet(pkt);
1252        }
1253
1254        // Should be complete now.
1255        assert_eq!(
1256            receiver2.state(),
1257            SnapshotReceiverState::Complete,
1258            "bead_id={TEST_BEAD_ID} case=e2e_resume_complete"
1259        );
1260    }
1261
1262    #[test]
1263    fn test_e2e_bd_1hi_15_compliance() {
1264        // Full compliance test.
1265        let page_size = 256_u32;
1266        let n_pages = 30_u32;
1267        let original_pages = make_pages(page_size, &(1..=n_pages).collect::<Vec<_>>());
1268        let mut pages = original_pages;
1269
1270        let config = SenderConfig {
1271            symbol_size: 256,
1272            max_isi_multiplier: 1,
1273        };
1274        let mut sender = SnapshotSender::prepare(page_size, &mut pages, config).expect("prepare");
1275
1276        // Verify sender state.
1277        assert!(
1278            sender.num_blocks() >= 1,
1279            "bead_id={TEST_BEAD_ID} case=compliance_has_blocks"
1280        );
1281        assert!(
1282            sender.total_source_symbols() > 0,
1283            "bead_id={TEST_BEAD_ID} case=compliance_has_symbols"
1284        );
1285
1286        let mut packets = Vec::new();
1287        while let Some(pkt) = sender.next_packet() {
1288            packets.push(pkt);
1289        }
1290
1291        let mut receiver = SnapshotReceiver::new(sender.num_blocks(), page_size);
1292        assert_eq!(receiver.state(), SnapshotReceiverState::Waiting);
1293
1294        for pkt in &packets {
1295            let _ = receiver.process_packet(pkt);
1296        }
1297        assert_eq!(receiver.state(), SnapshotReceiverState::Complete);
1298
1299        let blocks = receiver.take_decoded_blocks();
1300        let total_decoded: usize = blocks.iter().map(|b| b.pages.len()).sum();
1301        assert_eq!(
1302            total_decoded, n_pages as usize,
1303            "bead_id={TEST_BEAD_ID} case=compliance_all_pages_decoded"
1304        );
1305
1306        // Verify resume state.
1307        assert!(
1308            receiver.resume_state().all_decoded(),
1309            "bead_id={TEST_BEAD_ID} case=compliance_resume_all_decoded"
1310        );
1311    }
1312
1313    // -----------------------------------------------------------------------
1314    // Property tests
1315    // -----------------------------------------------------------------------
1316
1317    #[test]
1318    fn prop_partition_covers_all_pages() {
1319        for p in [1_u32, 10, 100, 1000, 56_403, 56_404, 100_000] {
1320            let blocks = partition_source_blocks(p).expect("partition");
1321            let total: u32 = blocks.iter().map(|b| b.num_pages).sum();
1322            assert_eq!(
1323                total, p,
1324                "bead_id={TEST_BEAD_ID} case=prop_partition_covers p={p}"
1325            );
1326        }
1327    }
1328
1329    #[test]
1330    fn prop_partition_block_sizes_valid() {
1331        for p in [1_u32, 56_403, 56_404, 200_000] {
1332            let blocks = partition_source_blocks(p).expect("partition");
1333            for block in &blocks {
1334                assert!(
1335                    block.num_pages <= K_MAX,
1336                    "bead_id={TEST_BEAD_ID} case=prop_block_size p={p} block={} num_pages={}",
1337                    block.index,
1338                    block.num_pages
1339                );
1340            }
1341        }
1342    }
1343
1344    // -----------------------------------------------------------------------
1345    // Compliance gate tests
1346    // -----------------------------------------------------------------------
1347
1348    #[test]
1349    fn test_bd_1hi_15_unit_compliance_gate() {
1350        // Verify all required types exist.
1351        let _ = SnapshotReceiverState::Waiting;
1352        let _ = SnapshotReceiverState::Receiving;
1353        let _ = SnapshotReceiverState::Complete;
1354
1355        let _ = SnapshotPacketResult::Accepted;
1356        let _ = SnapshotPacketResult::Duplicate;
1357        let _ = SnapshotPacketResult::Rejected;
1358        let _ = SnapshotPacketResult::AlreadyComplete;
1359
1360        let resume = ResumeState::new(3);
1361        assert_eq!(resume.total_blocks, 3);
1362        assert!(!resume.all_decoded());
1363        assert_eq!(resume.decoded_count(), 0);
1364
1365        // Verify BlockResumeState serialization.
1366        let block = BlockResumeState::new(0);
1367        let bytes = block.to_bytes();
1368        let (restored, _) = BlockResumeState::from_bytes(&bytes).expect("deser");
1369        assert_eq!(restored.block_id, 0);
1370    }
1371
1372    #[test]
1373    fn prop_bd_1hi_15_structure_compliance() {
1374        // Verify snapshot sender + receiver integration.
1375        let page_size = 128_u32;
1376        let mut pages = make_pages(page_size, &[1, 2]);
1377        let config = SenderConfig {
1378            symbol_size: 128,
1379            max_isi_multiplier: 1,
1380        };
1381        let mut sender = SnapshotSender::prepare(page_size, &mut pages, config).expect("prepare");
1382        assert!(sender.num_blocks() >= 1);
1383
1384        let mut packets = Vec::new();
1385        while let Some(pkt) = sender.next_packet() {
1386            packets.push(pkt);
1387        }
1388
1389        let mut receiver = SnapshotReceiver::new(sender.num_blocks(), page_size);
1390        for pkt in &packets {
1391            let _ = receiver.process_packet(pkt);
1392        }
1393        assert_eq!(receiver.state(), SnapshotReceiverState::Complete);
1394    }
1395}