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