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 search_start = (self.consecutive_completed_height + 1) as usize;
273 let mut matched = false;
274 let search_end = core::cmp::min(search_start + self.window.window, self.total_parts);
275 for i in search_start..search_end {
276 if let Some(ref h) = self.hashmap[i] {
277 if *h == part_hash {
278 if self.parts[i].is_none() {
279 self.parts[i] = Some(part_data.to_vec());
280 self.rtt_rxd_bytes += part_data.len();
281 self.received_count += 1;
282 self.outstanding_parts = self.outstanding_parts.saturating_sub(1);
283
284 let mut cp = (self.consecutive_completed_height + 1) as usize;
286 while cp < self.total_parts && self.parts[cp].is_some() {
287 self.consecutive_completed_height = cp as isize;
288 cp += 1;
289 }
290
291 matched = true;
292 }
293 break;
294 }
295 }
296 }
297
298 let mut actions = Vec::new();
299
300 if self.received_count == self.total_parts {
302 actions.push(ResourceAction::ProgressUpdate {
303 received: self.received_count,
304 total: self.total_parts,
305 });
306 return actions;
308 }
309
310 if matched {
311 actions.push(ResourceAction::ProgressUpdate {
312 received: self.received_count,
313 total: self.total_parts,
314 });
315 }
316
317 if self.outstanding_parts == 0 && self.received_count < self.total_parts {
319 self.window.on_window_complete();
321
322 if self.req_sent > 0.0 {
324 let rtt = now - self.req_sent;
325 let req_transferred = self.rtt_rxd_bytes - self.rtt_rxd_bytes_at_part_req;
326 if rtt > 0.0 {
327 self.req_data_rtt_rate = req_transferred as f64 / rtt;
328 self.rtt_rxd_bytes_at_part_req = self.rtt_rxd_bytes;
329 self.window.update_data_rate(self.req_data_rtt_rate);
330 }
331 }
332
333 let next_actions = self.request_next(now);
334 actions.extend(next_actions);
335 }
336
337 actions
338 }
339
340 pub fn handle_hashmap_update(&mut self, hmu_data: &[u8], now: f64) -> Vec<ResourceAction> {
344 if self.status == ResourceStatus::Failed || !self.waiting_for_hmu {
345 return vec![];
346 }
347
348 if hmu_data.len() < 32 || hmu_data[..32] != self.resource_hash {
349 return vec![];
350 }
351 if hmu_data.len() == 32 {
352 return self.cancel();
353 }
354
355 let payload = &hmu_data[32..];
356 let value = match crate::msgpack::unpack_exact(payload) {
357 Ok(v) => v,
358 Err(_) => return self.cancel(),
359 };
360
361 let arr = match value.as_array() {
362 Some(a) if a.len() == 2 => a,
363 _ => return self.cancel(),
364 };
365
366 let segment = match arr[0].as_uint() {
367 Some(s) => match usize::try_from(s) {
368 Ok(segment) => segment,
369 Err(_) => return self.cancel(),
370 },
371 None => return self.cancel(),
372 };
373
374 let hashmap_bytes = match arr[1].as_bin() {
375 Some(b) => b,
376 None => return self.cancel(),
377 };
378
379 if hashmap_bytes.is_empty() || !hashmap_bytes.len().is_multiple_of(RESOURCE_MAPHASH_LEN) {
380 return self.cancel();
381 }
382
383 let seg_len = RESOURCE_HASHMAP_MAX_LEN;
384 let segment_start = match segment.checked_mul(seg_len) {
385 Some(start) if start < self.total_parts => start,
386 _ => return self.cancel(),
387 };
388 let num_hashes = hashmap_bytes.len() / RESOURCE_MAPHASH_LEN;
389 if num_hashes > seg_len || segment_start.saturating_add(num_hashes) > self.total_parts {
390 return self.cancel();
391 }
392
393 self.last_activity = now;
394 self.retries_left = self.max_retries;
395 self.status = ResourceStatus::Transferring;
396
397 for i in 0..num_hashes {
399 let idx = segment_start + i;
400 let start = i * RESOURCE_MAPHASH_LEN;
401 let end = start + RESOURCE_MAPHASH_LEN;
402 if self.hashmap[idx].is_none() {
403 self.hashmap_height += 1;
404 }
405 let mut h = [0u8; RESOURCE_MAPHASH_LEN];
406 h.copy_from_slice(&hashmap_bytes[start..end]);
407 self.hashmap[idx] = Some(h);
408 }
409
410 self.waiting_for_hmu = false;
411 self.request_next(now)
412 }
413
414 pub fn request_next(&mut self, now: f64) -> Vec<ResourceAction> {
416 if self.status == ResourceStatus::Failed || self.waiting_for_hmu {
417 return vec![];
418 }
419
420 self.outstanding_parts = 0;
421 let mut hashmap_exhausted = RESOURCE_HASHMAP_IS_NOT_EXHAUSTED;
422 let mut requested_hashes = Vec::new();
423
424 let pn_start = (self.consecutive_completed_height + 1) as usize;
425 let search_end = core::cmp::min(pn_start + self.window.window, self.total_parts);
426 let mut i = 0;
427
428 for pn in pn_start..search_end {
429 if self.parts[pn].is_none() {
430 match self.hashmap[pn] {
431 Some(ref h) => {
432 requested_hashes.extend_from_slice(h);
433 self.outstanding_parts += 1;
434 i += 1;
435 }
436 None => {
437 hashmap_exhausted = RESOURCE_HASHMAP_IS_EXHAUSTED;
438 }
439 }
440 }
441 if i >= self.window.window || hashmap_exhausted == RESOURCE_HASHMAP_IS_EXHAUSTED {
442 break;
443 }
444 }
445
446 let mut request_data = Vec::new();
447 request_data.push(hashmap_exhausted);
448 if hashmap_exhausted == RESOURCE_HASHMAP_IS_EXHAUSTED {
449 if self.hashmap_height > 0 {
451 if let Some(ref last_hash) = self.hashmap[self.hashmap_height - 1] {
452 request_data.extend_from_slice(last_hash);
453 } else {
454 request_data.extend_from_slice(&[0u8; RESOURCE_MAPHASH_LEN]);
455 }
456 } else {
457 request_data.extend_from_slice(&[0u8; RESOURCE_MAPHASH_LEN]);
458 }
459 self.waiting_for_hmu = true;
460 }
461
462 request_data.extend_from_slice(&self.resource_hash);
463 request_data.extend_from_slice(&requested_hashes);
464
465 self.last_activity = now;
466 self.req_sent = now;
467 self.req_sent_bytes = request_data.len();
468 self.req_resp = None;
469
470 vec![ResourceAction::SendRequest(request_data)]
471 }
472
473 #[allow(clippy::type_complexity)]
475 pub fn assemble(
476 &mut self,
477 decrypt_fn: &dyn Fn(&[u8]) -> Result<Vec<u8>, ()>,
478 compressor: &dyn Compressor,
479 ) -> Vec<ResourceAction> {
480 if self.received_count != self.total_parts {
481 return vec![ResourceAction::Failed(ResourceError::InvalidState)];
482 }
483
484 self.status = ResourceStatus::Assembling;
485
486 let mut stream = Vec::new();
488 for part in &self.parts {
489 match part {
490 Some(data) => stream.extend_from_slice(data),
491 None => {
492 self.status = ResourceStatus::Failed;
493 return vec![ResourceAction::Failed(ResourceError::InvalidState)];
494 }
495 }
496 }
497
498 let decrypted = if self.flags.encrypted {
500 match decrypt_fn(&stream) {
501 Ok(d) => d,
502 Err(_) => {
503 self.status = ResourceStatus::Failed;
504 return vec![ResourceAction::Failed(ResourceError::DecryptionFailed)];
505 }
506 }
507 } else {
508 stream
509 };
510
511 if decrypted.len() < RESOURCE_RANDOM_HASH_SIZE {
513 return self.corrupt_actions(ResourceError::InvalidPart);
514 }
515 let data_after_random = &decrypted[RESOURCE_RANDOM_HASH_SIZE..];
516
517 let decompressed = if self.flags.compressed {
519 match compressor.decompress_bounded(data_after_random, self.max_decompressed_size) {
520 Ok(d) => d,
521 Err(crate::buffer::types::DecompressError::TooLarge) => {
522 return self.corrupt_actions(ResourceError::TooLarge);
523 }
524 Err(crate::buffer::types::DecompressError::InvalidData) => {
525 return self.corrupt_actions(ResourceError::DecompressionFailed);
526 }
527 }
528 } else {
529 data_after_random.to_vec()
530 };
531
532 let calculated_hash = compute_resource_hash(&decompressed, &self.random_hash);
534 if calculated_hash.as_slice() != self.resource_hash.as_slice() {
535 return self.corrupt_actions(ResourceError::HashMismatch);
536 }
537
538 let expected_proof = compute_expected_proof(&decompressed, &calculated_hash);
540 let proof_data = build_proof_data(&calculated_hash, &expected_proof);
541
542 let (data, metadata) = if self.has_metadata && self.segment_index == 1 {
544 match extract_metadata(&decompressed) {
545 Some((meta, rest)) => (rest, Some(meta)),
546 None => return self.corrupt_actions(ResourceError::InvalidPart),
547 }
548 } else {
549 (decompressed, None)
550 };
551
552 self.status = ResourceStatus::Complete;
553
554 vec![
555 ResourceAction::SendProof(proof_data),
556 ResourceAction::DataReceived { data, metadata },
557 ResourceAction::Completed,
558 ]
559 }
560
561 pub fn handle_cancel(&mut self) -> Vec<ResourceAction> {
563 if self.status < ResourceStatus::Complete {
564 self.status = ResourceStatus::Failed;
565 return vec![ResourceAction::Failed(ResourceError::Rejected)];
566 }
567 vec![]
568 }
569
570 #[allow(clippy::type_complexity)]
572 pub fn tick(
573 &mut self,
574 now: f64,
575 decrypt_fn: &dyn Fn(&[u8]) -> Result<Vec<u8>, ()>,
576 compressor: &dyn Compressor,
577 ) -> Vec<ResourceAction> {
578 if self.status >= ResourceStatus::Assembling {
579 return vec![];
580 }
581
582 if self.status == ResourceStatus::Transferring {
583 if self.received_count == self.total_parts {
585 return self.assemble(decrypt_fn, compressor);
586 }
587
588 let eifr = self.compute_eifr();
590 let retries_used = self.max_retries - self.retries_left;
591 let extra_wait = retries_used as f64 * RESOURCE_PER_RETRY_DELAY;
592 let expected_hmu_wait =
593 if eifr > 0.0 && (self.waiting_for_hmu || self.outstanding_parts == 0) {
594 (self.sdu as f64 * 8.0 * RESOURCE_HMU_WAIT_FACTOR) / eifr
595 } else {
596 0.0
597 };
598 let expected_tof = if self.outstanding_parts > 0 && eifr > 0.0 {
599 (self.outstanding_parts as f64 * self.sdu as f64 * 8.0) / eifr
600 } else if eifr > 0.0 {
601 (3.0 * self.sdu as f64) / eifr
602 } else {
603 10.0 };
605
606 let sleep_time = self.last_activity
607 + self.part_timeout_factor * expected_tof
608 + expected_hmu_wait
609 + RESOURCE_RETRY_GRACE_TIME
610 + extra_wait;
611
612 if now > sleep_time {
613 if self.retries_left > 0 {
614 self.window.on_timeout();
616 self.retries_left -= 1;
617 self.waiting_for_hmu = false;
618 return self.request_next(now);
619 } else {
620 self.status = ResourceStatus::Failed;
621 return vec![ResourceAction::Failed(ResourceError::MaxRetriesExceeded)];
622 }
623 }
624 }
625
626 vec![]
627 }
628
629 fn compute_eifr(&mut self) -> f64 {
631 let eifr = if self.req_data_rtt_rate > 0.0 {
632 self.req_data_rtt_rate * 8.0
633 } else if let Some(prev) = self.previous_eifr {
634 prev
635 } else {
636 let rtt = self.rtt.unwrap_or(self.link_rtt);
638 if rtt > 0.0 {
639 (self.sdu as f64 * 8.0) / rtt
640 } else {
641 10000.0
642 }
643 };
644 self.eifr = Some(eifr);
645 eifr
646 }
647
648 pub fn progress(&self) -> (usize, usize) {
650 (self.received_count, self.total_parts)
651 }
652
653 pub fn get_transfer_state(&self) -> (usize, Option<f64>) {
655 (self.window.window, self.eifr)
656 }
657}
658
659#[cfg(test)]
660mod tests {
661 use super::*;
662 use crate::buffer::types::{Compressor, DecompressError, NoopCompressor};
663 use crate::resource::advertisement::ResourceAdvertisement;
664 use crate::resource::sender::ResourceSender;
665
666 fn identity_encrypt(data: &[u8]) -> Vec<u8> {
667 data.to_vec()
668 }
669
670 fn identity_decrypt(data: &[u8]) -> Result<Vec<u8>, ()> {
671 Ok(data.to_vec())
672 }
673
674 struct ExpandingCompressor;
675
676 impl Compressor for ExpandingCompressor {
677 fn compress(&self, data: &[u8]) -> Option<Vec<u8>> {
678 Some(data[..data.len() / 2].to_vec())
679 }
680
681 fn decompress(&self, data: &[u8]) -> Option<Vec<u8>> {
682 self.decompress_bounded(data, usize::MAX).ok()
683 }
684
685 fn decompress_bounded(
686 &self,
687 data: &[u8],
688 max_output_size: usize,
689 ) -> Result<Vec<u8>, DecompressError> {
690 let mut out = data.to_vec();
691 out.extend_from_slice(data);
692 if out.len() > max_output_size {
693 return Err(DecompressError::TooLarge);
694 }
695 Ok(out)
696 }
697 }
698
699 fn base_timeout(receiver: &ResourceReceiver, eifr: f64) -> f64 {
700 let expected_tof = if receiver.outstanding_parts > 0 {
701 (receiver.outstanding_parts as f64 * receiver.sdu as f64 * 8.0) / eifr
702 } else {
703 (3.0 * receiver.sdu as f64) / eifr
704 };
705
706 receiver.last_activity
707 + receiver.part_timeout_factor * expected_tof
708 + RESOURCE_RETRY_GRACE_TIME
709 }
710
711 fn hmu_timeout(receiver: &ResourceReceiver, eifr: f64) -> f64 {
712 let expected_hmu_wait = (receiver.sdu as f64 * 8.0 * RESOURCE_HMU_WAIT_FACTOR) / eifr;
713 base_timeout(receiver, eifr) + expected_hmu_wait
714 }
715
716 fn make_sender_receiver() -> (ResourceSender, ResourceReceiver) {
717 let mut rng = rns_crypto::FixedRng::new(&[0x42; 64]);
718 let data = b"Hello, Resource Transfer!";
719
720 let sender = ResourceSender::new(
721 data,
722 None,
723 RESOURCE_SDU,
724 &identity_encrypt,
725 &NoopCompressor,
726 &mut rng,
727 1000.0,
728 false,
729 false,
730 None,
731 1,
732 1,
733 None,
734 0.5,
735 6.0,
736 )
737 .unwrap();
738
739 let adv_data = sender.get_advertisement(0);
740 let receiver =
741 ResourceReceiver::from_advertisement(&adv_data, RESOURCE_SDU, 0.5, 1000.0, None, None)
742 .unwrap();
743
744 (sender, receiver)
745 }
746
747 #[test]
748 fn test_from_advertisement() {
749 let (sender, receiver) = make_sender_receiver();
750 assert_eq!(receiver.total_parts, sender.total_parts());
751 assert_eq!(receiver.transfer_size, sender.transfer_size as u64);
752 assert_eq!(receiver.resource_hash, sender.resource_hash.to_vec());
753 assert!(!receiver.advertisement_packet.is_empty());
754 assert_eq!(
755 receiver.max_decompressed_size,
756 RESOURCE_AUTO_COMPRESS_MAX_SIZE
757 );
758 }
759
760 #[test]
761 fn test_from_advertisement_rejects_huge_part_count() {
762 let adv = ResourceAdvertisement {
763 transfer_size: 1024,
764 data_size: 1024,
765 num_parts: u64::MAX,
766 resource_hash: vec![0x11; 32],
767 random_hash: vec![0x22; RESOURCE_RANDOM_HASH_SIZE],
768 original_hash: vec![0x33; 32],
769 hashmap: vec![0x44; RESOURCE_MAPHASH_LEN],
770 flags: AdvFlags {
771 encrypted: false,
772 compressed: false,
773 split: false,
774 is_request: false,
775 is_response: false,
776 has_metadata: false,
777 },
778 segment_index: 1,
779 total_segments: 1,
780 request_id: None,
781 };
782
783 let err = match ResourceReceiver::from_advertisement(
784 &adv.pack(0),
785 RESOURCE_SDU,
786 0.5,
787 1000.0,
788 None,
789 None,
790 ) {
791 Ok(_) => panic!("huge advertisement should be rejected"),
792 Err(err) => err,
793 };
794
795 assert_eq!(err, ResourceError::TooLarge);
796 }
797
798 #[test]
799 fn test_from_advertisement_rejects_oversized_transfer() {
800 let adv = ResourceAdvertisement {
801 transfer_size: RESOURCE_AUTO_COMPRESS_MAX_SIZE as u64 + 1,
802 data_size: 1024,
803 num_parts: 1,
804 resource_hash: vec![0x11; 32],
805 random_hash: vec![0x22; RESOURCE_RANDOM_HASH_SIZE],
806 original_hash: vec![0x33; 32],
807 hashmap: vec![0x44; RESOURCE_MAPHASH_LEN],
808 flags: AdvFlags {
809 encrypted: false,
810 compressed: false,
811 split: false,
812 is_request: false,
813 is_response: false,
814 has_metadata: false,
815 },
816 segment_index: 1,
817 total_segments: 1,
818 request_id: None,
819 };
820
821 let err = match ResourceReceiver::from_advertisement(
822 &adv.pack(0),
823 RESOURCE_SDU,
824 0.5,
825 1000.0,
826 None,
827 None,
828 ) {
829 Ok(_) => panic!("oversized advertisement should be rejected"),
830 Err(err) => err,
831 };
832
833 assert_eq!(err, ResourceError::InvalidAdvertisement);
834 }
835
836 #[test]
837 fn test_accept() {
838 let (_, mut receiver) = make_sender_receiver();
839 let actions = receiver.accept(1000.0);
840 assert_eq!(receiver.status, ResourceStatus::Transferring);
841 assert!(!actions.is_empty());
842 assert!(actions
843 .iter()
844 .any(|a| matches!(a, ResourceAction::SendRequest(_))));
845 }
846
847 #[test]
848 fn test_reject() {
849 let (_, mut receiver) = make_sender_receiver();
850 let actions = receiver.reject();
851 assert_eq!(receiver.status, ResourceStatus::Rejected);
852 assert!(actions
853 .iter()
854 .any(|a| matches!(a, ResourceAction::SendCancelReceiver(_))));
855 }
856
857 #[test]
858 fn test_cancel_is_distinct_from_advertisement_rejection() {
859 let (_, mut receiver) = make_sender_receiver();
860 receiver.accept(1000.0);
861 let actions = receiver.cancel();
862 assert_eq!(receiver.status, ResourceStatus::Failed);
863 assert!(actions
864 .iter()
865 .any(|a| matches!(a, ResourceAction::SendCancelReceiver(_))));
866 assert!(receiver.cancel().is_empty());
867 }
868
869 fn hmu(resource_hash: &[u8], segment: u64, hashmap: Vec<u8>) -> Vec<u8> {
870 let payload = crate::msgpack::pack(&crate::msgpack::Value::Array(vec![
871 crate::msgpack::Value::UInt(segment),
872 crate::msgpack::Value::Bin(hashmap),
873 ]));
874 let mut data = resource_hash.to_vec();
875 data.extend_from_slice(&payload);
876 data
877 }
878
879 #[test]
880 fn test_unsolicited_hmu_is_ignored_without_refreshing_timeout_state() {
881 let (_, mut receiver) = make_sender_receiver();
882 receiver.status = ResourceStatus::Transferring;
883 receiver.waiting_for_hmu = false;
884 receiver.last_activity = 1000.0;
885 receiver.retries_left = 1;
886 let update = hmu(&receiver.resource_hash, 0, vec![0x99; RESOURCE_MAPHASH_LEN]);
887
888 let actions = receiver.handle_hashmap_update(&update, 2000.0);
889
890 assert!(actions.is_empty());
891 assert_eq!(receiver.status, ResourceStatus::Transferring);
892 assert!(!receiver.waiting_for_hmu);
893 assert_eq!(receiver.last_activity, 1000.0);
894 assert_eq!(receiver.retries_left, 1);
895 }
896
897 #[test]
898 fn test_empty_hmu_cancels_waiting_receiver() {
899 let (_, mut receiver) = make_sender_receiver();
900 receiver.status = ResourceStatus::Transferring;
901 receiver.waiting_for_hmu = true;
902 let update = hmu(&receiver.resource_hash, 1, Vec::new());
903
904 let actions = receiver.handle_hashmap_update(&update, 1001.0);
905
906 assert_eq!(receiver.status, ResourceStatus::Failed);
907 assert!(actions
908 .iter()
909 .any(|action| matches!(action, ResourceAction::SendCancelReceiver(_))));
910 assert!(!actions
911 .iter()
912 .any(|action| matches!(action, ResourceAction::SendRequest(_))));
913 }
914
915 #[test]
916 fn test_malformed_hmu_cancels_waiting_receiver() {
917 let (_, mut receiver) = make_sender_receiver();
918 receiver.status = ResourceStatus::Transferring;
919 receiver.waiting_for_hmu = true;
920 let mut malformed = receiver.resource_hash.clone();
921 malformed.extend_from_slice(&[0xc1]);
922
923 let actions = receiver.handle_hashmap_update(&malformed, 1001.0);
924
925 assert_eq!(receiver.status, ResourceStatus::Failed);
926 assert!(actions
927 .iter()
928 .any(|action| matches!(action, ResourceAction::SendCancelReceiver(_))));
929 }
930
931 #[test]
932 fn test_hmu_without_payload_cancels_waiting_receiver() {
933 let (_, mut receiver) = make_sender_receiver();
934 receiver.status = ResourceStatus::Transferring;
935 receiver.waiting_for_hmu = true;
936 let update = receiver.resource_hash.clone();
937
938 let actions = receiver.handle_hashmap_update(&update, 1001.0);
939
940 assert_eq!(receiver.status, ResourceStatus::Failed);
941 assert!(actions
942 .iter()
943 .any(|action| matches!(action, ResourceAction::SendCancelReceiver(_))));
944 }
945
946 #[test]
947 fn test_hmu_for_another_resource_is_ignored() {
948 let (_, mut receiver) = make_sender_receiver();
949 receiver.status = ResourceStatus::Transferring;
950 receiver.waiting_for_hmu = true;
951 receiver.last_activity = 1000.0;
952 let update = hmu(&[0xEE; 32], 0, vec![0x99; RESOURCE_MAPHASH_LEN]);
953
954 let actions = receiver.handle_hashmap_update(&update, 2000.0);
955
956 assert!(actions.is_empty());
957 assert!(receiver.waiting_for_hmu);
958 assert_eq!(receiver.last_activity, 1000.0);
959 }
960
961 #[test]
962 fn test_misaligned_hmu_hashmap_cancels_waiting_receiver() {
963 let (_, mut receiver) = make_sender_receiver();
964 receiver.status = ResourceStatus::Transferring;
965 receiver.waiting_for_hmu = true;
966 let update = hmu(
967 &receiver.resource_hash,
968 0,
969 vec![0x99; RESOURCE_MAPHASH_LEN + 1],
970 );
971
972 let actions = receiver.handle_hashmap_update(&update, 1001.0);
973
974 assert_eq!(receiver.status, ResourceStatus::Failed);
975 assert!(actions
976 .iter()
977 .any(|action| matches!(action, ResourceAction::SendCancelReceiver(_))));
978 }
979
980 #[test]
981 fn test_out_of_range_hmu_segment_cancels_waiting_receiver() {
982 let (_, mut receiver) = make_sender_receiver();
983 receiver.status = ResourceStatus::Transferring;
984 receiver.waiting_for_hmu = true;
985 let update = hmu(
986 &receiver.resource_hash,
987 u64::MAX,
988 vec![0x99; RESOURCE_MAPHASH_LEN],
989 );
990
991 let actions = receiver.handle_hashmap_update(&update, 1001.0);
992
993 assert_eq!(receiver.status, ResourceStatus::Failed);
994 assert!(actions
995 .iter()
996 .any(|action| matches!(action, ResourceAction::SendCancelReceiver(_))));
997 }
998
999 #[test]
1000 fn test_tick_is_inert_once_receiver_starts_assembling() {
1001 for status in [
1002 ResourceStatus::Assembling,
1003 ResourceStatus::Complete,
1004 ResourceStatus::Failed,
1005 ResourceStatus::Corrupt,
1006 ResourceStatus::Rejected,
1007 ] {
1008 let (_, mut receiver) = make_sender_receiver();
1009 receiver.status = status;
1010 receiver.last_activity = 1.0;
1011 receiver.retries_left = 0;
1012
1013 let actions = receiver.tick(10_000.0, &identity_decrypt, &NoopCompressor);
1014
1015 assert!(actions.is_empty(), "tick emitted an action for {status:?}");
1016 assert_eq!(receiver.status, status);
1017 assert_eq!(receiver.last_activity, 1.0);
1018 }
1019 }
1020
1021 #[test]
1022 fn test_receive_part_stores() {
1023 let (mut sender, mut receiver) = make_sender_receiver();
1024 receiver.accept(1000.0);
1025
1026 let mut request = Vec::new();
1029 request.push(RESOURCE_HASHMAP_IS_NOT_EXHAUSTED);
1030 request.extend_from_slice(&sender.resource_hash);
1031 request.extend_from_slice(&sender.part_hashes[0]);
1032
1033 let send_actions = sender.handle_request(&request, 1001.0);
1034 let part_data = send_actions
1035 .iter()
1036 .find_map(|a| match a {
1037 ResourceAction::SendPart(d) => Some(d.clone()),
1038 _ => None,
1039 })
1040 .unwrap();
1041
1042 receiver.req_sent = 1000.5;
1044 let _actions = receiver.receive_part(&part_data, 1001.0);
1045 assert_eq!(receiver.received_count, 1);
1046 }
1047
1048 #[test]
1049 fn test_consecutive_completed_height() {
1050 let (sender, mut receiver) = make_sender_receiver();
1051 receiver.accept(1000.0);
1052
1053 if sender.total_parts() > 1 {
1055 assert_eq!(receiver.consecutive_completed_height, -1);
1057 }
1058 }
1059
1060 #[test]
1061 fn receive_part_search_window_matches_request_after_completed_height() {
1062 let mut rng = rns_crypto::FixedRng::new(&[0x5a; 64]);
1063 let data: Vec<u8> = (0..48).collect();
1064 let mut sender = ResourceSender::new(
1065 &data,
1066 None,
1067 8,
1068 &identity_encrypt,
1069 &NoopCompressor,
1070 &mut rng,
1071 1000.0,
1072 false,
1073 false,
1074 None,
1075 1,
1076 1,
1077 None,
1078 0.5,
1079 6.0,
1080 )
1081 .unwrap();
1082 assert!(sender.total_parts() > 5);
1083
1084 let advertisement = sender.get_advertisement(0);
1085 let mut receiver =
1086 ResourceReceiver::from_advertisement(&advertisement, 8, 0.5, 1000.0, None, None)
1087 .unwrap();
1088 let first_request = receiver
1089 .accept(1000.0)
1090 .into_iter()
1091 .find_map(|action| match action {
1092 ResourceAction::SendRequest(data) => Some(data),
1093 _ => None,
1094 })
1095 .unwrap();
1096 let first_parts = sender.handle_request(&first_request, 1001.0);
1097 let first_part = first_parts
1098 .iter()
1099 .find_map(|action| match action {
1100 ResourceAction::SendPart(data)
1101 if map_hash(data, &receiver.random_hash) == sender.part_hashes[0] =>
1102 {
1103 Some(data.clone())
1104 }
1105 _ => None,
1106 })
1107 .unwrap();
1108 receiver.receive_part(&first_part, 1002.0);
1109 assert_eq!(receiver.consecutive_completed_height, 0);
1110
1111 let next_request = receiver
1112 .request_next(1003.0)
1113 .into_iter()
1114 .find_map(|action| match action {
1115 ResourceAction::SendRequest(data) => Some(data),
1116 _ => None,
1117 })
1118 .unwrap();
1119 let next_parts = sender.handle_request(&next_request, 1004.0);
1120 let window_tail_index = receiver.window.window;
1121 let window_tail = next_parts
1122 .iter()
1123 .find_map(|action| match action {
1124 ResourceAction::SendPart(data)
1125 if map_hash(data, &receiver.random_hash)
1126 == sender.part_hashes[window_tail_index] =>
1127 {
1128 Some(data.clone())
1129 }
1130 _ => None,
1131 })
1132 .expect("request must include the inclusive tail of the next window");
1133
1134 let actions = receiver.receive_part(&window_tail, 1005.0);
1135 assert!(receiver.parts[window_tail_index].is_some());
1136 assert_eq!(receiver.received_count, 2);
1137 assert_eq!(receiver.consecutive_completed_height, 0);
1138 assert!(actions.iter().any(|action| matches!(
1139 action,
1140 ResourceAction::ProgressUpdate {
1141 received: 2,
1142 total,
1143 } if *total == receiver.total_parts
1144 )));
1145 }
1146
1147 #[test]
1148 fn test_handle_cancel() {
1149 let (_, mut receiver) = make_sender_receiver();
1150 receiver.accept(1000.0);
1151 let _actions = receiver.handle_cancel();
1152 assert_eq!(receiver.status, ResourceStatus::Failed);
1153 }
1154
1155 #[test]
1156 fn test_assemble_compressed_resource_rejects_oversized_decompression() {
1157 let data = b"oversized!";
1158 let mut rng = rns_crypto::FixedRng::new(&[0x93; 64]);
1159
1160 let mut sender = ResourceSender::new(
1161 data,
1162 None,
1163 RESOURCE_SDU,
1164 &identity_encrypt,
1165 &ExpandingCompressor,
1166 &mut rng,
1167 1000.0,
1168 true,
1169 false,
1170 None,
1171 1,
1172 1,
1173 None,
1174 0.5,
1175 6.0,
1176 )
1177 .unwrap();
1178
1179 let adv = sender.get_advertisement(0);
1180 let mut receiver =
1181 ResourceReceiver::from_advertisement(&adv, RESOURCE_SDU, 0.5, 1000.0, None, None)
1182 .unwrap();
1183 receiver.max_decompressed_size = data.len() - 1;
1184
1185 let request_data = receiver
1186 .accept(1001.0)
1187 .into_iter()
1188 .find_map(|a| match a {
1189 ResourceAction::SendRequest(d) => Some(d),
1190 _ => None,
1191 })
1192 .unwrap();
1193
1194 let send_actions = sender.handle_request(&request_data, 1002.0);
1195 receiver.req_sent = 1001.0;
1196 for action in &send_actions {
1197 if let ResourceAction::SendPart(part_data) = action {
1198 receiver.receive_part(part_data, 1003.0);
1199 }
1200 }
1201
1202 let assemble_actions = receiver.assemble(&identity_decrypt, &ExpandingCompressor);
1203 assert_eq!(receiver.status, ResourceStatus::Corrupt);
1204 assert!(assemble_actions
1205 .iter()
1206 .any(|a| matches!(a, ResourceAction::SendCancelReceiver(_))));
1207 assert!(assemble_actions
1208 .iter()
1209 .any(|a| matches!(a, ResourceAction::Failed(ResourceError::TooLarge))));
1210 assert!(assemble_actions
1211 .iter()
1212 .any(|a| matches!(a, ResourceAction::TeardownLink)));
1213 }
1214
1215 #[test]
1216 fn test_full_transfer_small_data() {
1217 let data = b"small data";
1219 let mut rng = rns_crypto::FixedRng::new(&[0x77; 64]);
1220
1221 let mut sender = ResourceSender::new(
1222 data,
1223 None,
1224 RESOURCE_SDU,
1225 &identity_encrypt,
1226 &NoopCompressor,
1227 &mut rng,
1228 1000.0,
1229 false,
1230 false,
1231 None,
1232 1,
1233 1,
1234 None,
1235 0.5,
1236 6.0,
1237 )
1238 .unwrap();
1239
1240 let adv = sender.get_advertisement(0);
1241 let mut receiver =
1242 ResourceReceiver::from_advertisement(&adv, RESOURCE_SDU, 0.5, 1000.0, None, None)
1243 .unwrap();
1244
1245 let req_actions = receiver.accept(1001.0);
1247 assert_eq!(receiver.status, ResourceStatus::Transferring);
1248
1249 let request_data = req_actions
1251 .iter()
1252 .find_map(|a| match a {
1253 ResourceAction::SendRequest(d) => Some(d.clone()),
1254 _ => None,
1255 })
1256 .unwrap();
1257
1258 let send_actions = sender.handle_request(&request_data, 1002.0);
1260
1261 receiver.req_sent = 1001.0;
1263 for action in &send_actions {
1264 if let ResourceAction::SendPart(part_data) = action {
1265 receiver.receive_part(part_data, 1003.0);
1266 }
1267 }
1268
1269 assert_eq!(receiver.received_count, receiver.total_parts);
1270
1271 let assemble_actions = receiver.assemble(&identity_decrypt, &NoopCompressor);
1273
1274 let has_proof = assemble_actions
1276 .iter()
1277 .any(|a| matches!(a, ResourceAction::SendProof(_)));
1278 let has_data = assemble_actions
1279 .iter()
1280 .any(|a| matches!(a, ResourceAction::DataReceived { .. }));
1281 let has_complete = assemble_actions
1282 .iter()
1283 .any(|a| matches!(a, ResourceAction::Completed));
1284
1285 assert!(has_proof, "Should send proof");
1286 assert!(has_data, "Should return data");
1287 assert!(has_complete, "Should be completed");
1288
1289 let received_data = assemble_actions
1291 .iter()
1292 .find_map(|a| match a {
1293 ResourceAction::DataReceived { data, .. } => Some(data.clone()),
1294 _ => None,
1295 })
1296 .unwrap();
1297 assert_eq!(received_data, data);
1298
1299 let proof_data = assemble_actions
1301 .iter()
1302 .find_map(|a| match a {
1303 ResourceAction::SendProof(d) => Some(d.clone()),
1304 _ => None,
1305 })
1306 .unwrap();
1307
1308 let _proof_actions = sender.handle_proof(&proof_data, 1004.0);
1309 assert_eq!(sender.status, ResourceStatus::Complete);
1310 }
1311
1312 #[test]
1313 fn test_full_transfer_with_metadata() {
1314 let data = b"data with metadata";
1315 let metadata = b"some metadata";
1316 let mut rng = rns_crypto::FixedRng::new(&[0x88; 64]);
1317
1318 let mut sender = ResourceSender::new(
1319 data,
1320 Some(metadata),
1321 RESOURCE_SDU,
1322 &identity_encrypt,
1323 &NoopCompressor,
1324 &mut rng,
1325 1000.0,
1326 false,
1327 false,
1328 None,
1329 1,
1330 1,
1331 None,
1332 0.5,
1333 6.0,
1334 )
1335 .unwrap();
1336
1337 assert!(sender.flags.has_metadata);
1338
1339 let adv = sender.get_advertisement(0);
1340 let mut receiver =
1341 ResourceReceiver::from_advertisement(&adv, RESOURCE_SDU, 0.5, 1000.0, None, None)
1342 .unwrap();
1343
1344 assert!(receiver.has_metadata);
1345
1346 let req_actions = receiver.accept(1001.0);
1348 let request_data = req_actions
1349 .iter()
1350 .find_map(|a| match a {
1351 ResourceAction::SendRequest(d) => Some(d.clone()),
1352 _ => None,
1353 })
1354 .unwrap();
1355
1356 let send_actions = sender.handle_request(&request_data, 1002.0);
1357 receiver.req_sent = 1001.0;
1358 for action in &send_actions {
1359 if let ResourceAction::SendPart(part_data) = action {
1360 receiver.receive_part(part_data, 1003.0);
1361 }
1362 }
1363
1364 let assemble_actions = receiver.assemble(&identity_decrypt, &NoopCompressor);
1365
1366 let (recv_data, recv_meta) = assemble_actions
1367 .iter()
1368 .find_map(|a| match a {
1369 ResourceAction::DataReceived { data, metadata } => {
1370 Some((data.clone(), metadata.clone()))
1371 }
1372 _ => None,
1373 })
1374 .unwrap();
1375
1376 assert_eq!(recv_data, data);
1377 assert_eq!(recv_meta.unwrap(), metadata);
1378 }
1379
1380 #[test]
1381 fn test_previous_window_restore() {
1382 let (_, _receiver) = make_sender_receiver();
1383 let adv_data = {
1385 let mut rng = rns_crypto::FixedRng::new(&[0x42; 64]);
1386 let sender = ResourceSender::new(
1387 b"test",
1388 None,
1389 RESOURCE_SDU,
1390 &identity_encrypt,
1391 &NoopCompressor,
1392 &mut rng,
1393 1000.0,
1394 false,
1395 false,
1396 None,
1397 1,
1398 1,
1399 None,
1400 0.5,
1401 6.0,
1402 )
1403 .unwrap();
1404 sender.get_advertisement(0)
1405 };
1406
1407 let receiver = ResourceReceiver::from_advertisement(
1408 &adv_data,
1409 RESOURCE_SDU,
1410 0.5,
1411 1000.0,
1412 Some(8),
1413 Some(50000.0),
1414 )
1415 .unwrap();
1416 assert_eq!(receiver.window.window, 8);
1417 }
1418
1419 #[test]
1420 fn test_tick_timeout_retry() {
1421 let (_, mut receiver) = make_sender_receiver();
1422 receiver.accept(1000.0);
1423 receiver.rtt = Some(0.1);
1424
1425 let actions = receiver.tick(9999.0, &identity_decrypt, &NoopCompressor);
1427 assert!(!actions.is_empty() || receiver.retries_left < RESOURCE_MAX_RETRIES);
1429 }
1430
1431 #[test]
1432 fn test_tick_waiting_for_hmu_gets_extra_timeout() {
1433 let (_, mut receiver) = make_sender_receiver();
1434 receiver.accept(1000.0);
1435 receiver.waiting_for_hmu = true;
1436 receiver.outstanding_parts = 0;
1437 let eifr = 10_000.0;
1438 receiver.previous_eifr = Some(eifr);
1439 receiver.last_activity = 1000.0;
1440
1441 let old_timeout = base_timeout(&receiver, eifr);
1442 let now = old_timeout + 0.01;
1443
1444 let actions = receiver.tick(now, &identity_decrypt, &NoopCompressor);
1445
1446 assert!(actions.is_empty(), "receiver should keep waiting for HMU");
1447 assert_eq!(receiver.retries_left, RESOURCE_MAX_RETRIES);
1448 assert_eq!(receiver.status, ResourceStatus::Transferring);
1449 }
1450
1451 #[test]
1452 fn test_tick_zero_outstanding_parts_gets_extra_timeout_without_hmu_flag() {
1453 let (_, mut receiver) = make_sender_receiver();
1454 receiver.accept(1000.0);
1455 receiver.waiting_for_hmu = false;
1456 receiver.outstanding_parts = 0;
1457 let eifr = 10_000.0;
1458 receiver.previous_eifr = Some(eifr);
1459 receiver.last_activity = 1000.0;
1460
1461 let old_timeout = base_timeout(&receiver, eifr);
1462 let now = old_timeout + 0.01;
1463
1464 let actions = receiver.tick(now, &identity_decrypt, &NoopCompressor);
1465
1466 assert!(
1467 actions.is_empty(),
1468 "receiver should keep waiting for follow-up hashmap data"
1469 );
1470 assert_eq!(receiver.retries_left, RESOURCE_MAX_RETRIES);
1471 assert_eq!(receiver.status, ResourceStatus::Transferring);
1472 }
1473
1474 #[test]
1475 fn test_tick_waiting_for_hmu_retries_after_extended_timeout() {
1476 let (_, mut receiver) = make_sender_receiver();
1477 receiver.accept(1000.0);
1478 receiver.waiting_for_hmu = true;
1479 receiver.outstanding_parts = 0;
1480 let eifr = 10_000.0;
1481 receiver.previous_eifr = Some(eifr);
1482 receiver.last_activity = 1000.0;
1483
1484 let now = hmu_timeout(&receiver, eifr) + 0.01;
1485 let _actions = receiver.tick(now, &identity_decrypt, &NoopCompressor);
1486
1487 assert_eq!(receiver.retries_left, RESOURCE_MAX_RETRIES - 1);
1488 assert!(!receiver.waiting_for_hmu);
1489 }
1490
1491 #[test]
1492 fn test_tick_inflight_parts_do_not_get_hmu_timeout_extension() {
1493 let (_, mut receiver) = make_sender_receiver();
1494 receiver.accept(1000.0);
1495 receiver.waiting_for_hmu = false;
1496 receiver.outstanding_parts = 2;
1497 let eifr = 10_000.0;
1498 receiver.previous_eifr = Some(eifr);
1499 receiver.last_activity = 1000.0;
1500
1501 let now = base_timeout(&receiver, eifr) + 0.01;
1502 let _actions = receiver.tick(now, &identity_decrypt, &NoopCompressor);
1503
1504 assert_eq!(receiver.retries_left, RESOURCE_MAX_RETRIES - 1);
1505 }
1506
1507 #[test]
1508 fn test_tick_max_retries_exceeded() {
1509 let (_, mut receiver) = make_sender_receiver();
1510 receiver.accept(1000.0);
1511 receiver.retries_left = 0;
1512 receiver.rtt = Some(0.001);
1513 receiver.eifr = Some(100000.0);
1514
1515 let _actions = receiver.tick(9999.0, &identity_decrypt, &NoopCompressor);
1516 assert_eq!(receiver.status, ResourceStatus::Failed);
1517 }
1518}