Skip to main content

asupersync/
encoding.rs

1//! RaptorQ encoding pipeline (Phase 0).
2//!
3//! This module provides a deterministic, streaming encoder that splits input
4//! bytes into source symbols and produces a configurable number of repair
5//! symbols per block. Repair symbols are generated via the systematic
6//! RaptorQ encoder (precode + LT) for deterministic RFC-6330-style behavior.
7
8use crate::config::EncodingConfig;
9use crate::error::{Error, ErrorKind};
10use crate::raptorq::systematic::{SystematicEncoder, SystematicParamError, SystematicParams};
11use crate::types::resource::{PoolExhausted, SymbolPool};
12use crate::types::{ObjectId, Symbol, SymbolId, SymbolKind};
13use std::cmp::min;
14
15/// The symbol ID format caps objects at 256 source blocks.
16pub(crate) const MAX_SOURCE_BLOCKS: usize = u8::MAX as usize + 1;
17
18/// Returns the maximum object size supported by the byte-based block contract.
19#[must_use]
20#[inline]
21pub(crate) fn max_object_size(max_block_size: usize) -> usize {
22    max_block_size.saturating_mul(MAX_SOURCE_BLOCKS)
23}
24
25/// Errors produced by the encoding pipeline.
26#[derive(Debug, thiserror::Error)]
27pub enum EncodingError {
28    /// Input data exceeds the configured maximum.
29    #[error("data too large: {size} bytes exceeds limit {limit}")]
30    DataTooLarge {
31        /// Input size in bytes.
32        size: usize,
33        /// Maximum allowed size in bytes.
34        limit: usize,
35    },
36    /// The symbol pool could not supply a buffer.
37    #[error("symbol pool exhausted")]
38    PoolExhausted,
39    /// Configuration is invalid or inconsistent.
40    #[error("invalid configuration: {reason}")]
41    InvalidConfig {
42        /// Reason for invalid configuration.
43        reason: String,
44    },
45    /// The encoding computation failed.
46    #[error("encoding failed: {details}")]
47    ComputationFailed {
48        /// Details of the failure.
49        details: String,
50    },
51}
52
53impl From<PoolExhausted> for EncodingError {
54    #[inline]
55    fn from(_: PoolExhausted) -> Self {
56        Self::PoolExhausted
57    }
58}
59
60impl From<EncodingError> for Error {
61    fn from(err: EncodingError) -> Self {
62        match &err {
63            EncodingError::DataTooLarge { .. } | EncodingError::InvalidConfig { .. } => {
64                Self::new(ErrorKind::InvalidEncodingParams)
65            }
66            EncodingError::PoolExhausted => {
67                Self::new(ErrorKind::EncodingFailed).with_message("symbol pool exhausted")
68            }
69            EncodingError::ComputationFailed { .. } => Self::new(ErrorKind::EncodingFailed),
70        }
71        .with_message(err.to_string())
72    }
73}
74
75/// Encoder output with metadata.
76#[derive(Debug, Clone)]
77pub struct EncodedSymbol {
78    symbol: Symbol,
79}
80
81impl EncodedSymbol {
82    /// Creates a new encoded symbol wrapper.
83    #[must_use]
84    #[inline]
85    pub const fn new(symbol: Symbol) -> Self {
86        Self { symbol }
87    }
88
89    /// Returns the underlying symbol.
90    #[must_use]
91    #[inline]
92    pub const fn symbol(&self) -> &Symbol {
93        &self.symbol
94    }
95
96    /// Consumes the wrapper and returns the symbol.
97    #[must_use]
98    #[inline]
99    pub fn into_symbol(self) -> Symbol {
100        self.symbol
101    }
102
103    /// Returns the symbol ID.
104    #[must_use]
105    #[inline]
106    pub const fn id(&self) -> SymbolId {
107        self.symbol.id()
108    }
109
110    /// Returns the symbol kind.
111    #[must_use]
112    #[inline]
113    pub const fn kind(&self) -> SymbolKind {
114        self.symbol.kind()
115    }
116}
117
118/// Statistics for the most recent encoding run.
119#[derive(Debug, Clone, Copy, Default)]
120pub struct EncodingStats {
121    /// Input bytes consumed.
122    pub bytes_in: usize,
123    /// Number of blocks encoded.
124    pub blocks: usize,
125    /// Source symbols emitted.
126    pub source_symbols: usize,
127    /// Repair symbols emitted.
128    pub repair_symbols: usize,
129}
130
131impl EncodingStats {
132    fn reset_for(&mut self, bytes_in: usize, blocks: usize) {
133        *self = Self {
134            bytes_in,
135            blocks,
136            source_symbols: 0,
137            repair_symbols: 0,
138        };
139    }
140}
141
142/// Main encoding pipeline.
143#[derive(Debug)]
144pub struct EncodingPipeline {
145    config: EncodingConfig,
146    pool: SymbolPool,
147    stats: EncodingStats,
148}
149
150impl EncodingPipeline {
151    /// Creates a new encoding pipeline.
152    #[must_use]
153    #[inline]
154    pub fn new(config: EncodingConfig, pool: SymbolPool) -> Self {
155        Self {
156            config,
157            pool,
158            stats: EncodingStats::default(),
159        }
160    }
161
162    /// Returns encoding statistics for the most recent run.
163    #[must_use]
164    #[inline]
165    pub const fn stats(&self) -> EncodingStats {
166        self.stats
167    }
168
169    /// Resets internal statistics.
170    #[inline]
171    pub fn reset(&mut self) {
172        self.stats = EncodingStats::default();
173    }
174
175    /// Encodes data using the configured repair overhead.
176    pub fn encode<'a>(&'a mut self, object_id: ObjectId, data: &'a [u8]) -> EncodingIterator<'a> {
177        self.encode_internal(object_id, data, None)
178    }
179
180    /// Encodes data with an explicit repair count per block.
181    pub fn encode_with_repair<'a>(
182        &'a mut self,
183        object_id: ObjectId,
184        data: &'a [u8],
185        repair_count: usize,
186    ) -> EncodingIterator<'a> {
187        self.encode_internal(object_id, data, Some(repair_count))
188    }
189
190    /// Encodes only repair symbols for each block, starting at `first_repair`.
191    ///
192    /// `first_repair` is zero-based relative to the block's first repair symbol,
193    /// so `first_repair == 0` emits public ESI `K`, `first_repair == 1` emits
194    /// public ESI `K + 1`, and so on. This avoids emitting or copying
195    /// systematic source symbols when a transport only needs fresh RaptorQ
196    /// repair.
197    #[cfg_attr(target_arch = "wasm32", allow(dead_code))]
198    pub(crate) fn encode_repair_range<'a>(
199        &'a mut self,
200        object_id: ObjectId,
201        data: &'a [u8],
202        first_repair: usize,
203        repair_count: usize,
204    ) -> RepairEncodingIterator<'a> {
205        let (blocks, symbol_size, plan_error) = match self.plan_blocks(data) {
206            Ok((blocks, symbol_size)) => (blocks, symbol_size, None),
207            Err(err) => (Vec::new(), 0, Some(err)),
208        };
209
210        self.stats.reset_for(data.len(), blocks.len());
211
212        RepairEncodingIterator {
213            pipeline: self,
214            object_id,
215            data,
216            blocks,
217            block_index: 0,
218            repair_index: 0,
219            first_repair,
220            repair_count,
221            symbol_size,
222            plan_error,
223            systematic_encoder: None,
224            systematic_block_index: None,
225        }
226    }
227
228    /// Encodes one source block with its already-assigned source block number.
229    ///
230    /// This is used by bounded-memory transports that read one block from disk
231    /// at a time but must preserve the same object/SBN layout as
232    /// [`Self::encode_with_repair`] would have produced for the whole object.
233    #[cfg_attr(target_arch = "wasm32", allow(dead_code))]
234    pub(crate) fn encode_single_block_with_repair<'a>(
235        &'a mut self,
236        object_id: ObjectId,
237        sbn: u8,
238        data: &'a [u8],
239        repair_count: usize,
240    ) -> EncodingIterator<'a> {
241        self.encode_single_block_internal(object_id, sbn, data, Some(repair_count))
242    }
243
244    /// Encodes only repair symbols for one source block.
245    #[cfg_attr(target_arch = "wasm32", allow(dead_code))]
246    pub(crate) fn encode_single_block_repair_range<'a>(
247        &'a mut self,
248        object_id: ObjectId,
249        sbn: u8,
250        data: &'a [u8],
251        first_repair: usize,
252        repair_count: usize,
253    ) -> RepairEncodingIterator<'a> {
254        let (blocks, symbol_size, plan_error) = match self.plan_single_block(sbn, data) {
255            Ok((blocks, symbol_size)) => (blocks, symbol_size, None),
256            Err(err) => (Vec::new(), 0, Some(err)),
257        };
258
259        self.stats.reset_for(data.len(), blocks.len());
260
261        RepairEncodingIterator {
262            pipeline: self,
263            object_id,
264            data,
265            blocks,
266            block_index: 0,
267            repair_index: 0,
268            first_repair,
269            repair_count,
270            symbol_size,
271            plan_error,
272            systematic_encoder: None,
273            systematic_block_index: None,
274        }
275    }
276
277    fn encode_internal<'a>(
278        &'a mut self,
279        object_id: ObjectId,
280        data: &'a [u8],
281        repair_override: Option<usize>,
282    ) -> EncodingIterator<'a> {
283        let (blocks, symbol_size, plan_error) = match self.plan_blocks(data) {
284            Ok((blocks, symbol_size)) => (blocks, symbol_size, None),
285            Err(err) => (Vec::new(), 0, Some(err)),
286        };
287
288        self.stats.reset_for(data.len(), blocks.len());
289
290        EncodingIterator {
291            pipeline: self,
292            object_id,
293            data,
294            blocks,
295            block_index: 0,
296            esi: 0,
297            symbol_size,
298            repair_override,
299            plan_error,
300            systematic_encoder: None,
301            systematic_block_index: None,
302        }
303    }
304
305    #[cfg_attr(target_arch = "wasm32", allow(dead_code))]
306    fn encode_single_block_internal<'a>(
307        &'a mut self,
308        object_id: ObjectId,
309        sbn: u8,
310        data: &'a [u8],
311        repair_override: Option<usize>,
312    ) -> EncodingIterator<'a> {
313        let (blocks, symbol_size, plan_error) = match self.plan_single_block(sbn, data) {
314            Ok((blocks, symbol_size)) => (blocks, symbol_size, None),
315            Err(err) => (Vec::new(), 0, Some(err)),
316        };
317
318        self.stats.reset_for(data.len(), blocks.len());
319
320        EncodingIterator {
321            pipeline: self,
322            object_id,
323            data,
324            blocks,
325            block_index: 0,
326            esi: 0,
327            symbol_size,
328            repair_override,
329            plan_error,
330            systematic_encoder: None,
331            systematic_block_index: None,
332        }
333    }
334
335    fn plan_blocks(&self, data: &[u8]) -> Result<(Vec<BlockPlan>, usize), EncodingError> {
336        let symbol_size = self.validate_config()?;
337
338        if data.is_empty() {
339            return Ok((Vec::new(), symbol_size));
340        }
341
342        let max_total = max_object_size(self.config.max_block_size);
343        if data.len() > max_total {
344            return Err(EncodingError::DataTooLarge {
345                size: data.len(),
346                limit: max_total,
347            });
348        }
349
350        let mut blocks = Vec::new();
351        let mut offset = 0;
352        let mut sbn: u8 = 0;
353
354        while offset < data.len() {
355            let len = min(self.config.max_block_size, data.len() - offset);
356            let k = len.div_ceil(symbol_size);
357            validate_source_block_k(len, symbol_size, k)?;
358            blocks.push(BlockPlan {
359                sbn,
360                start: offset,
361                len,
362                k,
363            });
364            offset += len;
365            sbn = sbn.wrapping_add(1);
366        }
367
368        Ok((blocks, symbol_size))
369    }
370
371    #[cfg_attr(target_arch = "wasm32", allow(dead_code))]
372    fn plan_single_block(
373        &self,
374        sbn: u8,
375        data: &[u8],
376    ) -> Result<(Vec<BlockPlan>, usize), EncodingError> {
377        let symbol_size = self.validate_config()?;
378
379        if data.is_empty() {
380            return Ok((Vec::new(), symbol_size));
381        }
382        if data.len() > self.config.max_block_size {
383            return Err(EncodingError::DataTooLarge {
384                size: data.len(),
385                limit: self.config.max_block_size,
386            });
387        }
388
389        let k = data.len().div_ceil(symbol_size);
390        validate_source_block_k(data.len(), symbol_size, k)?;
391        Ok((
392            vec![BlockPlan {
393                sbn,
394                start: 0,
395                len: data.len(),
396                k,
397            }],
398            symbol_size,
399        ))
400    }
401
402    fn validate_config(&self) -> Result<usize, EncodingError> {
403        let symbol_size = usize::from(self.config.symbol_size);
404        if symbol_size == 0 {
405            return Err(EncodingError::InvalidConfig {
406                reason: "symbol_size must be non-zero".to_string(),
407            });
408        }
409
410        if self.config.max_block_size == 0 {
411            return Err(EncodingError::InvalidConfig {
412                reason: "max_block_size must be non-zero".to_string(),
413            });
414        }
415
416        if !self.config.repair_overhead.is_finite() || self.config.repair_overhead < 1.0 {
417            return Err(EncodingError::InvalidConfig {
418                reason: "repair_overhead must be finite and >= 1.0".to_string(),
419            });
420        }
421
422        if self.pool_enabled() && self.pool.config().symbol_size != self.config.symbol_size {
423            return Err(EncodingError::InvalidConfig {
424                reason: format!(
425                    "pool symbol_size {} does not match encoding symbol_size {}",
426                    self.pool.config().symbol_size,
427                    self.config.symbol_size
428                ),
429            });
430        }
431
432        Ok(symbol_size)
433    }
434
435    fn pool_enabled(&self) -> bool {
436        let config = self.pool.config();
437        config.max_size > 0 || config.initial_size > 0 || config.allow_growth
438    }
439
440    fn allocate_buffer(&mut self, symbol_size: usize) -> Result<Vec<u8>, EncodingError> {
441        if self.pool_enabled() {
442            let buffer = self.pool.allocate()?;
443            if buffer.len() != symbol_size {
444                return Err(EncodingError::InvalidConfig {
445                    reason: format!(
446                        "pool buffer size {} does not match symbol_size {}",
447                        buffer.len(),
448                        symbol_size
449                    ),
450                });
451            }
452            Ok(Vec::from(buffer.into_boxed_slice()))
453        } else {
454            Ok(vec![0_u8; symbol_size])
455        }
456    }
457}
458
459fn validate_source_block_k(
460    block_len: usize,
461    symbol_size: usize,
462    k: usize,
463) -> Result<(), EncodingError> {
464    SystematicParams::try_for_source_block(k, symbol_size)
465        .map(|_| ())
466        .map_err(|err| match err {
467        SystematicParamError::UnsupportedSourceBlockSize {
468            requested,
469            max_supported,
470        } => EncodingError::InvalidConfig {
471            reason: format!(
472                "block of {block_len} bytes with symbol_size {symbol_size} requires unsupported source block K={requested}; supported range is 1..={max_supported}"
473            ),
474        },
475        SystematicParamError::KPrimeExceedsU32 {
476            k_prime,
477            max_u32,
478        } => EncodingError::InvalidConfig {
479            reason: format!(
480                "block of {block_len} bytes with symbol_size {symbol_size} requires K'={k_prime} which exceeds u32::MAX ({max_u32}); ESI calculations would overflow"
481            ),
482        },
483        SystematicParamError::RfcTableInvariantViolation {
484            invariant,
485            details,
486        } => EncodingError::InvalidConfig {
487            reason: format!(
488                "block of {block_len} bytes with symbol_size {symbol_size} triggers RFC 6330 table invariant violation: {invariant} - {details}"
489            ),
490        },
491    })
492}
493
494/// Iterator over encoded symbols.
495pub struct EncodingIterator<'a> {
496    pipeline: &'a mut EncodingPipeline,
497    object_id: ObjectId,
498    data: &'a [u8],
499    blocks: Vec<BlockPlan>,
500    block_index: usize,
501    esi: u32,
502    symbol_size: usize,
503    repair_override: Option<usize>,
504    plan_error: Option<EncodingError>,
505    systematic_encoder: Option<SystematicEncoder>,
506    systematic_block_index: Option<usize>,
507}
508
509impl Iterator for EncodingIterator<'_> {
510    type Item = Result<EncodedSymbol, EncodingError>;
511
512    fn next(&mut self) -> Option<Self::Item> {
513        if let Some(err) = self.plan_error.take() {
514            return Some(Err(err));
515        }
516
517        while self.block_index < self.blocks.len() {
518            let block = self.blocks[self.block_index].clone();
519            let k = u32::try_from(block.k).unwrap_or(u32::MAX);
520            if k == 0 {
521                self.block_index += 1;
522                self.esi = 0;
523                self.systematic_encoder = None;
524                self.systematic_block_index = None;
525                continue;
526            }
527
528            let repair = u32::try_from(self.repair_override.unwrap_or_else(|| {
529                compute_repair_count(block.k, self.pipeline.config.repair_overhead)
530            }))
531            .unwrap_or(u32::MAX);
532            let total = k.saturating_add(repair);
533
534            if self.esi >= total {
535                self.block_index += 1;
536                self.esi = 0;
537                self.systematic_encoder = None;
538                self.systematic_block_index = None;
539                continue;
540            }
541
542            let esi = self.esi;
543            self.esi = self.esi.saturating_add(1);
544
545            let result = if esi < k {
546                self.emit_source(&block, esi)
547            } else {
548                self.emit_repair(&block, esi)
549            };
550
551            return Some(result.map(EncodedSymbol::new));
552        }
553
554        None
555    }
556}
557
558impl EncodingIterator<'_> {
559    fn emit_source(&mut self, block: &BlockPlan, esi: u32) -> Result<Symbol, EncodingError> {
560        let mut buffer = self.pipeline.allocate_buffer(self.symbol_size)?;
561        let start = block.start + (esi as usize * self.symbol_size);
562        let end = min(start + self.symbol_size, block.end());
563        let copy_len = end.saturating_sub(start);
564
565        if copy_len < buffer.len() {
566            buffer.fill(0);
567        }
568
569        if copy_len > 0 {
570            let slice = &self.data[start..end];
571            buffer[..slice.len()].copy_from_slice(slice);
572        }
573
574        self.pipeline.stats.source_symbols += 1;
575        Ok(Symbol::new(
576            SymbolId::new(self.object_id, block.sbn, esi),
577            buffer,
578            SymbolKind::Source,
579        ))
580    }
581
582    fn emit_repair(&mut self, block: &BlockPlan, esi: u32) -> Result<Symbol, EncodingError> {
583        let mut buffer = self.pipeline.allocate_buffer(self.symbol_size)?;
584        buffer.fill(0);
585
586        let encoder = self.systematic_encoder_for(block)?;
587        let repair = encoder.repair_symbol(esi);
588        if repair.len() != self.symbol_size {
589            return Err(EncodingError::ComputationFailed {
590                details: format!(
591                    "systematic repair symbol size mismatch: expected {}, got {}",
592                    self.symbol_size,
593                    repair.len()
594                ),
595            });
596        }
597        buffer.copy_from_slice(&repair);
598
599        self.pipeline.stats.repair_symbols += 1;
600        Ok(Symbol::new(
601            SymbolId::new(self.object_id, block.sbn, esi),
602            buffer,
603            SymbolKind::Repair,
604        ))
605    }
606
607    fn systematic_encoder_for(
608        &mut self,
609        block: &BlockPlan,
610    ) -> Result<&SystematicEncoder, EncodingError> {
611        let needs_rebuild = self.systematic_block_index != Some(self.block_index);
612        if needs_rebuild {
613            let source_symbols = build_source_symbols(self.data, block, self.symbol_size);
614            let seed = seed_for_block(self.object_id, block.sbn);
615            let encoder = SystematicEncoder::new(&source_symbols, self.symbol_size, seed)
616                .ok_or_else(|| EncodingError::ComputationFailed {
617                    details: "systematic encoder failed: singular constraint matrix".to_string(),
618                })?;
619            self.systematic_encoder = Some(encoder);
620            self.systematic_block_index = Some(self.block_index);
621        }
622
623        Ok(self
624            .systematic_encoder
625            .as_ref()
626            .expect("systematic encoder must be initialized"))
627    }
628}
629
630/// Iterator over repair-only encoded symbols.
631pub struct RepairEncodingIterator<'a> {
632    pipeline: &'a mut EncodingPipeline,
633    object_id: ObjectId,
634    data: &'a [u8],
635    blocks: Vec<BlockPlan>,
636    block_index: usize,
637    repair_index: usize,
638    first_repair: usize,
639    repair_count: usize,
640    symbol_size: usize,
641    plan_error: Option<EncodingError>,
642    systematic_encoder: Option<SystematicEncoder>,
643    systematic_block_index: Option<usize>,
644}
645
646impl Iterator for RepairEncodingIterator<'_> {
647    type Item = Result<EncodedSymbol, EncodingError>;
648
649    fn next(&mut self) -> Option<Self::Item> {
650        if let Some(err) = self.plan_error.take() {
651            return Some(Err(err));
652        }
653        if self.repair_count == 0 {
654            return None;
655        }
656
657        while self.block_index < self.blocks.len() {
658            let block = self.blocks[self.block_index].clone();
659            if block.k == 0 {
660                self.advance_block();
661                continue;
662            }
663            if self.repair_index >= self.repair_count {
664                self.advance_block();
665                continue;
666            }
667
668            let esi = match repair_esi(block.k, self.first_repair, self.repair_index) {
669                Ok(esi) => esi,
670                Err(err) => return Some(Err(err)),
671            };
672            self.repair_index += 1;
673            return Some(self.emit_repair(&block, esi).map(EncodedSymbol::new));
674        }
675
676        None
677    }
678}
679
680impl RepairEncodingIterator<'_> {
681    fn advance_block(&mut self) {
682        self.block_index += 1;
683        self.repair_index = 0;
684        self.systematic_encoder = None;
685        self.systematic_block_index = None;
686    }
687
688    fn emit_repair(&mut self, block: &BlockPlan, esi: u32) -> Result<Symbol, EncodingError> {
689        let mut buffer = self.pipeline.allocate_buffer(self.symbol_size)?;
690        buffer.fill(0);
691
692        let encoder = self.systematic_encoder_for(block)?;
693        let repair = encoder.repair_symbol(esi);
694        if repair.len() != self.symbol_size {
695            return Err(EncodingError::ComputationFailed {
696                details: format!(
697                    "systematic repair symbol size mismatch: expected {}, got {}",
698                    self.symbol_size,
699                    repair.len()
700                ),
701            });
702        }
703        buffer.copy_from_slice(&repair);
704
705        self.pipeline.stats.repair_symbols += 1;
706        Ok(Symbol::new(
707            SymbolId::new(self.object_id, block.sbn, esi),
708            buffer,
709            SymbolKind::Repair,
710        ))
711    }
712
713    fn systematic_encoder_for(
714        &mut self,
715        block: &BlockPlan,
716    ) -> Result<&SystematicEncoder, EncodingError> {
717        let needs_rebuild = self.systematic_block_index != Some(self.block_index);
718        if needs_rebuild {
719            let source_symbols = build_source_symbols(self.data, block, self.symbol_size);
720            let seed = seed_for_block(self.object_id, block.sbn);
721            let encoder = SystematicEncoder::new(&source_symbols, self.symbol_size, seed)
722                .ok_or_else(|| EncodingError::ComputationFailed {
723                    details: "systematic encoder failed: singular constraint matrix".to_string(),
724                })?;
725            self.systematic_encoder = Some(encoder);
726            self.systematic_block_index = Some(self.block_index);
727        }
728
729        Ok(self
730            .systematic_encoder
731            .as_ref()
732            .expect("systematic encoder must be initialized"))
733    }
734}
735
736fn repair_esi(
737    block_k: usize,
738    first_repair: usize,
739    repair_index: usize,
740) -> Result<u32, EncodingError> {
741    let esi = block_k
742        .checked_add(first_repair)
743        .and_then(|esi| esi.checked_add(repair_index))
744        .ok_or_else(|| EncodingError::InvalidConfig {
745            reason: "repair ESI overflow".to_string(),
746        })?;
747    u32::try_from(esi).map_err(|_| EncodingError::InvalidConfig {
748        reason: format!("repair ESI {esi} exceeds u32::MAX"),
749    })
750}
751
752#[derive(Debug, Clone)]
753struct BlockPlan {
754    sbn: u8,
755    start: usize,
756    len: usize,
757    k: usize,
758}
759
760impl BlockPlan {
761    fn end(&self) -> usize {
762        self.start + self.len
763    }
764}
765
766#[allow(clippy::cast_precision_loss, clippy::cast_sign_loss)]
767fn compute_repair_count(k: usize, overhead: f64) -> usize {
768    let desired = ((k as f64) * overhead).ceil() as usize;
769    desired.saturating_sub(k)
770}
771
772fn seed_for_block(object_id: ObjectId, sbn: u8) -> u64 {
773    seed_for(object_id, sbn, 0)
774}
775
776fn seed_for(object_id: ObjectId, sbn: u8, esi: u32) -> u64 {
777    let obj = object_id.as_u128();
778    let hi = (obj >> 64) as u64;
779    let lo = obj as u64;
780    let mut seed = hi ^ lo.rotate_left(13);
781    seed ^= u64::from(sbn) << 56;
782    seed ^= u64::from(esi);
783    if seed == 0 { 1 } else { seed }
784}
785
786fn build_source_symbols(data: &[u8], block: &BlockPlan, symbol_size: usize) -> Vec<Vec<u8>> {
787    let mut symbols = Vec::with_capacity(block.k);
788    for idx in 0..block.k {
789        let mut buffer = vec![0u8; symbol_size];
790        let start = block.start + idx * symbol_size;
791        let end = min(start + symbol_size, block.end());
792        if start < end {
793            let slice = &data[start..end];
794            buffer[..slice.len()].copy_from_slice(slice);
795        }
796        symbols.push(buffer);
797    }
798    symbols
799}
800
801#[cfg(test)]
802mod tests {
803    #![allow(
804        clippy::pedantic,
805        clippy::nursery,
806        clippy::expect_fun_call,
807        clippy::map_unwrap_or,
808        clippy::cast_possible_wrap,
809        clippy::future_not_send
810    )]
811    use super::*;
812    use crate::types::ObjectId;
813    use crate::types::resource::PoolConfig;
814
815    fn test_config(
816        symbol_size: u16,
817        max_block_size: usize,
818        repair_overhead: f64,
819    ) -> EncodingConfig {
820        EncodingConfig {
821            symbol_size,
822            max_block_size,
823            repair_overhead,
824            encoding_parallelism: 1,
825            decoding_parallelism: 1,
826        }
827    }
828
829    fn pool_for(symbol_size: u16, max_size: usize) -> SymbolPool {
830        SymbolPool::new(PoolConfig {
831            symbol_size,
832            initial_size: max_size,
833            max_size,
834            allow_growth: false,
835            growth_increment: 0,
836        })
837    }
838
839    #[test]
840    fn test_encode_small_data() {
841        let mut pipeline = EncodingPipeline::new(
842            test_config(4, 16, 1.0),
843            SymbolPool::new(PoolConfig::default()),
844        );
845        let data = b"hello";
846        let symbols: Vec<_> = pipeline
847            .encode(ObjectId::new_for_test(1), data)
848            .collect::<Result<Vec<_>, _>>()
849            .unwrap();
850
851        assert_eq!(symbols.len(), 2);
852        assert_eq!(symbols[0].symbol().data().len(), 4);
853        assert_eq!(symbols[1].symbol().data().len(), 4);
854    }
855
856    #[test]
857    fn test_encode_exact_block_size() {
858        let mut pipeline = EncodingPipeline::new(
859            test_config(4, 8, 1.0),
860            SymbolPool::new(PoolConfig::default()),
861        );
862        let data = b"abcdefgh";
863        let symbols: Vec<_> = pipeline
864            .encode(ObjectId::new_for_test(2), data)
865            .collect::<Result<Vec<_>, _>>()
866            .unwrap();
867
868        assert_eq!(symbols.len(), 2);
869        assert!(symbols.iter().all(|s| s.kind() == SymbolKind::Source));
870    }
871
872    #[test]
873    fn test_encode_multiple_blocks() {
874        let mut pipeline = EncodingPipeline::new(
875            test_config(4, 8, 1.0),
876            SymbolPool::new(PoolConfig::default()),
877        );
878        let data = b"abcdefghijklmnop";
879        let symbols: Vec<_> = pipeline
880            .encode(ObjectId::new_for_test(3), data)
881            .collect::<Result<Vec<_>, _>>()
882            .unwrap();
883
884        let sbns: Vec<u8> = symbols.iter().map(|s| s.id().sbn()).collect();
885        assert!(sbns.contains(&0));
886        assert!(sbns.contains(&1));
887    }
888
889    #[test]
890    fn test_encode_multiple_blocks_preserves_non_aligned_boundaries() {
891        let mut pipeline = EncodingPipeline::new(
892            test_config(4, 6, 1.0),
893            SymbolPool::new(PoolConfig::default()),
894        );
895        let data = b"ABCDEFGHIJKLM";
896        let symbols: Vec<_> = pipeline
897            .encode(ObjectId::new_for_test(13), data)
898            .collect::<Result<Vec<_>, _>>()
899            .unwrap();
900
901        let expected = [
902            (0u8, 0u32, b"ABCD".to_vec()),
903            (0u8, 1u32, vec![b'E', b'F', 0, 0]),
904            (1u8, 0u32, b"GHIJ".to_vec()),
905            (1u8, 1u32, vec![b'K', b'L', 0, 0]),
906            (2u8, 0u32, vec![b'M', 0, 0, 0]),
907        ];
908
909        assert_eq!(symbols.len(), expected.len());
910        for (symbol, (expected_sbn, expected_esi, expected_bytes)) in symbols.iter().zip(expected) {
911            assert_eq!(symbol.kind(), SymbolKind::Source);
912            assert_eq!(symbol.id().sbn(), expected_sbn);
913            assert_eq!(symbol.id().esi(), expected_esi);
914            assert_eq!(symbol.symbol().data(), expected_bytes.as_slice());
915        }
916    }
917
918    #[test]
919    fn test_encode_empty_data() {
920        let mut pipeline = EncodingPipeline::new(
921            test_config(8, 32, 1.0),
922            SymbolPool::new(PoolConfig::default()),
923        );
924        let symbols: Vec<_> = pipeline
925            .encode(ObjectId::new_for_test(4), &[])
926            .collect::<Result<Vec<_>, _>>()
927            .unwrap();
928
929        assert!(symbols.is_empty());
930    }
931
932    #[test]
933    fn test_repair_overhead_respected() {
934        let mut pipeline = EncodingPipeline::new(
935            test_config(4, 16, 1.5),
936            SymbolPool::new(PoolConfig::default()),
937        );
938        let data = b"abcdefgh";
939        let symbols: Vec<_> = pipeline
940            .encode(ObjectId::new_for_test(5), data)
941            .collect::<Result<Vec<_>, _>>()
942            .unwrap();
943
944        let repair_count = symbols
945            .iter()
946            .filter(|s| s.kind() == SymbolKind::Repair)
947            .count();
948        assert_eq!(repair_count, 1);
949    }
950
951    #[test]
952    fn test_symbol_ids_unique() {
953        let mut pipeline = EncodingPipeline::new(
954            test_config(4, 16, 1.2),
955            SymbolPool::new(PoolConfig::default()),
956        );
957        let data = b"abcdefgh";
958        let symbols: Vec<_> = pipeline
959            .encode(ObjectId::new_for_test(6), data)
960            .collect::<Result<Vec<_>, _>>()
961            .unwrap();
962
963        let mut ids = symbols.iter().map(EncodedSymbol::id).collect::<Vec<_>>();
964        ids.sort_by_key(|id| (id.sbn(), id.esi()));
965        ids.dedup();
966        assert_eq!(ids.len(), symbols.len());
967    }
968
969    #[test]
970    fn test_data_too_large_error() {
971        let mut pipeline = EncodingPipeline::new(
972            test_config(1, 1, 1.0),
973            SymbolPool::new(PoolConfig::default()),
974        );
975        let data = vec![0_u8; 257];
976        let err = pipeline
977            .encode(ObjectId::new_for_test(7), &data)
978            .next()
979            .unwrap()
980            .unwrap_err();
981
982        assert!(matches!(err, EncodingError::DataTooLarge { .. }));
983    }
984
985    #[test]
986    fn test_pool_exhaustion_handling() {
987        let mut pipeline = EncodingPipeline::new(test_config(4, 16, 1.0), pool_for(4, 1));
988        let data = b"abcdefgh";
989        let mut iter = pipeline.encode(ObjectId::new_for_test(8), data);
990
991        let _ = iter.next().unwrap().unwrap();
992        let err = iter.next().unwrap().unwrap_err();
993        assert!(matches!(err, EncodingError::PoolExhausted));
994    }
995
996    #[test]
997    fn test_source_symbol_zero_pads_with_pool() {
998        let symbol_size = 4u16;
999        let mut pool = pool_for(symbol_size, 1);
1000        let mut buffer = pool.allocate().unwrap();
1001        buffer.as_mut_slice().fill(0xAA);
1002        pool.deallocate(buffer);
1003
1004        let mut pipeline = EncodingPipeline::new(test_config(symbol_size, 16, 1.0), pool);
1005        let data = [0x11u8];
1006        let symbols: Vec<_> = pipeline
1007            .encode(ObjectId::new_for_test(11), &data)
1008            .collect::<Result<Vec<_>, _>>()
1009            .unwrap();
1010
1011        assert_eq!(symbols.len(), 1);
1012        let bytes = symbols[0].symbol().data();
1013        assert_eq!(bytes[0], 0x11);
1014        assert!(bytes[1..].iter().all(|byte| *byte == 0));
1015    }
1016
1017    #[test]
1018    fn test_deterministic_output() {
1019        let data = b"deterministic";
1020        let object_id = ObjectId::new_for_test(9);
1021        let config = test_config(4, 16, 1.5);
1022
1023        let mut pipeline_a =
1024            EncodingPipeline::new(config.clone(), SymbolPool::new(PoolConfig::default()));
1025        let mut pipeline_b = EncodingPipeline::new(config, SymbolPool::new(PoolConfig::default()));
1026
1027        let symbols_a: Vec<_> = pipeline_a
1028            .encode(object_id, data)
1029            .collect::<Result<Vec<_>, _>>()
1030            .unwrap();
1031        let symbols_b: Vec<_> = pipeline_b
1032            .encode(object_id, data)
1033            .collect::<Result<Vec<_>, _>>()
1034            .unwrap();
1035
1036        let bytes_a: Vec<Vec<u8>> = symbols_a
1037            .iter()
1038            .map(|s| s.symbol().data().to_vec())
1039            .collect();
1040        let bytes_b: Vec<Vec<u8>> = symbols_b
1041            .iter()
1042            .map(|s| s.symbol().data().to_vec())
1043            .collect();
1044
1045        assert_eq!(bytes_a, bytes_b);
1046    }
1047
1048    #[test]
1049    fn test_repair_symbols_match_systematic_encoder() {
1050        let symbol_size = 8usize;
1051        let max_block_size = 64usize;
1052        let repair_count = 3usize;
1053        let data: Vec<u8> = (0u8..37).map(|i| i.wrapping_mul(7)).collect();
1054        let object_id = ObjectId::new_for_test(10);
1055
1056        let mut pipeline = EncodingPipeline::new(
1057            test_config(symbol_size as u16, max_block_size, 1.0),
1058            SymbolPool::new(PoolConfig::default()),
1059        );
1060        let symbols: Vec<_> = pipeline
1061            .encode_with_repair(object_id, &data, repair_count)
1062            .collect::<Result<Vec<_>, _>>()
1063            .unwrap();
1064
1065        let k = data.len().div_ceil(symbol_size);
1066        let block = BlockPlan {
1067            sbn: 0,
1068            start: 0,
1069            len: data.len(),
1070            k,
1071        };
1072        let source_symbols = build_source_symbols(&data, &block, symbol_size);
1073        let seed = seed_for_block(object_id, block.sbn);
1074        let encoder = SystematicEncoder::new(&source_symbols, symbol_size, seed).expect("encoder");
1075
1076        for sym in symbols.iter().filter(|s| s.kind() == SymbolKind::Repair) {
1077            let esi = sym.id().esi();
1078            let expected = encoder.repair_symbol(esi);
1079            assert_eq!(sym.symbol().data(), expected.as_slice());
1080        }
1081    }
1082
1083    #[test]
1084    fn test_encode_repair_range_skips_sources_and_preserves_repair_esis() {
1085        let symbol_size = 4usize;
1086        let max_block_size = 8usize;
1087        let data = b"ABCDEFGHIJKLMNOP";
1088        let object_id = ObjectId::new_for_test(14);
1089
1090        let mut pipeline = EncodingPipeline::new(
1091            test_config(symbol_size as u16, max_block_size, 1.0),
1092            SymbolPool::new(PoolConfig::default()),
1093        );
1094        let symbols: Vec<_> = pipeline
1095            .encode_repair_range(object_id, data, 1, 2)
1096            .collect::<Result<Vec<_>, _>>()
1097            .unwrap();
1098
1099        assert_eq!(symbols.len(), 4);
1100        assert!(symbols.iter().all(|s| s.kind() == SymbolKind::Repair));
1101        assert_eq!(pipeline.stats().source_symbols, 0);
1102        assert_eq!(pipeline.stats().repair_symbols, 4);
1103
1104        let ids: Vec<(u8, u32)> = symbols
1105            .iter()
1106            .map(|s| (s.id().sbn(), s.id().esi()))
1107            .collect();
1108        assert_eq!(ids, vec![(0, 3), (0, 4), (1, 3), (1, 4)]);
1109
1110        let first_block = BlockPlan {
1111            sbn: 0,
1112            start: 0,
1113            len: max_block_size,
1114            k: 2,
1115        };
1116        let source_symbols = build_source_symbols(data, &first_block, symbol_size);
1117        let encoder =
1118            SystematicEncoder::new(&source_symbols, symbol_size, seed_for_block(object_id, 0))
1119                .expect("encoder");
1120        assert_eq!(symbols[0].symbol().data(), encoder.repair_symbol(3));
1121        assert_eq!(symbols[1].symbol().data(), encoder.repair_symbol(4));
1122    }
1123
1124    #[test]
1125    fn test_rejects_block_above_systematic_k_limit_before_emission() {
1126        let mut pipeline = EncodingPipeline::new(
1127            test_config(8, 451_232, 1.1),
1128            SymbolPool::new(PoolConfig::default()),
1129        );
1130        let data = vec![0u8; 451_232];
1131
1132        let err = pipeline
1133            .encode_with_repair(ObjectId::new_for_test(12), &data, 1)
1134            .next()
1135            .expect("iterator should yield planning error")
1136            .unwrap_err();
1137
1138        assert!(matches!(err, EncodingError::InvalidConfig { .. }));
1139        assert!(
1140            err.to_string().contains("unsupported source block K=56404"),
1141            "unexpected error: {err}"
1142        );
1143        assert_eq!(pipeline.stats().source_symbols, 0);
1144        assert_eq!(pipeline.stats().repair_symbols, 0);
1145    }
1146
1147    // ========================================================================
1148    // Pure data-type trait coverage (wave 25)
1149    // ========================================================================
1150
1151    #[test]
1152    fn encoding_error_debug_display_data_too_large() {
1153        let err = EncodingError::DataTooLarge {
1154            size: 1024,
1155            limit: 512,
1156        };
1157        let dbg = format!("{err:?}");
1158        assert!(dbg.contains("DataTooLarge"));
1159        let disp = format!("{err}");
1160        assert!(disp.contains("1024"));
1161        assert!(disp.contains("512"));
1162    }
1163
1164    #[test]
1165    fn encoding_error_display_pool_exhausted() {
1166        let err = EncodingError::PoolExhausted;
1167        let disp = format!("{err}");
1168        assert!(disp.contains("pool") || disp.contains("exhausted"));
1169    }
1170
1171    #[test]
1172    fn encoding_error_display_invalid_config() {
1173        let err = EncodingError::InvalidConfig {
1174            reason: "symbol_size must be non-zero".into(),
1175        };
1176        let disp = format!("{err}");
1177        assert!(disp.contains("symbol_size"));
1178    }
1179
1180    #[test]
1181    fn encoding_error_display_computation_failed() {
1182        let err = EncodingError::ComputationFailed {
1183            details: "singular matrix".into(),
1184        };
1185        let disp = format!("{err}");
1186        assert!(disp.contains("singular matrix"));
1187    }
1188
1189    #[test]
1190    fn encoding_error_is_std_error() {
1191        let err = EncodingError::PoolExhausted;
1192        let dyn_err: &dyn std::error::Error = &err;
1193        assert!(!dyn_err.to_string().is_empty());
1194    }
1195
1196    #[test]
1197    fn encoding_error_from_pool_exhausted() {
1198        let pool_err = PoolExhausted;
1199        let encoding_err: EncodingError = pool_err.into();
1200        assert!(matches!(encoding_err, EncodingError::PoolExhausted));
1201    }
1202
1203    #[test]
1204    fn encoding_error_into_error() {
1205        let err = EncodingError::DataTooLarge {
1206            size: 100,
1207            limit: 50,
1208        };
1209        let generic: Error = err.into();
1210        let msg = format!("{generic}");
1211        assert!(!msg.is_empty());
1212
1213        let err2 = EncodingError::PoolExhausted;
1214        let generic2: Error = err2.into();
1215        assert!(!format!("{generic2}").is_empty());
1216
1217        let err3 = EncodingError::InvalidConfig {
1218            reason: "bad".into(),
1219        };
1220        let generic3: Error = err3.into();
1221        assert!(!format!("{generic3}").is_empty());
1222
1223        let err4 = EncodingError::ComputationFailed {
1224            details: "fail".into(),
1225        };
1226        let generic4: Error = err4.into();
1227        assert!(!format!("{generic4}").is_empty());
1228    }
1229
1230    #[test]
1231    fn encoding_stats_debug_clone_copy_default() {
1232        let stats = EncodingStats::default();
1233        assert_eq!(stats.bytes_in, 0);
1234        assert_eq!(stats.blocks, 0);
1235        assert_eq!(stats.source_symbols, 0);
1236        assert_eq!(stats.repair_symbols, 0);
1237        let dbg = format!("{stats:?}");
1238        assert!(dbg.contains("EncodingStats"));
1239        let s2 = stats; // Copy
1240        assert_eq!(s2.bytes_in, stats.bytes_in);
1241    }
1242
1243    #[test]
1244    fn encoding_stats_reset_for() {
1245        let mut stats = EncodingStats {
1246            source_symbols: 10,
1247            repair_symbols: 5,
1248            ..EncodingStats::default()
1249        };
1250        stats.reset_for(1024, 4);
1251        assert_eq!(stats.bytes_in, 1024);
1252        assert_eq!(stats.blocks, 4);
1253        assert_eq!(stats.source_symbols, 0);
1254        assert_eq!(stats.repair_symbols, 0);
1255    }
1256
1257    #[test]
1258    fn encoded_symbol_debug_clone_accessors() {
1259        let sym = Symbol::new(
1260            SymbolId::new(ObjectId::new_for_test(1), 0, 0),
1261            vec![1, 2, 3, 4],
1262            SymbolKind::Source,
1263        );
1264        let encoded = EncodedSymbol::new(sym);
1265        let dbg = format!("{encoded:?}");
1266        assert!(dbg.contains("EncodedSymbol"));
1267        assert_eq!(encoded.kind(), SymbolKind::Source);
1268
1269        let cloned = encoded.clone();
1270        assert_eq!(cloned.symbol().data(), encoded.symbol().data());
1271
1272        let id = encoded.id();
1273        assert_eq!(id.sbn(), 0);
1274        assert_eq!(id.esi(), 0);
1275
1276        let sym_back = encoded.into_symbol();
1277        assert_eq!(sym_back.data(), &[1, 2, 3, 4]);
1278    }
1279
1280    #[test]
1281    fn block_plan_debug_clone_end() {
1282        let plan = BlockPlan {
1283            sbn: 0,
1284            start: 100,
1285            len: 50,
1286            k: 5,
1287        };
1288        let dbg = format!("{plan:?}");
1289        assert!(dbg.contains("BlockPlan"));
1290        let plan2 = plan;
1291        assert_eq!(plan2.end(), 150);
1292        assert_eq!(plan2.sbn, 0);
1293        assert_eq!(plan2.k, 5);
1294    }
1295
1296    #[test]
1297    fn compute_repair_count_cases() {
1298        // overhead 1.0 means 0 repair
1299        assert_eq!(compute_repair_count(10, 1.0), 0);
1300        // overhead 1.5 means ceil(10*1.5)=15, so 5 repair
1301        assert_eq!(compute_repair_count(10, 1.5), 5);
1302        // overhead 2.0 means ceil(10*2.0)=20, so 10 repair
1303        assert_eq!(compute_repair_count(10, 2.0), 10);
1304        // k=1 with overhead 1.5 means ceil(1.5)=2, so 1 repair
1305        assert_eq!(compute_repair_count(1, 1.5), 1);
1306    }
1307
1308    #[test]
1309    fn compute_repair_count_large_k_does_not_truncate() {
1310        // Regression guard: casting k through u32 would wrap this to zero.
1311        let k = (u32::MAX as usize) + 1;
1312        assert_eq!(compute_repair_count(k, 1.25), k / 4);
1313    }
1314
1315    #[test]
1316    fn seed_for_block_deterministic() {
1317        let id = ObjectId::new_for_test(42);
1318        let s1 = seed_for_block(id, 0);
1319        let s2 = seed_for_block(id, 0);
1320        assert_eq!(s1, s2);
1321        // Different blocks should (almost certainly) yield different seeds
1322        let s3 = seed_for_block(id, 1);
1323        assert_ne!(s1, s3);
1324    }
1325
1326    #[test]
1327    fn repair_overhead_nan_rejected() {
1328        let mut pipeline = EncodingPipeline::new(
1329            test_config(4, 16, f64::NAN),
1330            SymbolPool::new(PoolConfig::default()),
1331        );
1332        let err = pipeline
1333            .encode(ObjectId::new_for_test(100), b"test")
1334            .next()
1335            .unwrap()
1336            .unwrap_err();
1337        assert!(matches!(err, EncodingError::InvalidConfig { .. }));
1338    }
1339
1340    #[test]
1341    fn repair_overhead_infinity_rejected() {
1342        let mut pipeline = EncodingPipeline::new(
1343            test_config(4, 16, f64::INFINITY),
1344            SymbolPool::new(PoolConfig::default()),
1345        );
1346        let err = pipeline
1347            .encode(ObjectId::new_for_test(101), b"test")
1348            .next()
1349            .unwrap()
1350            .unwrap_err();
1351        assert!(matches!(err, EncodingError::InvalidConfig { .. }));
1352    }
1353
1354    #[test]
1355    fn empty_payload_still_rejects_zero_symbol_size() {
1356        let mut pipeline = EncodingPipeline::new(
1357            test_config(0, 16, 1.0),
1358            SymbolPool::new(PoolConfig::default()),
1359        );
1360        let err = pipeline
1361            .encode(ObjectId::new_for_test(102), &[])
1362            .next()
1363            .unwrap()
1364            .unwrap_err();
1365        assert!(matches!(err, EncodingError::InvalidConfig { .. }));
1366    }
1367
1368    #[test]
1369    fn empty_payload_still_rejects_invalid_repair_overhead() {
1370        let mut pipeline = EncodingPipeline::new(
1371            test_config(4, 16, f64::NAN),
1372            SymbolPool::new(PoolConfig::default()),
1373        );
1374        let err = pipeline
1375            .encode(ObjectId::new_for_test(103), &[])
1376            .next()
1377            .unwrap()
1378            .unwrap_err();
1379        assert!(matches!(err, EncodingError::InvalidConfig { .. }));
1380    }
1381
1382    #[test]
1383    fn empty_payload_still_rejects_pool_symbol_size_mismatch() {
1384        let mut pipeline = EncodingPipeline::new(test_config(4, 16, 1.0), pool_for(8, 1));
1385        let err = pipeline
1386            .encode(ObjectId::new_for_test(104), &[])
1387            .next()
1388            .unwrap()
1389            .unwrap_err();
1390        assert!(matches!(err, EncodingError::InvalidConfig { .. }));
1391    }
1392
1393    #[test]
1394    fn encoding_pipeline_stats_and_reset() {
1395        let mut pipeline = EncodingPipeline::new(
1396            test_config(4, 16, 1.0),
1397            SymbolPool::new(PoolConfig::default()),
1398        );
1399
1400        let stats = pipeline.stats();
1401        assert_eq!(stats.bytes_in, 0);
1402
1403        let _: Vec<_> = pipeline
1404            .encode(ObjectId::new_for_test(99), b"test")
1405            .collect::<Result<Vec<_>, _>>()
1406            .unwrap();
1407
1408        let stats = pipeline.stats();
1409        assert!(stats.source_symbols > 0);
1410
1411        pipeline.reset();
1412        let stats = pipeline.stats();
1413        assert_eq!(stats.bytes_in, 0);
1414        assert_eq!(stats.source_symbols, 0);
1415    }
1416}