Skip to main content

rns_core/resource/
receiver.rs

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
12/// Resource receiver state machine.
13///
14/// Unpacks advertisements, requests parts, receives parts, assembles data.
15/// Returns `Vec<ResourceAction>` — no I/O, no callbacks.
16pub struct ResourceReceiver {
17    /// Current status
18    pub status: ResourceStatus,
19    /// Resource hash (from advertisement, 32 bytes)
20    pub resource_hash: Vec<u8>,
21    /// Random hash (from advertisement)
22    pub random_hash: Vec<u8>,
23    /// Original hash
24    pub original_hash: Vec<u8>,
25    /// Flags
26    pub flags: AdvFlags,
27    /// Transfer size (encrypted)
28    pub transfer_size: u64,
29    /// Total uncompressed data size
30    pub data_size: u64,
31    /// Total parts
32    pub total_parts: usize,
33    /// Received parts data (None = not yet received)
34    parts: Vec<Option<Vec<u8>>>,
35    /// Hashmap: part index -> map_hash (None if not yet known)
36    hashmap: Vec<Option<[u8; RESOURCE_MAPHASH_LEN]>>,
37    /// Number of hashmap entries populated
38    hashmap_height: usize,
39    /// Whether we're waiting for a hashmap update
40    pub waiting_for_hmu: bool,
41    /// Number of parts received
42    pub received_count: usize,
43    /// Outstanding parts in current window request
44    pub outstanding_parts: usize,
45    /// Consecutive completed height (-1 means none)
46    consecutive_completed_height: isize,
47    /// SDU size
48    sdu: usize,
49    /// Link RTT estimate (from link establishment)
50    link_rtt: f64,
51    /// Retries left
52    pub retries_left: usize,
53    /// Max retries
54    max_retries: usize,
55    /// RTT estimate
56    pub rtt: Option<f64>,
57    /// Part timeout factor
58    part_timeout_factor: f64,
59    /// Last activity timestamp
60    pub last_activity: f64,
61    /// Request sent timestamp
62    pub req_sent: f64,
63    /// Request sent bytes
64    req_sent_bytes: usize,
65    /// Request response timestamp
66    req_resp: Option<f64>,
67    /// RTT received bytes
68    rtt_rxd_bytes: usize,
69    /// RTT received bytes at part request
70    rtt_rxd_bytes_at_part_req: usize,
71    /// Request response RTT rate
72    req_resp_rtt_rate: f64,
73    /// Request data RTT rate
74    req_data_rtt_rate: f64,
75    /// EIFR
76    pub eifr: Option<f64>,
77    /// Previous EIFR from prior transfer
78    previous_eifr: Option<f64>,
79    /// Segment index
80    pub segment_index: u64,
81    /// Total segments
82    pub total_segments: u64,
83    /// Has metadata
84    pub has_metadata: bool,
85    /// Request ID
86    pub request_id: Option<Vec<u8>>,
87    /// Window state
88    pub window: WindowState,
89    /// Original advertisement bytes for this incoming transfer.
90    pub advertisement_packet: Vec<u8>,
91    /// Maximum allowed decompressed size for compressed resource payloads.
92    pub max_decompressed_size: usize,
93}
94
95impl ResourceReceiver {
96    /// Create a receiver from an advertisement packet.
97    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        // Validate resource_hash is 32 bytes
108        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        // Split advertisements carry the logical size of the complete
115        // Resource in data_size, while each receiver only assembles one
116        // bounded segment. Aggregate receive policy is enforced by rns-net.
117        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        // Populate initial hashmap from advertisement
137        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    /// Accept the advertised resource. Begins transfer.
198    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    /// Reject the advertised resource.
205    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    /// Cancel an accepted incoming transfer.
213    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    /// Receive a part. Matches by map hash and stores it.
234    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        // Update RTT on first part of window
243        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        // Compute map hash for this part
268        let part_hash = map_hash(part_data, &self.random_hash);
269
270        // Align with request_next(): the completed height itself is not part
271        // of the outstanding window, which starts at the following index.
272        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                        // Walk forward to extend consecutive height
285                        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        // Check if all parts received
301        if self.received_count == self.total_parts {
302            actions.push(ResourceAction::ProgressUpdate {
303                received: self.received_count,
304                total: self.total_parts,
305            });
306            // Assembly will be triggered by caller
307            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        // Request next window when outstanding is 0
318        if self.outstanding_parts == 0 && self.received_count < self.total_parts {
319            // Window complete — grow
320            self.window.on_window_complete();
321
322            // Update data rate
323            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    /// Handle a hashmap update packet.
341    ///
342    /// HMU format: [resource_hash: 32 bytes][msgpack([segment, hashmap])]
343    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        // Populate hashmap slots
398        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    /// Build and return request for next window of parts.
415    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            // Append last known map hash
450            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    /// Assemble received parts, decrypt, decompress, verify hash.
474    #[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        // Join all parts
487        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        // Decrypt
499        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        // Strip random hash prefix
512        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        // Decompress
518        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        // Verify hash
533        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        // Compute proof before metadata extraction (proof uses full decompressed data)
539        let expected_proof = compute_expected_proof(&decompressed, &calculated_hash);
540        let proof_data = build_proof_data(&calculated_hash, &expected_proof);
541
542        // Extract metadata if present
543        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    /// Handle cancel from sender (RESOURCE_ICL).
562    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    /// Periodic tick. Checks for timeouts.
571    #[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            // Check if all parts received — trigger assembly
584            if self.received_count == self.total_parts {
585                return self.assemble(decrypt_fn, compressor);
586            }
587
588            // Compute timeout
589            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 // fallback
604            };
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                    // Timeout — shrink window, retry
615                    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    /// Compute EIFR (expected inflight rate) and update self.eifr.
630    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            // Fallback: use link_rtt as establishment cost estimate
637            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    /// Get current progress as (received, total).
649    pub fn progress(&self) -> (usize, usize) {
650        (self.received_count, self.total_parts)
651    }
652
653    /// Get window and EIFR for passing to next transfer.
654    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        // Get part data from sender (we use identity encryption so parts ARE the raw data)
1027        // Request first part
1028        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        // Give it to receiver
1043        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        // Simulate receiving parts in order for a multi-part resource
1054        if sender.total_parts() > 1 {
1055            // This only applies to multi-part transfers
1056            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        // End-to-end: sender creates, receiver accepts, parts flow, assembly completes
1218        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        // Accept
1246        let req_actions = receiver.accept(1001.0);
1247        assert_eq!(receiver.status, ResourceStatus::Transferring);
1248
1249        // Get request data
1250        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        // Sender handles request
1259        let send_actions = sender.handle_request(&request_data, 1002.0);
1260
1261        // Feed all parts to receiver
1262        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        // Assemble
1272        let assemble_actions = receiver.assemble(&identity_decrypt, &NoopCompressor);
1273
1274        // Check for proof and data
1275        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        // Verify data matches
1290        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        // Verify proof validates on sender side
1300        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        // Transfer all parts
1347        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        // Create with previous window
1384        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        // Way past timeout
1426        let actions = receiver.tick(9999.0, &identity_decrypt, &NoopCompressor);
1427        // Should have retried (window decreased, request_next called)
1428        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}