1use alloc::vec;
2use alloc::vec::Vec;
3
4use super::advertisement::ResourceAdvertisement;
5use super::parts::{extract_metadata, map_hash};
6use super::proof::{build_proof_data, compute_expected_proof, compute_resource_hash};
7use super::types::*;
8use super::window::WindowState;
9use crate::buffer::types::Compressor;
10use crate::constants::*;
11
12pub struct ResourceReceiver {
17 pub status: ResourceStatus,
19 pub resource_hash: Vec<u8>,
21 pub random_hash: Vec<u8>,
23 pub original_hash: Vec<u8>,
25 pub flags: AdvFlags,
27 pub transfer_size: u64,
29 pub data_size: u64,
31 pub total_parts: usize,
33 parts: Vec<Option<Vec<u8>>>,
35 hashmap: Vec<Option<[u8; RESOURCE_MAPHASH_LEN]>>,
37 hashmap_height: usize,
39 pub waiting_for_hmu: bool,
41 pub received_count: usize,
43 pub outstanding_parts: usize,
45 consecutive_completed_height: isize,
47 sdu: usize,
49 link_rtt: f64,
51 pub retries_left: usize,
53 max_retries: usize,
55 pub rtt: Option<f64>,
57 part_timeout_factor: f64,
59 pub last_activity: f64,
61 pub req_sent: f64,
63 req_sent_bytes: usize,
65 req_resp: Option<f64>,
67 rtt_rxd_bytes: usize,
69 rtt_rxd_bytes_at_part_req: usize,
71 req_resp_rtt_rate: f64,
73 req_data_rtt_rate: f64,
75 pub eifr: Option<f64>,
77 previous_eifr: Option<f64>,
79 pub segment_index: u64,
81 pub total_segments: u64,
83 pub has_metadata: bool,
85 pub request_id: Option<Vec<u8>>,
87 pub window: WindowState,
89 pub advertisement_packet: Vec<u8>,
91 pub max_decompressed_size: usize,
93}
94
95impl ResourceReceiver {
96 pub fn from_advertisement(
98 adv_data: &[u8],
99 sdu: usize,
100 link_rtt: f64,
101 now: f64,
102 previous_window: Option<usize>,
103 previous_eifr: Option<f64>,
104 ) -> Result<Self, ResourceError> {
105 let adv = ResourceAdvertisement::unpack(adv_data)?;
106
107 if adv.resource_hash.len() != 32 {
109 return Err(ResourceError::InvalidAdvertisement);
110 }
111 if adv.random_hash.len() != RESOURCE_RANDOM_HASH_SIZE || adv.original_hash.len() != 32 {
112 return Err(ResourceError::InvalidAdvertisement);
113 }
114 if adv.transfer_size > RESOURCE_AUTO_COMPRESS_MAX_SIZE as u64
118 || (!adv.flags.split && adv.data_size > RESOURCE_AUTO_COMPRESS_MAX_SIZE as u64)
119 {
120 return Err(ResourceError::TooLarge);
121 }
122
123 let total_parts = usize::try_from(adv.num_parts).map_err(|_| ResourceError::TooLarge)?;
124 if total_parts == 0 {
125 return Err(ResourceError::InvalidAdvertisement);
126 }
127 let max_parts = RESOURCE_AUTO_COMPRESS_MAX_SIZE
128 .div_ceil(sdu.max(1))
129 .saturating_add(RESOURCE_COLLISION_GUARD_SIZE);
130 if total_parts > max_parts {
131 return Err(ResourceError::TooLarge);
132 }
133 let parts_vec: Vec<Option<Vec<u8>>> = vec![None; total_parts];
134 let mut hashmap_vec: Vec<Option<[u8; RESOURCE_MAPHASH_LEN]>> = vec![None; total_parts];
135
136 let initial_hashes = adv.hashmap.len() / RESOURCE_MAPHASH_LEN;
138 let mut hashmap_height = 0;
139 for (i, slot) in hashmap_vec.iter_mut().enumerate().take(initial_hashes) {
140 if i < total_parts {
141 let start = i * RESOURCE_MAPHASH_LEN;
142 let end = start + RESOURCE_MAPHASH_LEN;
143 let mut h = [0u8; RESOURCE_MAPHASH_LEN];
144 h.copy_from_slice(&adv.hashmap[start..end]);
145 *slot = Some(h);
146 hashmap_height += 1;
147 }
148 }
149
150 let mut window_state = WindowState::new();
151 if let Some(prev_w) = previous_window {
152 window_state.restore(prev_w);
153 }
154
155 Ok(ResourceReceiver {
156 status: ResourceStatus::None,
157 resource_hash: adv.resource_hash,
158 random_hash: adv.random_hash,
159 original_hash: adv.original_hash,
160 flags: adv.flags,
161 transfer_size: adv.transfer_size,
162 data_size: adv.data_size,
163 total_parts,
164 parts: parts_vec,
165 hashmap: hashmap_vec,
166 hashmap_height,
167 waiting_for_hmu: false,
168 received_count: 0,
169 outstanding_parts: 0,
170 consecutive_completed_height: -1,
171 sdu,
172 link_rtt,
173 retries_left: RESOURCE_MAX_RETRIES,
174 max_retries: RESOURCE_MAX_RETRIES,
175 rtt: None,
176 part_timeout_factor: RESOURCE_PART_TIMEOUT_FACTOR,
177 last_activity: now,
178 req_sent: 0.0,
179 req_sent_bytes: 0,
180 req_resp: None,
181 rtt_rxd_bytes: 0,
182 rtt_rxd_bytes_at_part_req: 0,
183 req_resp_rtt_rate: 0.0,
184 req_data_rtt_rate: 0.0,
185 eifr: None,
186 previous_eifr,
187 segment_index: adv.segment_index,
188 total_segments: adv.total_segments,
189 has_metadata: adv.flags.has_metadata,
190 request_id: adv.request_id,
191 window: window_state,
192 advertisement_packet: adv_data.to_vec(),
193 max_decompressed_size: RESOURCE_AUTO_COMPRESS_MAX_SIZE,
194 })
195 }
196
197 pub fn accept(&mut self, now: f64) -> Vec<ResourceAction> {
199 self.status = ResourceStatus::Transferring;
200 self.last_activity = now;
201 self.request_next(now)
202 }
203
204 pub fn reject(&mut self) -> Vec<ResourceAction> {
206 self.status = ResourceStatus::Rejected;
207 vec![ResourceAction::SendCancelReceiver(
208 self.resource_hash.clone(),
209 )]
210 }
211
212 pub fn cancel(&mut self) -> Vec<ResourceAction> {
214 if self.status < ResourceStatus::Complete {
215 self.status = ResourceStatus::Failed;
216 vec![ResourceAction::SendCancelReceiver(
217 self.resource_hash.clone(),
218 )]
219 } else {
220 Vec::new()
221 }
222 }
223
224 fn corrupt_actions(&mut self, error: ResourceError) -> Vec<ResourceAction> {
225 self.status = ResourceStatus::Corrupt;
226 vec![
227 ResourceAction::SendCancelReceiver(self.resource_hash.clone()),
228 ResourceAction::Failed(error),
229 ResourceAction::TeardownLink,
230 ]
231 }
232
233 pub fn receive_part(&mut self, part_data: &[u8], now: f64) -> Vec<ResourceAction> {
235 if self.status == ResourceStatus::Failed {
236 return vec![];
237 }
238
239 self.last_activity = now;
240 self.retries_left = self.max_retries;
241
242 if self.req_resp.is_none() {
244 self.req_resp = Some(now);
245 let rtt = now - self.req_sent;
246 self.part_timeout_factor = RESOURCE_PART_TIMEOUT_FACTOR_AFTER_RTT;
247
248 if self.rtt.is_none() {
249 self.rtt = Some(rtt);
250 } else if let Some(current_rtt) = self.rtt {
251 if rtt < current_rtt {
252 self.rtt = Some(f64::max(current_rtt - current_rtt * 0.05, rtt));
253 } else {
254 self.rtt = Some(f64::min(current_rtt + current_rtt * 0.05, rtt));
255 }
256 }
257
258 if rtt > 0.0 {
259 let req_resp_cost = part_data.len() + self.req_sent_bytes;
260 self.req_resp_rtt_rate = req_resp_cost as f64 / rtt;
261 self.window.update_req_resp_rate(self.req_resp_rtt_rate);
262 }
263 }
264
265 self.status = ResourceStatus::Transferring;
266
267 let part_hash = map_hash(part_data, &self.random_hash);
269
270 let consecutive_idx = if self.consecutive_completed_height >= 0 {
272 self.consecutive_completed_height as usize
273 } else {
274 0
275 };
276
277 let mut matched = false;
278 let search_end = core::cmp::min(consecutive_idx + self.window.window, self.total_parts);
279 for i in consecutive_idx..search_end {
280 if let Some(ref h) = self.hashmap[i] {
281 if *h == part_hash {
282 if self.parts[i].is_none() {
283 self.parts[i] = Some(part_data.to_vec());
284 self.rtt_rxd_bytes += part_data.len();
285 self.received_count += 1;
286 self.outstanding_parts = self.outstanding_parts.saturating_sub(1);
287
288 if i as isize == self.consecutive_completed_height + 1 {
290 self.consecutive_completed_height = i as isize;
291 }
292
293 let mut cp = (self.consecutive_completed_height + 1) as usize;
295 while cp < self.total_parts && self.parts[cp].is_some() {
296 self.consecutive_completed_height = cp as isize;
297 cp += 1;
298 }
299
300 matched = true;
301 }
302 break;
303 }
304 }
305 }
306
307 let mut actions = Vec::new();
308
309 if self.received_count == self.total_parts {
311 actions.push(ResourceAction::ProgressUpdate {
312 received: self.received_count,
313 total: self.total_parts,
314 });
315 return actions;
317 }
318
319 if matched {
320 actions.push(ResourceAction::ProgressUpdate {
321 received: self.received_count,
322 total: self.total_parts,
323 });
324 }
325
326 if self.outstanding_parts == 0 && self.received_count < self.total_parts {
328 self.window.on_window_complete();
330
331 if self.req_sent > 0.0 {
333 let rtt = now - self.req_sent;
334 let req_transferred = self.rtt_rxd_bytes - self.rtt_rxd_bytes_at_part_req;
335 if rtt > 0.0 {
336 self.req_data_rtt_rate = req_transferred as f64 / rtt;
337 self.rtt_rxd_bytes_at_part_req = self.rtt_rxd_bytes;
338 self.window.update_data_rate(self.req_data_rtt_rate);
339 }
340 }
341
342 let next_actions = self.request_next(now);
343 actions.extend(next_actions);
344 }
345
346 actions
347 }
348
349 pub fn handle_hashmap_update(&mut self, hmu_data: &[u8], now: f64) -> Vec<ResourceAction> {
353 if self.status == ResourceStatus::Failed || !self.waiting_for_hmu {
354 return vec![];
355 }
356
357 if hmu_data.len() < 32 || hmu_data[..32] != self.resource_hash {
358 return vec![];
359 }
360 if hmu_data.len() == 32 {
361 return self.cancel();
362 }
363
364 let payload = &hmu_data[32..];
365 let value = match crate::msgpack::unpack_exact(payload) {
366 Ok(v) => v,
367 Err(_) => return self.cancel(),
368 };
369
370 let arr = match value.as_array() {
371 Some(a) if a.len() == 2 => a,
372 _ => return self.cancel(),
373 };
374
375 let segment = match arr[0].as_uint() {
376 Some(s) => match usize::try_from(s) {
377 Ok(segment) => segment,
378 Err(_) => return self.cancel(),
379 },
380 None => return self.cancel(),
381 };
382
383 let hashmap_bytes = match arr[1].as_bin() {
384 Some(b) => b,
385 None => return self.cancel(),
386 };
387
388 if hashmap_bytes.is_empty() || !hashmap_bytes.len().is_multiple_of(RESOURCE_MAPHASH_LEN) {
389 return self.cancel();
390 }
391
392 let seg_len = RESOURCE_HASHMAP_MAX_LEN;
393 let segment_start = match segment.checked_mul(seg_len) {
394 Some(start) if start < self.total_parts => start,
395 _ => return self.cancel(),
396 };
397 let num_hashes = hashmap_bytes.len() / RESOURCE_MAPHASH_LEN;
398 if num_hashes > seg_len || segment_start.saturating_add(num_hashes) > self.total_parts {
399 return self.cancel();
400 }
401
402 self.last_activity = now;
403 self.retries_left = self.max_retries;
404 self.status = ResourceStatus::Transferring;
405
406 for i in 0..num_hashes {
408 let idx = segment_start + i;
409 let start = i * RESOURCE_MAPHASH_LEN;
410 let end = start + RESOURCE_MAPHASH_LEN;
411 if self.hashmap[idx].is_none() {
412 self.hashmap_height += 1;
413 }
414 let mut h = [0u8; RESOURCE_MAPHASH_LEN];
415 h.copy_from_slice(&hashmap_bytes[start..end]);
416 self.hashmap[idx] = Some(h);
417 }
418
419 self.waiting_for_hmu = false;
420 self.request_next(now)
421 }
422
423 pub fn request_next(&mut self, now: f64) -> Vec<ResourceAction> {
425 if self.status == ResourceStatus::Failed || self.waiting_for_hmu {
426 return vec![];
427 }
428
429 self.outstanding_parts = 0;
430 let mut hashmap_exhausted = RESOURCE_HASHMAP_IS_NOT_EXHAUSTED;
431 let mut requested_hashes = Vec::new();
432
433 let pn_start = (self.consecutive_completed_height + 1) as usize;
434 let search_end = core::cmp::min(pn_start + self.window.window, self.total_parts);
435 let mut i = 0;
436
437 for pn in pn_start..search_end {
438 if self.parts[pn].is_none() {
439 match self.hashmap[pn] {
440 Some(ref h) => {
441 requested_hashes.extend_from_slice(h);
442 self.outstanding_parts += 1;
443 i += 1;
444 }
445 None => {
446 hashmap_exhausted = RESOURCE_HASHMAP_IS_EXHAUSTED;
447 }
448 }
449 }
450 if i >= self.window.window || hashmap_exhausted == RESOURCE_HASHMAP_IS_EXHAUSTED {
451 break;
452 }
453 }
454
455 let mut request_data = Vec::new();
456 request_data.push(hashmap_exhausted);
457 if hashmap_exhausted == RESOURCE_HASHMAP_IS_EXHAUSTED {
458 if self.hashmap_height > 0 {
460 if let Some(ref last_hash) = self.hashmap[self.hashmap_height - 1] {
461 request_data.extend_from_slice(last_hash);
462 } else {
463 request_data.extend_from_slice(&[0u8; RESOURCE_MAPHASH_LEN]);
464 }
465 } else {
466 request_data.extend_from_slice(&[0u8; RESOURCE_MAPHASH_LEN]);
467 }
468 self.waiting_for_hmu = true;
469 }
470
471 request_data.extend_from_slice(&self.resource_hash);
472 request_data.extend_from_slice(&requested_hashes);
473
474 self.last_activity = now;
475 self.req_sent = now;
476 self.req_sent_bytes = request_data.len();
477 self.req_resp = None;
478
479 vec![ResourceAction::SendRequest(request_data)]
480 }
481
482 #[allow(clippy::type_complexity)]
484 pub fn assemble(
485 &mut self,
486 decrypt_fn: &dyn Fn(&[u8]) -> Result<Vec<u8>, ()>,
487 compressor: &dyn Compressor,
488 ) -> Vec<ResourceAction> {
489 if self.received_count != self.total_parts {
490 return vec![ResourceAction::Failed(ResourceError::InvalidState)];
491 }
492
493 self.status = ResourceStatus::Assembling;
494
495 let mut stream = Vec::new();
497 for part in &self.parts {
498 match part {
499 Some(data) => stream.extend_from_slice(data),
500 None => {
501 self.status = ResourceStatus::Failed;
502 return vec![ResourceAction::Failed(ResourceError::InvalidState)];
503 }
504 }
505 }
506
507 let decrypted = if self.flags.encrypted {
509 match decrypt_fn(&stream) {
510 Ok(d) => d,
511 Err(_) => {
512 self.status = ResourceStatus::Failed;
513 return vec![ResourceAction::Failed(ResourceError::DecryptionFailed)];
514 }
515 }
516 } else {
517 stream
518 };
519
520 if decrypted.len() < RESOURCE_RANDOM_HASH_SIZE {
522 return self.corrupt_actions(ResourceError::InvalidPart);
523 }
524 let data_after_random = &decrypted[RESOURCE_RANDOM_HASH_SIZE..];
525
526 let decompressed = if self.flags.compressed {
528 match compressor.decompress_bounded(data_after_random, self.max_decompressed_size) {
529 Ok(d) => d,
530 Err(crate::buffer::types::DecompressError::TooLarge) => {
531 return self.corrupt_actions(ResourceError::TooLarge);
532 }
533 Err(crate::buffer::types::DecompressError::InvalidData) => {
534 return self.corrupt_actions(ResourceError::DecompressionFailed);
535 }
536 }
537 } else {
538 data_after_random.to_vec()
539 };
540
541 let calculated_hash = compute_resource_hash(&decompressed, &self.random_hash);
543 if calculated_hash.as_slice() != self.resource_hash.as_slice() {
544 return self.corrupt_actions(ResourceError::HashMismatch);
545 }
546
547 let expected_proof = compute_expected_proof(&decompressed, &calculated_hash);
549 let proof_data = build_proof_data(&calculated_hash, &expected_proof);
550
551 let (data, metadata) = if self.has_metadata && self.segment_index == 1 {
553 match extract_metadata(&decompressed) {
554 Some((meta, rest)) => (rest, Some(meta)),
555 None => return self.corrupt_actions(ResourceError::InvalidPart),
556 }
557 } else {
558 (decompressed, None)
559 };
560
561 self.status = ResourceStatus::Complete;
562
563 vec![
564 ResourceAction::SendProof(proof_data),
565 ResourceAction::DataReceived { data, metadata },
566 ResourceAction::Completed,
567 ]
568 }
569
570 pub fn handle_cancel(&mut self) -> Vec<ResourceAction> {
572 if self.status < ResourceStatus::Complete {
573 self.status = ResourceStatus::Failed;
574 return vec![ResourceAction::Failed(ResourceError::Rejected)];
575 }
576 vec![]
577 }
578
579 #[allow(clippy::type_complexity)]
581 pub fn tick(
582 &mut self,
583 now: f64,
584 decrypt_fn: &dyn Fn(&[u8]) -> Result<Vec<u8>, ()>,
585 compressor: &dyn Compressor,
586 ) -> Vec<ResourceAction> {
587 if self.status >= ResourceStatus::Assembling {
588 return vec![];
589 }
590
591 if self.status == ResourceStatus::Transferring {
592 if self.received_count == self.total_parts {
594 return self.assemble(decrypt_fn, compressor);
595 }
596
597 let eifr = self.compute_eifr();
599 let retries_used = self.max_retries - self.retries_left;
600 let extra_wait = retries_used as f64 * RESOURCE_PER_RETRY_DELAY;
601 let expected_hmu_wait =
602 if eifr > 0.0 && (self.waiting_for_hmu || self.outstanding_parts == 0) {
603 (self.sdu as f64 * 8.0 * RESOURCE_HMU_WAIT_FACTOR) / eifr
604 } else {
605 0.0
606 };
607 let expected_tof = if self.outstanding_parts > 0 && eifr > 0.0 {
608 (self.outstanding_parts as f64 * self.sdu as f64 * 8.0) / eifr
609 } else if eifr > 0.0 {
610 (3.0 * self.sdu as f64) / eifr
611 } else {
612 10.0 };
614
615 let sleep_time = self.last_activity
616 + self.part_timeout_factor * expected_tof
617 + expected_hmu_wait
618 + RESOURCE_RETRY_GRACE_TIME
619 + extra_wait;
620
621 if now > sleep_time {
622 if self.retries_left > 0 {
623 self.window.on_timeout();
625 self.retries_left -= 1;
626 self.waiting_for_hmu = false;
627 return self.request_next(now);
628 } else {
629 self.status = ResourceStatus::Failed;
630 return vec![ResourceAction::Failed(ResourceError::MaxRetriesExceeded)];
631 }
632 }
633 }
634
635 vec![]
636 }
637
638 fn compute_eifr(&mut self) -> f64 {
640 let eifr = if self.req_data_rtt_rate > 0.0 {
641 self.req_data_rtt_rate * 8.0
642 } else if let Some(prev) = self.previous_eifr {
643 prev
644 } else {
645 let rtt = self.rtt.unwrap_or(self.link_rtt);
647 if rtt > 0.0 {
648 (self.sdu as f64 * 8.0) / rtt
649 } else {
650 10000.0
651 }
652 };
653 self.eifr = Some(eifr);
654 eifr
655 }
656
657 pub fn progress(&self) -> (usize, usize) {
659 (self.received_count, self.total_parts)
660 }
661
662 pub fn get_transfer_state(&self) -> (usize, Option<f64>) {
664 (self.window.window, self.eifr)
665 }
666}
667
668#[cfg(test)]
669mod tests {
670 use super::*;
671 use crate::buffer::types::{Compressor, DecompressError, NoopCompressor};
672 use crate::resource::advertisement::ResourceAdvertisement;
673 use crate::resource::sender::ResourceSender;
674
675 fn identity_encrypt(data: &[u8]) -> Vec<u8> {
676 data.to_vec()
677 }
678
679 fn identity_decrypt(data: &[u8]) -> Result<Vec<u8>, ()> {
680 Ok(data.to_vec())
681 }
682
683 struct ExpandingCompressor;
684
685 impl Compressor for ExpandingCompressor {
686 fn compress(&self, data: &[u8]) -> Option<Vec<u8>> {
687 Some(data[..data.len() / 2].to_vec())
688 }
689
690 fn decompress(&self, data: &[u8]) -> Option<Vec<u8>> {
691 self.decompress_bounded(data, usize::MAX).ok()
692 }
693
694 fn decompress_bounded(
695 &self,
696 data: &[u8],
697 max_output_size: usize,
698 ) -> Result<Vec<u8>, DecompressError> {
699 let mut out = data.to_vec();
700 out.extend_from_slice(data);
701 if out.len() > max_output_size {
702 return Err(DecompressError::TooLarge);
703 }
704 Ok(out)
705 }
706 }
707
708 fn base_timeout(receiver: &ResourceReceiver, eifr: f64) -> f64 {
709 let expected_tof = if receiver.outstanding_parts > 0 {
710 (receiver.outstanding_parts as f64 * receiver.sdu as f64 * 8.0) / eifr
711 } else {
712 (3.0 * receiver.sdu as f64) / eifr
713 };
714
715 receiver.last_activity
716 + receiver.part_timeout_factor * expected_tof
717 + RESOURCE_RETRY_GRACE_TIME
718 }
719
720 fn hmu_timeout(receiver: &ResourceReceiver, eifr: f64) -> f64 {
721 let expected_hmu_wait = (receiver.sdu as f64 * 8.0 * RESOURCE_HMU_WAIT_FACTOR) / eifr;
722 base_timeout(receiver, eifr) + expected_hmu_wait
723 }
724
725 fn make_sender_receiver() -> (ResourceSender, ResourceReceiver) {
726 let mut rng = rns_crypto::FixedRng::new(&[0x42; 64]);
727 let data = b"Hello, Resource Transfer!";
728
729 let sender = ResourceSender::new(
730 data,
731 None,
732 RESOURCE_SDU,
733 &identity_encrypt,
734 &NoopCompressor,
735 &mut rng,
736 1000.0,
737 false,
738 false,
739 None,
740 1,
741 1,
742 None,
743 0.5,
744 6.0,
745 )
746 .unwrap();
747
748 let adv_data = sender.get_advertisement(0);
749 let receiver =
750 ResourceReceiver::from_advertisement(&adv_data, RESOURCE_SDU, 0.5, 1000.0, None, None)
751 .unwrap();
752
753 (sender, receiver)
754 }
755
756 #[test]
757 fn test_from_advertisement() {
758 let (sender, receiver) = make_sender_receiver();
759 assert_eq!(receiver.total_parts, sender.total_parts());
760 assert_eq!(receiver.transfer_size, sender.transfer_size as u64);
761 assert_eq!(receiver.resource_hash, sender.resource_hash.to_vec());
762 assert!(!receiver.advertisement_packet.is_empty());
763 assert_eq!(
764 receiver.max_decompressed_size,
765 RESOURCE_AUTO_COMPRESS_MAX_SIZE
766 );
767 }
768
769 #[test]
770 fn test_from_advertisement_rejects_huge_part_count() {
771 let adv = ResourceAdvertisement {
772 transfer_size: 1024,
773 data_size: 1024,
774 num_parts: u64::MAX,
775 resource_hash: vec![0x11; 32],
776 random_hash: vec![0x22; RESOURCE_RANDOM_HASH_SIZE],
777 original_hash: vec![0x33; 32],
778 hashmap: vec![0x44; RESOURCE_MAPHASH_LEN],
779 flags: AdvFlags {
780 encrypted: false,
781 compressed: false,
782 split: false,
783 is_request: false,
784 is_response: false,
785 has_metadata: false,
786 },
787 segment_index: 1,
788 total_segments: 1,
789 request_id: None,
790 };
791
792 let err = match ResourceReceiver::from_advertisement(
793 &adv.pack(0),
794 RESOURCE_SDU,
795 0.5,
796 1000.0,
797 None,
798 None,
799 ) {
800 Ok(_) => panic!("huge advertisement should be rejected"),
801 Err(err) => err,
802 };
803
804 assert_eq!(err, ResourceError::TooLarge);
805 }
806
807 #[test]
808 fn test_from_advertisement_rejects_oversized_transfer() {
809 let adv = ResourceAdvertisement {
810 transfer_size: RESOURCE_AUTO_COMPRESS_MAX_SIZE as u64 + 1,
811 data_size: 1024,
812 num_parts: 1,
813 resource_hash: vec![0x11; 32],
814 random_hash: vec![0x22; RESOURCE_RANDOM_HASH_SIZE],
815 original_hash: vec![0x33; 32],
816 hashmap: vec![0x44; RESOURCE_MAPHASH_LEN],
817 flags: AdvFlags {
818 encrypted: false,
819 compressed: false,
820 split: false,
821 is_request: false,
822 is_response: false,
823 has_metadata: false,
824 },
825 segment_index: 1,
826 total_segments: 1,
827 request_id: None,
828 };
829
830 let err = match ResourceReceiver::from_advertisement(
831 &adv.pack(0),
832 RESOURCE_SDU,
833 0.5,
834 1000.0,
835 None,
836 None,
837 ) {
838 Ok(_) => panic!("oversized advertisement should be rejected"),
839 Err(err) => err,
840 };
841
842 assert_eq!(err, ResourceError::InvalidAdvertisement);
843 }
844
845 #[test]
846 fn test_accept() {
847 let (_, mut receiver) = make_sender_receiver();
848 let actions = receiver.accept(1000.0);
849 assert_eq!(receiver.status, ResourceStatus::Transferring);
850 assert!(!actions.is_empty());
851 assert!(actions
852 .iter()
853 .any(|a| matches!(a, ResourceAction::SendRequest(_))));
854 }
855
856 #[test]
857 fn test_reject() {
858 let (_, mut receiver) = make_sender_receiver();
859 let actions = receiver.reject();
860 assert_eq!(receiver.status, ResourceStatus::Rejected);
861 assert!(actions
862 .iter()
863 .any(|a| matches!(a, ResourceAction::SendCancelReceiver(_))));
864 }
865
866 #[test]
867 fn test_cancel_is_distinct_from_advertisement_rejection() {
868 let (_, mut receiver) = make_sender_receiver();
869 receiver.accept(1000.0);
870 let actions = receiver.cancel();
871 assert_eq!(receiver.status, ResourceStatus::Failed);
872 assert!(actions
873 .iter()
874 .any(|a| matches!(a, ResourceAction::SendCancelReceiver(_))));
875 assert!(receiver.cancel().is_empty());
876 }
877
878 fn hmu(resource_hash: &[u8], segment: u64, hashmap: Vec<u8>) -> Vec<u8> {
879 let payload = crate::msgpack::pack(&crate::msgpack::Value::Array(vec![
880 crate::msgpack::Value::UInt(segment),
881 crate::msgpack::Value::Bin(hashmap),
882 ]));
883 let mut data = resource_hash.to_vec();
884 data.extend_from_slice(&payload);
885 data
886 }
887
888 #[test]
889 fn test_unsolicited_hmu_is_ignored_without_refreshing_timeout_state() {
890 let (_, mut receiver) = make_sender_receiver();
891 receiver.status = ResourceStatus::Transferring;
892 receiver.waiting_for_hmu = false;
893 receiver.last_activity = 1000.0;
894 receiver.retries_left = 1;
895 let update = hmu(&receiver.resource_hash, 0, vec![0x99; RESOURCE_MAPHASH_LEN]);
896
897 let actions = receiver.handle_hashmap_update(&update, 2000.0);
898
899 assert!(actions.is_empty());
900 assert_eq!(receiver.status, ResourceStatus::Transferring);
901 assert!(!receiver.waiting_for_hmu);
902 assert_eq!(receiver.last_activity, 1000.0);
903 assert_eq!(receiver.retries_left, 1);
904 }
905
906 #[test]
907 fn test_empty_hmu_cancels_waiting_receiver() {
908 let (_, mut receiver) = make_sender_receiver();
909 receiver.status = ResourceStatus::Transferring;
910 receiver.waiting_for_hmu = true;
911 let update = hmu(&receiver.resource_hash, 1, Vec::new());
912
913 let actions = receiver.handle_hashmap_update(&update, 1001.0);
914
915 assert_eq!(receiver.status, ResourceStatus::Failed);
916 assert!(actions
917 .iter()
918 .any(|action| matches!(action, ResourceAction::SendCancelReceiver(_))));
919 assert!(!actions
920 .iter()
921 .any(|action| matches!(action, ResourceAction::SendRequest(_))));
922 }
923
924 #[test]
925 fn test_malformed_hmu_cancels_waiting_receiver() {
926 let (_, mut receiver) = make_sender_receiver();
927 receiver.status = ResourceStatus::Transferring;
928 receiver.waiting_for_hmu = true;
929 let mut malformed = receiver.resource_hash.clone();
930 malformed.extend_from_slice(&[0xc1]);
931
932 let actions = receiver.handle_hashmap_update(&malformed, 1001.0);
933
934 assert_eq!(receiver.status, ResourceStatus::Failed);
935 assert!(actions
936 .iter()
937 .any(|action| matches!(action, ResourceAction::SendCancelReceiver(_))));
938 }
939
940 #[test]
941 fn test_hmu_without_payload_cancels_waiting_receiver() {
942 let (_, mut receiver) = make_sender_receiver();
943 receiver.status = ResourceStatus::Transferring;
944 receiver.waiting_for_hmu = true;
945 let update = receiver.resource_hash.clone();
946
947 let actions = receiver.handle_hashmap_update(&update, 1001.0);
948
949 assert_eq!(receiver.status, ResourceStatus::Failed);
950 assert!(actions
951 .iter()
952 .any(|action| matches!(action, ResourceAction::SendCancelReceiver(_))));
953 }
954
955 #[test]
956 fn test_hmu_for_another_resource_is_ignored() {
957 let (_, mut receiver) = make_sender_receiver();
958 receiver.status = ResourceStatus::Transferring;
959 receiver.waiting_for_hmu = true;
960 receiver.last_activity = 1000.0;
961 let update = hmu(&[0xEE; 32], 0, vec![0x99; RESOURCE_MAPHASH_LEN]);
962
963 let actions = receiver.handle_hashmap_update(&update, 2000.0);
964
965 assert!(actions.is_empty());
966 assert!(receiver.waiting_for_hmu);
967 assert_eq!(receiver.last_activity, 1000.0);
968 }
969
970 #[test]
971 fn test_misaligned_hmu_hashmap_cancels_waiting_receiver() {
972 let (_, mut receiver) = make_sender_receiver();
973 receiver.status = ResourceStatus::Transferring;
974 receiver.waiting_for_hmu = true;
975 let update = hmu(
976 &receiver.resource_hash,
977 0,
978 vec![0x99; RESOURCE_MAPHASH_LEN + 1],
979 );
980
981 let actions = receiver.handle_hashmap_update(&update, 1001.0);
982
983 assert_eq!(receiver.status, ResourceStatus::Failed);
984 assert!(actions
985 .iter()
986 .any(|action| matches!(action, ResourceAction::SendCancelReceiver(_))));
987 }
988
989 #[test]
990 fn test_out_of_range_hmu_segment_cancels_waiting_receiver() {
991 let (_, mut receiver) = make_sender_receiver();
992 receiver.status = ResourceStatus::Transferring;
993 receiver.waiting_for_hmu = true;
994 let update = hmu(
995 &receiver.resource_hash,
996 u64::MAX,
997 vec![0x99; RESOURCE_MAPHASH_LEN],
998 );
999
1000 let actions = receiver.handle_hashmap_update(&update, 1001.0);
1001
1002 assert_eq!(receiver.status, ResourceStatus::Failed);
1003 assert!(actions
1004 .iter()
1005 .any(|action| matches!(action, ResourceAction::SendCancelReceiver(_))));
1006 }
1007
1008 #[test]
1009 fn test_tick_is_inert_once_receiver_starts_assembling() {
1010 for status in [
1011 ResourceStatus::Assembling,
1012 ResourceStatus::Complete,
1013 ResourceStatus::Failed,
1014 ResourceStatus::Corrupt,
1015 ResourceStatus::Rejected,
1016 ] {
1017 let (_, mut receiver) = make_sender_receiver();
1018 receiver.status = status;
1019 receiver.last_activity = 1.0;
1020 receiver.retries_left = 0;
1021
1022 let actions = receiver.tick(10_000.0, &identity_decrypt, &NoopCompressor);
1023
1024 assert!(actions.is_empty(), "tick emitted an action for {status:?}");
1025 assert_eq!(receiver.status, status);
1026 assert_eq!(receiver.last_activity, 1.0);
1027 }
1028 }
1029
1030 #[test]
1031 fn test_receive_part_stores() {
1032 let (mut sender, mut receiver) = make_sender_receiver();
1033 receiver.accept(1000.0);
1034
1035 let mut request = Vec::new();
1038 request.push(RESOURCE_HASHMAP_IS_NOT_EXHAUSTED);
1039 request.extend_from_slice(&sender.resource_hash);
1040 request.extend_from_slice(&sender.part_hashes[0]);
1041
1042 let send_actions = sender.handle_request(&request, 1001.0);
1043 let part_data = send_actions
1044 .iter()
1045 .find_map(|a| match a {
1046 ResourceAction::SendPart(d) => Some(d.clone()),
1047 _ => None,
1048 })
1049 .unwrap();
1050
1051 receiver.req_sent = 1000.5;
1053 let _actions = receiver.receive_part(&part_data, 1001.0);
1054 assert_eq!(receiver.received_count, 1);
1055 }
1056
1057 #[test]
1058 fn test_consecutive_completed_height() {
1059 let (sender, mut receiver) = make_sender_receiver();
1060 receiver.accept(1000.0);
1061
1062 if sender.total_parts() > 1 {
1064 assert_eq!(receiver.consecutive_completed_height, -1);
1066 }
1067 }
1068
1069 #[test]
1070 fn test_handle_cancel() {
1071 let (_, mut receiver) = make_sender_receiver();
1072 receiver.accept(1000.0);
1073 let _actions = receiver.handle_cancel();
1074 assert_eq!(receiver.status, ResourceStatus::Failed);
1075 }
1076
1077 #[test]
1078 fn test_assemble_compressed_resource_rejects_oversized_decompression() {
1079 let data = b"oversized!";
1080 let mut rng = rns_crypto::FixedRng::new(&[0x93; 64]);
1081
1082 let mut sender = ResourceSender::new(
1083 data,
1084 None,
1085 RESOURCE_SDU,
1086 &identity_encrypt,
1087 &ExpandingCompressor,
1088 &mut rng,
1089 1000.0,
1090 true,
1091 false,
1092 None,
1093 1,
1094 1,
1095 None,
1096 0.5,
1097 6.0,
1098 )
1099 .unwrap();
1100
1101 let adv = sender.get_advertisement(0);
1102 let mut receiver =
1103 ResourceReceiver::from_advertisement(&adv, RESOURCE_SDU, 0.5, 1000.0, None, None)
1104 .unwrap();
1105 receiver.max_decompressed_size = data.len() - 1;
1106
1107 let request_data = receiver
1108 .accept(1001.0)
1109 .into_iter()
1110 .find_map(|a| match a {
1111 ResourceAction::SendRequest(d) => Some(d),
1112 _ => None,
1113 })
1114 .unwrap();
1115
1116 let send_actions = sender.handle_request(&request_data, 1002.0);
1117 receiver.req_sent = 1001.0;
1118 for action in &send_actions {
1119 if let ResourceAction::SendPart(part_data) = action {
1120 receiver.receive_part(part_data, 1003.0);
1121 }
1122 }
1123
1124 let assemble_actions = receiver.assemble(&identity_decrypt, &ExpandingCompressor);
1125 assert_eq!(receiver.status, ResourceStatus::Corrupt);
1126 assert!(assemble_actions
1127 .iter()
1128 .any(|a| matches!(a, ResourceAction::SendCancelReceiver(_))));
1129 assert!(assemble_actions
1130 .iter()
1131 .any(|a| matches!(a, ResourceAction::Failed(ResourceError::TooLarge))));
1132 assert!(assemble_actions
1133 .iter()
1134 .any(|a| matches!(a, ResourceAction::TeardownLink)));
1135 }
1136
1137 #[test]
1138 fn test_full_transfer_small_data() {
1139 let data = b"small data";
1141 let mut rng = rns_crypto::FixedRng::new(&[0x77; 64]);
1142
1143 let mut sender = ResourceSender::new(
1144 data,
1145 None,
1146 RESOURCE_SDU,
1147 &identity_encrypt,
1148 &NoopCompressor,
1149 &mut rng,
1150 1000.0,
1151 false,
1152 false,
1153 None,
1154 1,
1155 1,
1156 None,
1157 0.5,
1158 6.0,
1159 )
1160 .unwrap();
1161
1162 let adv = sender.get_advertisement(0);
1163 let mut receiver =
1164 ResourceReceiver::from_advertisement(&adv, RESOURCE_SDU, 0.5, 1000.0, None, None)
1165 .unwrap();
1166
1167 let req_actions = receiver.accept(1001.0);
1169 assert_eq!(receiver.status, ResourceStatus::Transferring);
1170
1171 let request_data = req_actions
1173 .iter()
1174 .find_map(|a| match a {
1175 ResourceAction::SendRequest(d) => Some(d.clone()),
1176 _ => None,
1177 })
1178 .unwrap();
1179
1180 let send_actions = sender.handle_request(&request_data, 1002.0);
1182
1183 receiver.req_sent = 1001.0;
1185 for action in &send_actions {
1186 if let ResourceAction::SendPart(part_data) = action {
1187 receiver.receive_part(part_data, 1003.0);
1188 }
1189 }
1190
1191 assert_eq!(receiver.received_count, receiver.total_parts);
1192
1193 let assemble_actions = receiver.assemble(&identity_decrypt, &NoopCompressor);
1195
1196 let has_proof = assemble_actions
1198 .iter()
1199 .any(|a| matches!(a, ResourceAction::SendProof(_)));
1200 let has_data = assemble_actions
1201 .iter()
1202 .any(|a| matches!(a, ResourceAction::DataReceived { .. }));
1203 let has_complete = assemble_actions
1204 .iter()
1205 .any(|a| matches!(a, ResourceAction::Completed));
1206
1207 assert!(has_proof, "Should send proof");
1208 assert!(has_data, "Should return data");
1209 assert!(has_complete, "Should be completed");
1210
1211 let received_data = assemble_actions
1213 .iter()
1214 .find_map(|a| match a {
1215 ResourceAction::DataReceived { data, .. } => Some(data.clone()),
1216 _ => None,
1217 })
1218 .unwrap();
1219 assert_eq!(received_data, data);
1220
1221 let proof_data = assemble_actions
1223 .iter()
1224 .find_map(|a| match a {
1225 ResourceAction::SendProof(d) => Some(d.clone()),
1226 _ => None,
1227 })
1228 .unwrap();
1229
1230 let _proof_actions = sender.handle_proof(&proof_data, 1004.0);
1231 assert_eq!(sender.status, ResourceStatus::Complete);
1232 }
1233
1234 #[test]
1235 fn test_full_transfer_with_metadata() {
1236 let data = b"data with metadata";
1237 let metadata = b"some metadata";
1238 let mut rng = rns_crypto::FixedRng::new(&[0x88; 64]);
1239
1240 let mut sender = ResourceSender::new(
1241 data,
1242 Some(metadata),
1243 RESOURCE_SDU,
1244 &identity_encrypt,
1245 &NoopCompressor,
1246 &mut rng,
1247 1000.0,
1248 false,
1249 false,
1250 None,
1251 1,
1252 1,
1253 None,
1254 0.5,
1255 6.0,
1256 )
1257 .unwrap();
1258
1259 assert!(sender.flags.has_metadata);
1260
1261 let adv = sender.get_advertisement(0);
1262 let mut receiver =
1263 ResourceReceiver::from_advertisement(&adv, RESOURCE_SDU, 0.5, 1000.0, None, None)
1264 .unwrap();
1265
1266 assert!(receiver.has_metadata);
1267
1268 let req_actions = receiver.accept(1001.0);
1270 let request_data = req_actions
1271 .iter()
1272 .find_map(|a| match a {
1273 ResourceAction::SendRequest(d) => Some(d.clone()),
1274 _ => None,
1275 })
1276 .unwrap();
1277
1278 let send_actions = sender.handle_request(&request_data, 1002.0);
1279 receiver.req_sent = 1001.0;
1280 for action in &send_actions {
1281 if let ResourceAction::SendPart(part_data) = action {
1282 receiver.receive_part(part_data, 1003.0);
1283 }
1284 }
1285
1286 let assemble_actions = receiver.assemble(&identity_decrypt, &NoopCompressor);
1287
1288 let (recv_data, recv_meta) = assemble_actions
1289 .iter()
1290 .find_map(|a| match a {
1291 ResourceAction::DataReceived { data, metadata } => {
1292 Some((data.clone(), metadata.clone()))
1293 }
1294 _ => None,
1295 })
1296 .unwrap();
1297
1298 assert_eq!(recv_data, data);
1299 assert_eq!(recv_meta.unwrap(), metadata);
1300 }
1301
1302 #[test]
1303 fn test_previous_window_restore() {
1304 let (_, _receiver) = make_sender_receiver();
1305 let adv_data = {
1307 let mut rng = rns_crypto::FixedRng::new(&[0x42; 64]);
1308 let sender = ResourceSender::new(
1309 b"test",
1310 None,
1311 RESOURCE_SDU,
1312 &identity_encrypt,
1313 &NoopCompressor,
1314 &mut rng,
1315 1000.0,
1316 false,
1317 false,
1318 None,
1319 1,
1320 1,
1321 None,
1322 0.5,
1323 6.0,
1324 )
1325 .unwrap();
1326 sender.get_advertisement(0)
1327 };
1328
1329 let receiver = ResourceReceiver::from_advertisement(
1330 &adv_data,
1331 RESOURCE_SDU,
1332 0.5,
1333 1000.0,
1334 Some(8),
1335 Some(50000.0),
1336 )
1337 .unwrap();
1338 assert_eq!(receiver.window.window, 8);
1339 }
1340
1341 #[test]
1342 fn test_tick_timeout_retry() {
1343 let (_, mut receiver) = make_sender_receiver();
1344 receiver.accept(1000.0);
1345 receiver.rtt = Some(0.1);
1346
1347 let actions = receiver.tick(9999.0, &identity_decrypt, &NoopCompressor);
1349 assert!(!actions.is_empty() || receiver.retries_left < RESOURCE_MAX_RETRIES);
1351 }
1352
1353 #[test]
1354 fn test_tick_waiting_for_hmu_gets_extra_timeout() {
1355 let (_, mut receiver) = make_sender_receiver();
1356 receiver.accept(1000.0);
1357 receiver.waiting_for_hmu = true;
1358 receiver.outstanding_parts = 0;
1359 let eifr = 10_000.0;
1360 receiver.previous_eifr = Some(eifr);
1361 receiver.last_activity = 1000.0;
1362
1363 let old_timeout = base_timeout(&receiver, eifr);
1364 let now = old_timeout + 0.01;
1365
1366 let actions = receiver.tick(now, &identity_decrypt, &NoopCompressor);
1367
1368 assert!(actions.is_empty(), "receiver should keep waiting for HMU");
1369 assert_eq!(receiver.retries_left, RESOURCE_MAX_RETRIES);
1370 assert_eq!(receiver.status, ResourceStatus::Transferring);
1371 }
1372
1373 #[test]
1374 fn test_tick_zero_outstanding_parts_gets_extra_timeout_without_hmu_flag() {
1375 let (_, mut receiver) = make_sender_receiver();
1376 receiver.accept(1000.0);
1377 receiver.waiting_for_hmu = false;
1378 receiver.outstanding_parts = 0;
1379 let eifr = 10_000.0;
1380 receiver.previous_eifr = Some(eifr);
1381 receiver.last_activity = 1000.0;
1382
1383 let old_timeout = base_timeout(&receiver, eifr);
1384 let now = old_timeout + 0.01;
1385
1386 let actions = receiver.tick(now, &identity_decrypt, &NoopCompressor);
1387
1388 assert!(
1389 actions.is_empty(),
1390 "receiver should keep waiting for follow-up hashmap data"
1391 );
1392 assert_eq!(receiver.retries_left, RESOURCE_MAX_RETRIES);
1393 assert_eq!(receiver.status, ResourceStatus::Transferring);
1394 }
1395
1396 #[test]
1397 fn test_tick_waiting_for_hmu_retries_after_extended_timeout() {
1398 let (_, mut receiver) = make_sender_receiver();
1399 receiver.accept(1000.0);
1400 receiver.waiting_for_hmu = true;
1401 receiver.outstanding_parts = 0;
1402 let eifr = 10_000.0;
1403 receiver.previous_eifr = Some(eifr);
1404 receiver.last_activity = 1000.0;
1405
1406 let now = hmu_timeout(&receiver, eifr) + 0.01;
1407 let _actions = receiver.tick(now, &identity_decrypt, &NoopCompressor);
1408
1409 assert_eq!(receiver.retries_left, RESOURCE_MAX_RETRIES - 1);
1410 assert!(!receiver.waiting_for_hmu);
1411 }
1412
1413 #[test]
1414 fn test_tick_inflight_parts_do_not_get_hmu_timeout_extension() {
1415 let (_, mut receiver) = make_sender_receiver();
1416 receiver.accept(1000.0);
1417 receiver.waiting_for_hmu = false;
1418 receiver.outstanding_parts = 2;
1419 let eifr = 10_000.0;
1420 receiver.previous_eifr = Some(eifr);
1421 receiver.last_activity = 1000.0;
1422
1423 let now = base_timeout(&receiver, eifr) + 0.01;
1424 let _actions = receiver.tick(now, &identity_decrypt, &NoopCompressor);
1425
1426 assert_eq!(receiver.retries_left, RESOURCE_MAX_RETRIES - 1);
1427 }
1428
1429 #[test]
1430 fn test_tick_max_retries_exceeded() {
1431 let (_, mut receiver) = make_sender_receiver();
1432 receiver.accept(1000.0);
1433 receiver.retries_left = 0;
1434 receiver.rtt = Some(0.001);
1435 receiver.eifr = Some(100000.0);
1436
1437 let _actions = receiver.tick(9999.0, &identity_decrypt, &NoopCompressor);
1438 assert_eq!(receiver.status, ResourceStatus::Failed);
1439 }
1440}