1use 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
15pub(crate) const MAX_SOURCE_BLOCKS: usize = u8::MAX as usize + 1;
17
18#[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#[derive(Debug, thiserror::Error)]
27pub enum EncodingError {
28 #[error("data too large: {size} bytes exceeds limit {limit}")]
30 DataTooLarge {
31 size: usize,
33 limit: usize,
35 },
36 #[error("symbol pool exhausted")]
38 PoolExhausted,
39 #[error("invalid configuration: {reason}")]
41 InvalidConfig {
42 reason: String,
44 },
45 #[error("encoding failed: {details}")]
47 ComputationFailed {
48 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#[derive(Debug, Clone)]
77pub struct EncodedSymbol {
78 symbol: Symbol,
79}
80
81impl EncodedSymbol {
82 #[must_use]
84 #[inline]
85 pub const fn new(symbol: Symbol) -> Self {
86 Self { symbol }
87 }
88
89 #[must_use]
91 #[inline]
92 pub const fn symbol(&self) -> &Symbol {
93 &self.symbol
94 }
95
96 #[must_use]
98 #[inline]
99 pub fn into_symbol(self) -> Symbol {
100 self.symbol
101 }
102
103 #[must_use]
105 #[inline]
106 pub const fn id(&self) -> SymbolId {
107 self.symbol.id()
108 }
109
110 #[must_use]
112 #[inline]
113 pub const fn kind(&self) -> SymbolKind {
114 self.symbol.kind()
115 }
116}
117
118#[derive(Debug, Clone, Copy, Default)]
120pub struct EncodingStats {
121 pub bytes_in: usize,
123 pub blocks: usize,
125 pub source_symbols: usize,
127 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#[derive(Debug)]
144pub struct EncodingPipeline {
145 config: EncodingConfig,
146 pool: SymbolPool,
147 stats: EncodingStats,
148}
149
150impl EncodingPipeline {
151 #[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 #[must_use]
164 #[inline]
165 pub const fn stats(&self) -> EncodingStats {
166 self.stats
167 }
168
169 #[inline]
171 pub fn reset(&mut self) {
172 self.stats = EncodingStats::default();
173 }
174
175 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 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 #[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 #[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 #[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
494pub 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
630pub 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 #[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; 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 assert_eq!(compute_repair_count(10, 1.0), 0);
1300 assert_eq!(compute_repair_count(10, 1.5), 5);
1302 assert_eq!(compute_repair_count(10, 2.0), 10);
1304 assert_eq!(compute_repair_count(1, 1.5), 1);
1306 }
1307
1308 #[test]
1309 fn compute_repair_count_large_k_does_not_truncate() {
1310 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 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}