Skip to content
File

Blob: firmware/vendor/sctp-proto/src/association/association_test.rs

rust4357 lines
1use crate::chunk::chunk_i_forward_tsn::ChunkIForwardTsnStream;
2use crate::chunk::{Chunk, chunk_init::ChunkInit};
3use crate::config::generate_snap_token;
4 
5use super::*;
6 
7const ACCEPT_CH_SIZE: usize = 16;
8 
9fn create_association(config: TransportConfig) -> Association {
10 Association::new(
11 None,
12 Arc::new(config),
13 1400,
14 0,
15 SocketAddr::from_str("0.0.0.0:0").unwrap(),
16 None,
17 Instant::now(),
18 )
19}
20 
21#[test]
22fn test_assoc_is_closing() {
23 let closing_states = [
24 AssociationState::ShutdownSent,
25 AssociationState::ShutdownAckSent,
26 AssociationState::ShutdownPending,
27 AssociationState::ShutdownReceived,
28 ];
29 
30 for state in [
31 AssociationState::Closed,
32 AssociationState::CookieWait,
33 AssociationState::CookieEchoed,
34 AssociationState::Established,
35 ] {
36 let a = Association {
37 state,
38 ..Default::default()
39 };
40 
41 assert!(!a.is_closing(), "{state} should not be closing");
42 }
43 
44 for state in closing_states {
45 let a = Association {
46 state,
47 ..Default::default()
48 };
49 
50 assert!(a.is_closing(), "{state} should be closing");
51 assert!(!a.is_closed(), "{state} should not be closed");
52 }
53}
54 
55fn outgoing_reset(rsn: u32, stream_id: StreamId) -> ChunkReconfig {
56 ChunkReconfig {
57 param_a: Some(Box::new(ParamOutgoingResetRequest {
58 reconfig_request_sequence_number: rsn,
59 stream_identifiers: vec![stream_id],
60 ..Default::default()
61 })),
62 ..Default::default()
63 }
64}
65 
66fn insert_active_reset(a: &mut Association, rsn: u32, stream_id: StreamId) {
67 a.reconfigs.insert(rsn, outgoing_reset(rsn, stream_id));
68 a.active_reconfig = Some(rsn);
69}
70 
71fn insert_queued_reset(a: &mut Association, rsn: u32, stream_id: StreamId) {
72 let reset = outgoing_reset(rsn, stream_id);
73 a.reconfigs.insert(rsn, reset.clone());
74 let packet = a.create_packet(vec![Box::new(reset)]);
75 a.control_queue.push_back(packet);
76}
77 
78fn reconfig_response_result(packets: &[Packet], rsn: u32) -> Option<ReconfigResult> {
79 packets.iter().find_map(|packet| {
80 packet.chunks.iter().find_map(|chunk| {
81 let reconfig = chunk.as_any().downcast_ref::<ChunkReconfig>()?;
82 reconfig
83 .param_a
84 .iter()
85 .chain(reconfig.param_b.iter())
86 .find_map(|param| {
87 let response = param.as_any().downcast_ref::<ParamReconfigResponse>()?;
88 (response.reconfig_response_sequence_number == rsn).then_some(response.result)
89 })
90 })
91 })
92}
93 
94#[test]
95fn test_reconfig_in_progress_timeout_does_not_consume_retry_budget() -> Result<()> {
96 let now = Instant::now();
97 let rsn = 7;
98 let mut a = create_association(
99 TransportConfig::default()
100 .with_max_init_retransmits(Some(0))
101 .with_rto_initial_ms(1),
102 );
103 
104 insert_active_reset(&mut a, rsn, 1);
105 a.timers.start(Timer::Reconfig, now, a.rto_mgr.get_rto());
106 
107 let response: Box<dyn Param + Send + Sync> = Box::new(ParamReconfigResponse {
108 reconfig_response_sequence_number: rsn,
109 result: ReconfigResult::InProgress,
110 });
111 a.handle_reconfig_param(&response, &mut vec![])?;
112 
113 // The outbound path restarts the timer after processing InProgress. RFC
114 // 6525 section 5.2.7 H2 requires the next expiry to retransmit without
115 // incrementing the error counter.
116 a.timers
117 .restart_if_stale(Timer::Reconfig, now, a.rto_mgr.get_rto());
118 let deadline = a.timers.get(Timer::Reconfig).unwrap();
119 a.handle_timeout(deadline);
120 
121 assert!(a.reconfigs.contains_key(&rsn));
122 assert!(a.will_retransmit_reconfig);
123 Ok(())
124}
125 
126#[test]
127fn test_reset_complete_only_for_successful_reconfig_response() -> Result<()> {
128 let rsn = 7;
129 let stream_id = 1;
130 
131 for result in [
132 ReconfigResult::SuccessNop,
133 ReconfigResult::SuccessPerformed,
134 ReconfigResult::Denied,
135 ReconfigResult::ErrorWrongSsn,
136 ReconfigResult::ErrorRequestAlreadyInProgress,
137 ReconfigResult::ErrorBadSequenceNumber,
138 ReconfigResult::InProgress,
139 ReconfigResult::Unknown,
140 ] {
141 let mut a = Association::default();
142 a.pending_reset_completions.insert(stream_id);
143 insert_active_reset(&mut a, rsn, stream_id);
144 
145 let response: Box<dyn Param + Send + Sync> = Box::new(ParamReconfigResponse {
146 reconfig_response_sequence_number: rsn,
147 result,
148 });
149 a.handle_reconfig_param(&response, &mut vec![])?;
150 
151 let completed = matches!(
152 a.poll(),
153 Some(Event::Stream(StreamEvent::ResetComplete { id })) if id == stream_id
154 );
155 let should_complete = matches!(
156 result,
157 ReconfigResult::SuccessNop | ReconfigResult::SuccessPerformed
158 );
159 assert_eq!(completed, should_complete, "unexpected result for {result}");
160 }
161 
162 Ok(())
163}
164 
165#[test]
166fn test_reconfig_retransmission_failure_is_terminal() {
167 let rsn = 7;
168 let stream_id = 1;
169 let mut a = Association::default();
170 a.pending_reset_completions.insert(stream_id);
171 insert_active_reset(&mut a, rsn, stream_id);
172 
173 a.on_retransmission_failure(Timer::Reconfig);
174 
175 assert!(matches!(
176 a.poll(),
177 Some(Event::Stream(StreamEvent::ResetFailed {
178 id,
179 reason: StreamResetError::Failed,
180 })) if id == stream_id
181 ));
182}
183 
184#[test]
185fn test_ambiguous_reset_failure_keeps_existing_stream_quarantined() -> Result<()> {
186 let stream_id = 1;
187 let mut a = Association {
188 state: AssociationState::Established,
189 my_next_tsn: 1,
190 ..Default::default()
191 };
192 assert!(
193 a.create_stream(stream_id, false, PayloadProtocolIdentifier::Binary)
194 .is_some()
195 );
196 
197 a.stream(stream_id)?.stop()?;
198 let _ = a.gather_outbound(Instant::now());
199 assert!(a.active_reconfig.is_some());
200 
201 a.on_retransmission_failure(Timer::Reconfig);
202 
203 assert!(a.failed_reset_streams.contains(&stream_id));
204 assert!(
205 !a.stream(stream_id)?.is_writable(),
206 "an unacknowledged reset may have succeeded at the peer, so old SSNs remain unsafe"
207 );
208 Ok(())
209}
210 
211#[test]
212fn test_denied_reset_emits_terminal_event_and_remains_quarantined() -> Result<()> {
213 let rsn = 7;
214 let stream_id = 1;
215 let mut a = Association::default();
216 a.pending_reset_completions.insert(stream_id);
217 insert_active_reset(&mut a, rsn, stream_id);
218 
219 let response: Box<dyn Param + Send + Sync> = Box::new(ParamReconfigResponse {
220 reconfig_response_sequence_number: rsn,
221 result: ReconfigResult::Denied,
222 });
223 a.handle_reconfig_param(&response, &mut vec![])?;
224 
225 assert!(matches!(
226 a.poll(),
227 Some(Event::Stream(StreamEvent::ResetFailed {
228 id,
229 reason: StreamResetError::Denied,
230 })) if id == stream_id
231 ));
232 assert!(!a.pending_reset_completions.contains(&stream_id));
233 assert!(
234 matches!(
235 a.open_stream(stream_id, PayloadProtocolIdentifier::Binary),
236 Err(Error::ErrStreamResetPending)
237 ),
238 "the quarantine must prevent the failed stream ID from being reused"
239 );
240 Ok(())
241}
242 
243#[test]
244fn test_reset_complete_preserves_stream_generations() {
245 let stream_id = 1;
246 let mut a = Association::default();
247 
248 assert!(
249 a.create_stream(stream_id, false, PayloadProtocolIdentifier::Binary)
250 .is_some()
251 );
252 a.unregister_stream(stream_id, true);
253 
254 // Incoming DATA can recreate an id while the previous reciprocal reset is
255 // still pending. A second reset of that id is a distinct generation.
256 assert!(a.get_or_create_stream(stream_id).is_some());
257 a.unregister_stream(stream_id, true);
258 
259 a.emit_reset_complete([stream_id]);
260 a.emit_reset_complete([stream_id]);
261 
262 let mut finished = 0;
263 let mut reset_complete = 0;
264 while let Some(event) = a.poll() {
265 match event {
266 Event::Stream(StreamEvent::Finished { id }) if id == stream_id => finished += 1,
267 Event::Stream(StreamEvent::ResetComplete { id }) if id == stream_id => {
268 reset_complete += 1;
269 }
270 _ => {}
271 }
272 }
273 
274 assert_eq!(finished, 2);
275 assert_eq!(reset_complete, 2);
276}
277 
278#[test]
279fn test_reset_complete_preserves_generations_through_responses() -> Result<()> {
280 let stream_id = 1;
281 let first_rsn = 7;
282 let second_rsn = 8;
283 let mut a = Association::default();
284 
285 // Two Finished events for the same stream ID represent two distinct
286 // incarnations. Each successful reset response must complete one.
287 a.pending_reset_completions.insert(stream_id);
288 a.pending_reset_completions.insert(stream_id);
289 a.reconfigs
290 .insert(first_rsn, outgoing_reset(first_rsn, stream_id));
291 a.reconfigs
292 .insert(second_rsn, outgoing_reset(second_rsn, stream_id));
293 
294 for rsn in [first_rsn, second_rsn] {
295 a.active_reconfig = Some(rsn);
296 let response: Box<dyn Param + Send + Sync> = Box::new(ParamReconfigResponse {
297 reconfig_response_sequence_number: rsn,
298 result: ReconfigResult::SuccessPerformed,
299 });
300 a.handle_reconfig_param(&response, &mut vec![])?;
301 }
302 
303 let reset_complete = core::iter::from_fn(|| a.poll())
304 .filter(|event| {
305 matches!(
306 event,
307 Event::Stream(StreamEvent::ResetComplete { id }) if *id == stream_id
308 )
309 })
310 .count();
311 
312 assert_eq!(
313 reset_complete, 2,
314 "each completed stream generation needs its own ResetComplete"
315 );
316 Ok(())
317}
318 
319#[test]
320fn test_overlapping_generations_success_then_denied_reports_success() -> Result<()> {
321 let stream_id = 1;
322 let mut a = Association::default();
323 
324 a.pending_reset_completions.insert(stream_id);
325 a.pending_reset_completions.insert(stream_id);
326 a.reconfigs.insert(7, outgoing_reset(7, stream_id));
327 a.reconfigs.insert(8, outgoing_reset(8, stream_id));
328 
329 for (rsn, result) in [
330 (7, ReconfigResult::SuccessPerformed),
331 (8, ReconfigResult::Denied),
332 ] {
333 a.active_reconfig = Some(rsn);
334 let response: Box<dyn Param + Send + Sync> = Box::new(ParamReconfigResponse {
335 reconfig_response_sequence_number: rsn,
336 result,
337 });
338 a.handle_reconfig_param(&response, &mut vec![])?;
339 }
340 
341 let reset_complete = core::iter::from_fn(|| a.poll())
342 .filter(|event| {
343 matches!(
344 event,
345 Event::Stream(StreamEvent::ResetComplete { id }) if *id == stream_id
346 )
347 })
348 .count();
349 
350 assert_eq!(
351 reset_complete, 1,
352 "the successful generation still needs its ResetComplete"
353 );
354 Ok(())
355}
356 
357#[test]
358fn test_overlapping_generations_denied_then_success_reports_only_success() -> Result<()> {
359 let stream_id = 1;
360 let mut a = Association::default();
361 
362 a.pending_reset_completions.insert(stream_id);
363 a.pending_reset_completions.insert(stream_id);
364 a.reconfigs.insert(7, outgoing_reset(7, stream_id));
365 a.reconfigs.insert(8, outgoing_reset(8, stream_id));
366 
367 for (rsn, result) in [
368 (7, ReconfigResult::Denied),
369 (8, ReconfigResult::SuccessPerformed),
370 ] {
371 a.active_reconfig = Some(rsn);
372 let response: Box<dyn Param + Send + Sync> = Box::new(ParamReconfigResponse {
373 reconfig_response_sequence_number: rsn,
374 result,
375 });
376 a.handle_reconfig_param(&response, &mut vec![])?;
377 }
378 
379 let reset_complete = core::iter::from_fn(|| a.poll())
380 .filter(|event| {
381 matches!(
382 event,
383 Event::Stream(StreamEvent::ResetComplete { id }) if *id == stream_id
384 )
385 })
386 .count();
387 
388 assert_eq!(
389 reset_complete, 1,
390 "the denied generation must not inherit a later successful completion"
391 );
392 Ok(())
393}
394 
395#[test]
396fn test_outgoing_reset_implicitly_acknowledges_pending_request() -> Result<()> {
397 let local_rsn = 7;
398 let stream_id = 1;
399 let mut a = Association::default();
400 
401 insert_active_reset(&mut a, local_rsn, stream_id);
402 a.timers
403 .start(Timer::Reconfig, Instant::now(), a.rto_mgr.get_rto());
404 
405 // RFC 6525 section 5.2.2 E1: the response sequence number carried by an
406 // incoming Outgoing Reset Request acknowledges our matching request.
407 let request: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
408 reconfig_request_sequence_number: 8,
409 reconfig_response_sequence_number: local_rsn,
410 sender_last_tsn: a.peer_last_tsn,
411 stream_identifiers: vec![],
412 });
413 a.handle_reconfig_param(&request, &mut vec![])?;
414 
415 assert!(
416 !a.reconfigs.contains_key(&local_rsn),
417 "the implicitly acknowledged request must no longer be in flight"
418 );
419 assert!(
420 a.timers.get(Timer::Reconfig).is_none(),
421 "the timer must stop after the final in-flight request is acknowledged"
422 );
423 Ok(())
424}
425 
426#[test]
427fn test_outgoing_reset_does_not_ack_unsent_reciprocal() -> Result<()> {
428 let stream_id = 1;
429 let mut a = Association {
430 my_next_tsn: 1,
431 ..Default::default()
432 };
433 assert!(
434 a.create_stream(stream_id, false, PayloadProtocolIdentifier::Binary)
435 .is_some()
436 );
437 
438 // Handling this request queues a reciprocal Outgoing Reset Request in
439 // `reply`, but it has not been transmitted and its timer is not running.
440 let request: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
441 reconfig_request_sequence_number: 7,
442 reconfig_response_sequence_number: u32::MAX,
443 sender_last_tsn: a.peer_last_tsn,
444 stream_identifiers: vec![stream_id],
445 });
446 let mut reply = vec![];
447 a.handle_reconfig_param(&request, &mut reply)?;
448 
449 let reciprocal_rsn = *a.reconfigs.keys().next().unwrap();
450 assert!(!reply.is_empty());
451 assert!(a.timers.get(Timer::Reconfig).is_none());
452 
453 // RFC 6525 section 5.2.2 E1 only acknowledges a request for which the
454 // Re-configuration Timer is running. A peer must not be able to pre-ack
455 // the queued reciprocal before the application polls it for transmission.
456 let premature_ack: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
457 reconfig_request_sequence_number: 8,
458 reconfig_response_sequence_number: reciprocal_rsn,
459 sender_last_tsn: a.peer_last_tsn,
460 stream_identifiers: vec![],
461 });
462 a.handle_reconfig_param(&premature_ack, &mut vec![])?;
463 
464 let reset_complete = core::iter::from_fn(|| a.poll())
465 .filter(|event| {
466 matches!(
467 event,
468 Event::Stream(StreamEvent::ResetComplete { id }) if *id == stream_id
469 )
470 })
471 .count();
472 assert_eq!(
473 (a.reconfigs.contains_key(&reciprocal_rsn), reset_complete),
474 (true, 0),
475 "an unsent reciprocal must remain pending and cannot complete a reset"
476 );
477 Ok(())
478}
479 
480#[test]
481fn test_reset_complete_does_not_override_newer_failure() -> Result<()> {
482 let stream_id = 1;
483 let mut a = Association::default();
484 
485 a.pending_reset_completions.insert(stream_id);
486 a.pending_reset_completions.insert(stream_id);
487 a.reconfigs.insert(7, outgoing_reset(7, stream_id));
488 a.reconfigs.insert(8, outgoing_reset(8, stream_id));
489 
490 for (rsn, result) in [
491 (7, ReconfigResult::SuccessPerformed),
492 (8, ReconfigResult::Denied),
493 ] {
494 a.active_reconfig = Some(rsn);
495 let response: Box<dyn Param + Send + Sync> = Box::new(ParamReconfigResponse {
496 reconfig_response_sequence_number: rsn,
497 result,
498 });
499 a.handle_reconfig_param(&response, &mut vec![])?;
500 }
501 
502 assert!(matches!(
503 a.poll(),
504 Some(Event::Stream(StreamEvent::ResetComplete { id })) if id == stream_id
505 ));
506 assert!(matches!(
507 a.poll(),
508 Some(Event::Stream(StreamEvent::ResetFailed {
509 id,
510 reason: StreamResetError::Denied,
511 })) if id == stream_id
512 ));
513 assert!(matches!(
514 a.open_stream(stream_id, PayloadProtocolIdentifier::Binary),
515 Err(Error::ErrStreamResetPending)
516 ));
517 Ok(())
518}
519 
520#[test]
521fn test_reconfig_response_does_not_ack_unsent_reciprocal() -> Result<()> {
522 let stream_id = 1;
523 let mut a = Association {
524 my_next_tsn: 1,
525 ..Default::default()
526 };
527 assert!(
528 a.create_stream(stream_id, false, PayloadProtocolIdentifier::Binary)
529 .is_some()
530 );
531 
532 let request: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
533 reconfig_request_sequence_number: 7,
534 reconfig_response_sequence_number: u32::MAX,
535 sender_last_tsn: a.peer_last_tsn,
536 stream_identifiers: vec![stream_id],
537 });
538 let mut reply = vec![];
539 a.handle_reconfig_param(&request, &mut reply)?;
540 
541 let reciprocal_rsn = *a.reconfigs.keys().next().unwrap();
542 assert!(!reply.is_empty());
543 assert!(a.timers.get(Timer::Reconfig).is_none());
544 
545 // The reciprocal is only queued in the reply and has not been transmitted.
546 // RFC 6525 H1 says to ignore a response for an RSN whose timer is not running.
547 let premature_response: Box<dyn Param + Send + Sync> = Box::new(ParamReconfigResponse {
548 reconfig_response_sequence_number: reciprocal_rsn,
549 result: ReconfigResult::SuccessPerformed,
550 });
551 a.handle_reconfig_param(&premature_response, &mut vec![])?;
552 
553 let reset_complete = core::iter::from_fn(|| a.poll())
554 .filter(|event| {
555 matches!(
556 event,
557 Event::Stream(StreamEvent::ResetComplete { id }) if *id == stream_id
558 )
559 })
560 .count();
561 assert_eq!(
562 (a.reconfigs.contains_key(&reciprocal_rsn), reset_complete),
563 (true, 0),
564 "RFC 6525 H1 requires a response for an RSN whose timer is not running to be ignored"
565 );
566 Ok(())
567}
568 
569#[test]
570fn test_outgoing_reset_does_not_ack_unsent_reciprocal_with_unrelated_timer() -> Result<()> {
571 let stream_id = 1;
572 let sent_rsn = 99;
573 let mut a = Association {
574 my_next_tsn: 1,
575 ..Default::default()
576 };
577 assert!(
578 a.create_stream(stream_id, false, PayloadProtocolIdentifier::Binary)
579 .is_some()
580 );
581 
582 insert_active_reset(&mut a, sent_rsn, 2);
583 a.timers
584 .start(Timer::Reconfig, Instant::now(), a.rto_mgr.get_rto());
585 
586 // The global timer belongs to the already-sent request above. Processing
587 // this incoming request queues a distinct reciprocal, but does not send it.
588 let request: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
589 reconfig_request_sequence_number: 7,
590 reconfig_response_sequence_number: u32::MAX,
591 sender_last_tsn: a.peer_last_tsn,
592 stream_identifiers: vec![stream_id],
593 });
594 let mut reply = vec![];
595 a.handle_reconfig_param(&request, &mut reply)?;
596 
597 let reciprocal_rsn = *a.reconfigs.keys().find(|&&rsn| rsn != sent_rsn).unwrap();
598 assert!(!reply.is_empty());
599 
600 let premature_ack: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
601 reconfig_request_sequence_number: 8,
602 reconfig_response_sequence_number: reciprocal_rsn,
603 sender_last_tsn: a.peer_last_tsn,
604 stream_identifiers: vec![3],
605 });
606 a.handle_reconfig_param(&premature_ack, &mut vec![])?;
607 
608 let reset_complete = core::iter::from_fn(|| a.poll())
609 .filter(|event| {
610 matches!(
611 event,
612 Event::Stream(StreamEvent::ResetComplete { id }) if *id == stream_id
613 )
614 })
615 .count();
616 assert_eq!(
617 (a.reconfigs.contains_key(&reciprocal_rsn), reset_complete),
618 (true, 0),
619 "a timer running for another RSN must not make an unsent reciprocal acknowledgeable"
620 );
621 Ok(())
622}
623 
624#[test]
625fn test_failed_generation_stays_quarantined_after_other_success() -> Result<()> {
626 let stream_id = 1;
627 let mut a = Association::default();
628 
629 a.pending_reset_completions.insert(stream_id);
630 a.pending_reset_completions.insert(stream_id);
631 a.reconfigs.insert(7, outgoing_reset(7, stream_id));
632 a.reconfigs.insert(8, outgoing_reset(8, stream_id));
633 
634 // The older generation completes, but the newer generation is denied.
635 // A success for the former cannot make the latter safe to reuse.
636 for (rsn, result) in [
637 (7, ReconfigResult::SuccessPerformed),
638 (8, ReconfigResult::Denied),
639 ] {
640 a.active_reconfig = Some(rsn);
641 let response: Box<dyn Param + Send + Sync> = Box::new(ParamReconfigResponse {
642 reconfig_response_sequence_number: rsn,
643 result,
644 });
645 a.handle_reconfig_param(&response, &mut vec![])?;
646 }
647 
648 assert!(matches!(
649 a.open_stream(stream_id, PayloadProtocolIdentifier::Binary),
650 Err(Error::ErrStreamResetPending)
651 ));
652 Ok(())
653}
654 
655#[test]
656fn test_second_reconfig_request_stays_buffered_while_timer_runs() -> Result<()> {
657 let stream_id = 1;
658 let sent_rsn = 99;
659 let mut a = Association {
660 state: AssociationState::Established,
661 my_next_tsn: 1,
662 ..Default::default()
663 };
664 assert!(
665 a.create_stream(stream_id, false, PayloadProtocolIdentifier::Binary)
666 .is_some()
667 );
668 
669 insert_active_reset(&mut a, sent_rsn, 2);
670 a.timers
671 .start(Timer::Reconfig, Instant::now(), a.rto_mgr.get_rto());
672 
673 // Processing the peer's request creates a reciprocal request while another
674 // local request is already in flight. RFC 6525 section 5.1.1 requires the
675 // reciprocal request to remain buffered until the running timer stops.
676 let request: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
677 reconfig_request_sequence_number: 7,
678 reconfig_response_sequence_number: u32::MAX,
679 sender_last_tsn: a.peer_last_tsn,
680 stream_identifiers: vec![stream_id],
681 });
682 let mut reply = vec![];
683 a.handle_reconfig_param(&request, &mut reply)?;
684 
685 let reciprocal_rsn = *a.reconfigs.keys().find(|&&rsn| rsn != sent_rsn).unwrap();
686 a.control_queue.extend(reply);
687 
688 assert!(a.poll_transmit(Instant::now()).is_some());
689 assert_eq!(
690 a.active_reconfig,
691 Some(sent_rsn),
692 "a second request must remain buffered while the first request's timer runs"
693 );
694 assert!(a.reconfigs.contains_key(&reciprocal_rsn));
695 assert_eq!(a.control_queue.len(), 1);
696 
697 let response: Box<dyn Param + Send + Sync> = Box::new(ParamReconfigResponse {
698 reconfig_response_sequence_number: sent_rsn,
699 result: ReconfigResult::SuccessPerformed,
700 });
701 a.handle_reconfig_param(&response, &mut vec![])?;
702 assert!(a.poll_transmit(Instant::now()).is_some());
703 assert_eq!(a.active_reconfig, Some(reciprocal_rsn));
704 assert!(a.control_queue.is_empty());
705 Ok(())
706}
707 
708#[test]
709fn test_peer_recreated_quarantined_stream_is_not_writable() -> Result<()> {
710 let stream_id = 1;
711 let mut a = Association::default();
712 
713 // The peer may start its new incoming generation before our reciprocal
714 // outgoing reset is acknowledged. Reading that generation is safe, but
715 // sending with a freshly initialized SSN is not yet safe.
716 a.pending_reset_completions.insert(stream_id);
717 insert_active_reset(&mut a, 7, stream_id);
718 
719 assert!(a.get_or_create_stream(stream_id).is_some());
720 assert!(
721 !a.stream(stream_id)?.is_writable(),
722 "a pending reciprocal reset must quarantine the outgoing direction"
723 );
724 Ok(())
725}
726 
727#[test]
728fn test_pending_reset_blocks_existing_stream_writes() -> Result<()> {
729 let stream_id = 1;
730 let mut a = Association::default();
731 assert!(
732 a.create_stream(stream_id, false, PayloadProtocolIdentifier::Binary)
733 .is_some()
734 );
735 insert_active_reset(&mut a, 7, stream_id);
736 
737 assert!(
738 !a.stream(stream_id)?.is_writable(),
739 "new SSNs must not be assigned while an outgoing reset is pending"
740 );
741 Ok(())
742}
743 
744#[test]
745fn test_reset_completion_does_not_reopen_finished_write_half() -> Result<()> {
746 let stream_id = 1;
747 let rsn = 7;
748 let mut a = Association::default();
749 assert!(
750 a.create_stream(stream_id, false, PayloadProtocolIdentifier::Binary)
751 .is_some()
752 );
753 insert_active_reset(&mut a, rsn, stream_id);
754 
755 // Finishing the write half is permanent, including while a reset request
756 // temporarily makes the stream non-writable for protocol reasons.
757 a.stream(stream_id)?.finish()?;
758 
759 let response: Box<dyn Param + Send + Sync> = Box::new(ParamReconfigResponse {
760 reconfig_response_sequence_number: rsn,
761 result: ReconfigResult::SuccessPerformed,
762 });
763 a.handle_reconfig_param(&response, &mut vec![])?;
764 
765 assert!(
766 !a.stream(stream_id)?.is_writable(),
767 "reset completion must not undo Stream::finish()"
768 );
769 Ok(())
770}
771 
772#[test]
773fn test_successful_outgoing_reset_restarts_stream_sequence_number() -> Result<()> {
774 let stream_id = 1;
775 let rsn = 7;
776 let mut a = Association::default();
777 assert!(
778 a.create_stream(stream_id, false, PayloadProtocolIdentifier::Binary)
779 .is_some()
780 );
781 a.streams.get_mut(&stream_id).unwrap().sequence_number = 9;
782 insert_active_reset(&mut a, rsn, stream_id);
783 
784 let response: Box<dyn Param + Send + Sync> = Box::new(ParamReconfigResponse {
785 reconfig_response_sequence_number: rsn,
786 result: ReconfigResult::SuccessPerformed,
787 });
788 a.handle_reconfig_param(&response, &mut vec![])?;
789 
790 assert_eq!(
791 a.streams.get(&stream_id).unwrap().sequence_number,
792 0,
793 "RFC 6525 section 5.2.7 H4 requires the affected outgoing SSN to reset"
794 );
795 Ok(())
796}
797 
798#[test]
799fn test_only_one_buffered_reconfig_is_sent_when_timer_is_idle() {
800 let mut a = Association {
801 state: AssociationState::Established,
802 ..Default::default()
803 };
804 
805 for (rsn, stream_id) in [(7, 1), (8, 2)] {
806 insert_queued_reset(&mut a, rsn, stream_id);
807 }
808 assert!(a.timers.get(Timer::Reconfig).is_none());
809 
810 assert!(a.poll_transmit(Instant::now()).is_some());
811 assert_eq!(a.active_reconfig, Some(7));
812 assert_eq!(a.control_queue.len(), 1);
813 assert!(a.reconfigs.contains_key(&8));
814 
815 let response: Box<dyn Param + Send + Sync> = Box::new(ParamReconfigResponse {
816 reconfig_response_sequence_number: 7,
817 result: ReconfigResult::SuccessPerformed,
818 });
819 a.handle_reconfig_param(&response, &mut vec![]).unwrap();
820 assert!(a.poll_transmit(Instant::now()).is_some());
821 assert_eq!(a.active_reconfig, Some(8));
822 assert!(a.control_queue.is_empty());
823}
824 
825#[test]
826fn test_close_discards_buffered_reconfig_requests() -> Result<()> {
827 let mut a = Association {
828 state: AssociationState::Established,
829 ..Default::default()
830 };
831 insert_active_reset(&mut a, 7, 1);
832 insert_queued_reset(&mut a, 8, 2);
833 a.timers
834 .start(Timer::Reconfig, Instant::now(), a.rto_mgr.get_rto());
835 
836 a.close()?;
837 
838 assert!(
839 a.poll_transmit(Instant::now()).is_none(),
840 "a closed association must not send a buffered stream-reset request"
841 );
842 assert!(a.active_reconfig.is_none());
843 assert!(a.reconfigs.is_empty());
844 Ok(())
845}
846 
847#[test]
848fn test_in_progress_keeps_later_request_buffered() -> Result<()> {
849 let mut a = Association {
850 state: AssociationState::Established,
851 ..Default::default()
852 };
853 
854 insert_active_reset(&mut a, 7, 1);
855 insert_queued_reset(&mut a, 8, 2);
856 a.timers
857 .start(Timer::Reconfig, Instant::now(), a.rto_mgr.get_rto());
858 
859 let response: Box<dyn Param + Send + Sync> = Box::new(ParamReconfigResponse {
860 reconfig_response_sequence_number: 7,
861 result: ReconfigResult::InProgress,
862 });
863 a.handle_reconfig_param(&response, &mut vec![])?;
864 assert!(a.reconfigs.contains_key(&7));
865 
866 let _ = a.poll_transmit(Instant::now());
867 assert_eq!(a.active_reconfig, Some(7));
868 assert_eq!(a.control_queue.len(), 1);
869 Ok(())
870}
871 
872#[test]
873fn test_local_reset_stays_queued_while_reconfig_timer_runs() -> Result<()> {
874 let mut a = Association {
875 state: AssociationState::Established,
876 my_next_tsn: 1,
877 ..Default::default()
878 };
879 
880 insert_active_reset(&mut a, 7, 1);
881 a.timers
882 .start(Timer::Reconfig, Instant::now(), a.rto_mgr.get_rto());
883 a.send_reset_request(2)?;
884 
885 let _ = a.poll_transmit(Instant::now());
886 assert_eq!(
887 a.reconfigs.len(),
888 1,
889 "a locally initiated reset must remain queued while another request is in flight"
890 );
891 assert_eq!(
892 a.pending_reset_streams.iter().copied().collect::<Vec<_>>(),
893 [2]
894 );
895 
896 let response: Box<dyn Param + Send + Sync> = Box::new(ParamReconfigResponse {
897 reconfig_response_sequence_number: 7,
898 result: ReconfigResult::SuccessPerformed,
899 });
900 a.handle_reconfig_param(&response, &mut vec![])?;
901 assert!(a.poll_transmit(Instant::now()).is_some());
902 assert_ne!(a.active_reconfig, Some(7));
903 assert!(a.active_reconfig.is_some());
904 assert!(a.pending_reset_streams.is_empty());
905 Ok(())
906}
907 
908#[test]
909fn test_local_reset_does_not_overtake_pending_data() -> Result<()> {
910 let mut a = Association {
911 state: AssociationState::Established,
912 my_next_tsn: 1,
913 ..Default::default()
914 };
915 a.inflight_queue.push_no_check(ChunkPayloadData {
916 tsn: 1,
917 user_data: Bytes::from_static(b"inflight"),
918 ..Default::default()
919 });
920 a.pending_queue.push(ChunkPayloadData {
921 stream_identifier: 1,
922 beginning_fragment: true,
923 ending_fragment: true,
924 user_data: Bytes::from_static(b"pending"),
925 ..Default::default()
926 });
927 a.send_reset_request(1)?;
928 
929 let _ = a.poll_transmit(Instant::now());
930 assert!(a.active_reconfig.is_none());
931 assert!(a.reconfigs.is_empty());
932 assert_eq!(a.pending_reset_streams.len(), 1);
933 Ok(())
934}
935 
936#[test]
937fn test_local_reset_is_not_blocked_by_unrelated_pending_data() -> Result<()> {
938 let reset_stream_id = 1;
939 let unrelated_stream_id = 2;
940 let mut a = Association {
941 state: AssociationState::Established,
942 my_next_tsn: 2,
943 cwnd: 0,
944 rwnd: 0,
945 mtu: 1400,
946 ..Default::default()
947 };
948 a.inflight_queue.push_no_check(ChunkPayloadData {
949 tsn: 1,
950 stream_identifier: unrelated_stream_id,
951 user_data: Bytes::from_static(b"inflight"),
952 ..Default::default()
953 });
954 a.pending_queue.push(ChunkPayloadData {
955 stream_identifier: unrelated_stream_id,
956 beginning_fragment: true,
957 ending_fragment: true,
958 user_data: Bytes::from_static(b"flow controlled"),
959 ..Default::default()
960 });
961 a.send_reset_request(reset_stream_id)?;
962 
963 let _ = a.gather_outbound(Instant::now());
964 
965 assert!(
966 a.active_reconfig.is_some(),
967 "flow-controlled DATA on another stream must not starve this reset"
968 );
969 assert!(
970 !a.pending_queue.is_empty(),
971 "the test requires the unrelated DATA to remain flow controlled"
972 );
973 Ok(())
974}
975 
976#[test]
977fn test_reset_sender_last_tsn_wraps_at_zero() -> Result<()> {
978 let stream_id = 1;
979 let mut a = Association {
980 state: AssociationState::Established,
981 my_next_tsn: 0,
982 ..Default::default()
983 };
984 a.send_reset_request(stream_id)?;
985 
986 let _ = a.gather_outbound(Instant::now());
987 
988 let reset = a
989 .reconfigs
990 .get(&a.active_reconfig.unwrap())
991 .unwrap()
992 .param_a
993 .as_ref()
994 .and_then(|param| param.as_any().downcast_ref::<ParamOutgoingResetRequest>())
995 .unwrap();
996 assert_eq!(
997 reset.sender_last_tsn,
998 u32::MAX,
999 "the TSN preceding zero must wrap to u32::MAX"
1000 );
1001 Ok(())
1002}
1003 
1004#[test]
1005fn test_reset_request_sequence_number_wraps() -> Result<()> {
1006 let stream_id = 1;
1007 let mut a = Association {
1008 state: AssociationState::Established,
1009 my_next_tsn: 1,
1010 my_next_rsn: u32::MAX,
1011 ..Default::default()
1012 };
1013 a.send_reset_request(stream_id)?;
1014 
1015 let _ = a.gather_outbound(Instant::now());
1016 
1017 assert_eq!(a.active_reconfig, Some(u32::MAX));
1018 assert_eq!(
1019 a.my_next_rsn, 0,
1020 "the re-configuration request sequence number must wrap"
1021 );
1022 Ok(())
1023}
1024 
1025#[test]
1026fn test_local_reset_acknowledges_latest_peer_request_sequence() -> Result<()> {
1027 let peer_rsn = 41;
1028 let mut a = Association {
1029 state: AssociationState::Established,
1030 my_next_tsn: 1,
1031 ..Default::default()
1032 };
1033 
1034 let peer_request: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
1035 reconfig_request_sequence_number: peer_rsn,
1036 reconfig_response_sequence_number: u32::MAX,
1037 sender_last_tsn: a.peer_last_tsn,
1038 stream_identifiers: vec![],
1039 });
1040 a.handle_reconfig_param(&peer_request, &mut vec![])?;
1041 assert_eq!(a.max_completed_reconfig_rsn, Some(peer_rsn));
1042 
1043 a.send_reset_request(1)?;
1044 let _ = a.gather_outbound(Instant::now());
1045 
1046 let reset = a
1047 .reconfigs
1048 .get(&a.active_reconfig.unwrap())
1049 .unwrap()
1050 .param_a
1051 .as_ref()
1052 .and_then(|param| param.as_any().downcast_ref::<ParamOutgoingResetRequest>())
1053 .unwrap();
1054 assert_eq!(
1055 reset.reconfig_response_sequence_number, peer_rsn,
1056 "RFC 6525 section 5.1.2 A4 requires the latest peer RSN in the response field"
1057 );
1058 Ok(())
1059}
1060 
1061#[test]
1062fn test_reciprocal_reset_covers_preexisting_pending_data() -> Result<()> {
1063 let stream_id = 1;
1064 let mut a = Association {
1065 state: AssociationState::Established,
1066 my_next_tsn: 1,
1067 cwnd: 1400,
1068 rwnd: 1400,
1069 mtu: 1400,
1070 ..Default::default()
1071 };
1072 assert!(
1073 a.create_stream(stream_id, false, PayloadProtocolIdentifier::Binary)
1074 .is_some()
1075 );
1076 a.pending_queue.push(ChunkPayloadData {
1077 stream_identifier: stream_id,
1078 stream_sequence_number: 4,
1079 beginning_fragment: true,
1080 ending_fragment: true,
1081 user_data: Bytes::from_static(b"pending"),
1082 ..Default::default()
1083 });
1084 
1085 // Receiving the peer's outgoing reset creates a reciprocal outgoing reset.
1086 // DATA accepted before that request must be assigned a TSN covered by the
1087 // reciprocal request's Sender's Last Assigned TSN boundary.
1088 let request: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
1089 reconfig_request_sequence_number: 7,
1090 reconfig_response_sequence_number: u32::MAX,
1091 sender_last_tsn: a.peer_last_tsn,
1092 stream_identifiers: vec![stream_id],
1093 });
1094 let mut reply = vec![];
1095 a.handle_reconfig_param(&request, &mut reply)?;
1096 a.control_queue.extend(reply);
1097 
1098 let _ = a.gather_outbound(Instant::now());
1099 
1100 let data_tsn = a.inflight_queue.get(1).unwrap().tsn;
1101 let reciprocal = a
1102 .reconfigs
1103 .get(&a.active_reconfig.unwrap())
1104 .unwrap()
1105 .param_a
1106 .as_ref()
1107 .and_then(|param| param.as_any().downcast_ref::<ParamOutgoingResetRequest>())
1108 .unwrap();
1109 assert!(
1110 sna32gte(reciprocal.sender_last_tsn, data_tsn),
1111 "the reciprocal reset boundary must cover DATA queued before the reset"
1112 );
1113 Ok(())
1114}
1115 
1116#[test]
1117fn test_data_above_deferred_reset_boundary_waits_for_new_generation() -> Result<()> {
1118 let stream_id = 1;
1119 let mut a = Association {
1120 state: AssociationState::Established,
1121 peer_last_tsn: 0,
1122 my_next_tsn: 1,
1123 max_receive_buffer_size: 1024,
1124 max_receive_message_size: 1024,
1125 max_payload_size: 1200,
1126 ..Default::default()
1127 };
1128 assert!(
1129 a.create_stream(stream_id, false, PayloadProtocolIdentifier::Binary)
1130 .is_some()
1131 );
1132 
1133 let request: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
1134 reconfig_request_sequence_number: 7,
1135 reconfig_response_sequence_number: u32::MAX,
1136 sender_last_tsn: 1,
1137 stream_identifiers: vec![stream_id],
1138 });
1139 a.handle_reconfig_param(&request, &mut vec![])?;
1140 assert!(a.reconfig_requests.contains_key(&7));
1141 
1142 // TSN 2 belongs to the post-reset generation, but arrives while TSN 1 is
1143 // still missing and the reset therefore remains InProgress.
1144 a.handle_data(&ChunkPayloadData {
1145 tsn: 2,
1146 stream_identifier: stream_id,
1147 stream_sequence_number: 0,
1148 beginning_fragment: true,
1149 ending_fragment: true,
1150 user_data: Bytes::from_static(b"new"),
1151 ..Default::default()
1152 })?;
1153 
1154 assert!(
1155 a.poll().is_none(),
1156 "post-reset DATA must not notify the old stream generation"
1157 );
1158 assert!(
1159 a.stream(stream_id)?.read()?.is_none(),
1160 "post-reset DATA must not be readable before the reset boundary arrives"
1161 );
1162 
1163 // Once TSN 1 arrives, its payload remains readable as the last message of
1164 // the old generation. Held TSN 2 is released only after that generation is
1165 // drained, into a newly created generation whose SSN starts at zero.
1166 a.handle_data(&ChunkPayloadData {
1167 tsn: 1,
1168 stream_identifier: stream_id,
1169 stream_sequence_number: 0,
1170 beginning_fragment: true,
1171 ending_fragment: true,
1172 user_data: Bytes::from_static(b"old"),
1173 ..Default::default()
1174 })?;
1175 
1176 let message = a.stream(stream_id)?.read()?.unwrap();
1177 let mut payload = [0; 3];
1178 assert_eq!(message.read(&mut payload)?, payload.len());
1179 assert_eq!(&payload, b"old");
1180 
1181 let message = a.stream(stream_id)?.read()?.unwrap();
1182 let mut payload = [0; 3];
1183 assert_eq!(message.read(&mut payload)?, payload.len());
1184 assert_eq!(&payload, b"new");
1185 Ok(())
1186}
1187 
1188#[test]
1189fn test_reset_boundary_data_remains_readable_before_finished() -> Result<()> {
1190 let stream_id = 1;
1191 let mut a = Association {
1192 state: AssociationState::Established,
1193 peer_last_tsn: 0,
1194 my_next_tsn: 1,
1195 max_receive_buffer_size: 1024,
1196 max_receive_message_size: 1024,
1197 max_payload_size: 1200,
1198 ..Default::default()
1199 };
1200 assert!(
1201 a.create_stream(stream_id, false, PayloadProtocolIdentifier::Binary)
1202 .is_some()
1203 );
1204 
1205 let request: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
1206 reconfig_request_sequence_number: 7,
1207 reconfig_response_sequence_number: u32::MAX,
1208 sender_last_tsn: 1,
1209 stream_identifiers: vec![stream_id],
1210 });
1211 a.handle_reconfig_param(&request, &mut vec![])?;
1212 
1213 a.handle_data(&ChunkPayloadData {
1214 tsn: 1,
1215 stream_identifier: stream_id,
1216 stream_sequence_number: 0,
1217 beginning_fragment: true,
1218 ending_fragment: true,
1219 user_data: Bytes::from_static(b"last"),
1220 ..Default::default()
1221 })?;
1222 
1223 loop {
1224 match a.poll() {
1225 Some(Event::Stream(StreamEvent::Readable { id })) if id == stream_id => break,
1226 Some(Event::Stream(StreamEvent::Finished { id })) if id == stream_id => {
1227 panic!("Finished must not overtake readable boundary DATA")
1228 }
1229 Some(_) => {}
1230 None => panic!("missing Readable event for boundary DATA"),
1231 }
1232 }
1233 
1234 assert!(
1235 a.poll().is_none(),
1236 "Finished must remain hidden until the readable boundary DATA is drained"
1237 );
1238 
1239 let message = a.stream(stream_id)?.read()?.unwrap();
1240 let mut payload = [0; 4];
1241 assert_eq!(message.read(&mut payload)?, payload.len());
1242 assert_eq!(&payload, b"last");
1243 assert!(matches!(
1244 core::iter::from_fn(|| a.poll()).find(|event| matches!(
1245 event,
1246 Event::Stream(StreamEvent::Finished { id }) if *id == stream_id
1247 )),
1248 Some(Event::Stream(StreamEvent::Finished { id })) if id == stream_id
1249 ));
1250 Ok(())
1251}
1252 
1253#[test]
1254fn test_successive_resets_preserve_each_unread_generation() -> Result<()> {
1255 let stream_id = 1;
1256 let mut a = Association {
1257 state: AssociationState::Established,
1258 peer_last_tsn: 0,
1259 my_next_tsn: 1,
1260 max_receive_buffer_size: 1024,
1261 max_receive_message_size: 1024,
1262 max_payload_size: 1200,
1263 ..Default::default()
1264 };
1265 assert!(
1266 a.create_stream(stream_id, false, PayloadProtocolIdentifier::Binary)
1267 .is_some()
1268 );
1269 
1270 let first_reset: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
1271 reconfig_request_sequence_number: 7,
1272 reconfig_response_sequence_number: u32::MAX,
1273 sender_last_tsn: 1,
1274 stream_identifiers: vec![stream_id],
1275 });
1276 a.handle_reconfig_param(&first_reset, &mut vec![])?;
1277 a.handle_data(&ChunkPayloadData {
1278 tsn: 1,
1279 stream_identifier: stream_id,
1280 stream_sequence_number: 0,
1281 beginning_fragment: true,
1282 ending_fragment: true,
1283 user_data: Bytes::from_static(b"gen-a"),
1284 ..Default::default()
1285 })?;
1286 
1287 // The successor arrives while generation A is still unread.
1288 a.handle_data(&ChunkPayloadData {
1289 tsn: 2,
1290 stream_identifier: stream_id,
1291 stream_sequence_number: 0,
1292 beginning_fragment: true,
1293 ending_fragment: true,
1294 user_data: Bytes::from_static(b"gen-b"),
1295 ..Default::default()
1296 })?;
1297 let second_reset: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
1298 reconfig_request_sequence_number: 8,
1299 reconfig_response_sequence_number: u32::MAX,
1300 sender_last_tsn: 2,
1301 stream_identifiers: vec![stream_id],
1302 });
1303 a.handle_reconfig_param(&second_reset, &mut vec![])?;
1304 
1305 // Model successful reciprocal handshakes for both reset generations.
1306 a.emit_reset_complete([stream_id]);
1307 a.emit_reset_complete([stream_id]);
1308 
1309 let first = a.stream(stream_id)?.read()?.unwrap();
1310 let mut payload = [0; 5];
1311 assert_eq!(first.read(&mut payload)?, payload.len());
1312 assert_eq!(&payload, b"gen-a");
1313 
1314 let mut finished = 0;
1315 let mut reset_complete = 0;
1316 while reset_complete == 0 {
1317 match a.poll() {
1318 Some(Event::Stream(StreamEvent::Finished { id })) if id == stream_id => finished += 1,
1319 Some(Event::Stream(StreamEvent::ResetComplete { id })) if id == stream_id => {
1320 reset_complete += 1;
1321 }
1322 Some(_) => {}
1323 None => panic!("generation A terminal event was blocked by unread generation B"),
1324 }
1325 }
1326 assert_eq!(finished, 1);
1327 
1328 assert!(
1329 core::iter::from_fn(|| a.poll()).any(|event| matches!(
1330 event,
1331 Event::Stream(StreamEvent::Readable { id }) if id == stream_id
1332 )),
1333 "generation B must advertise readability before its Finished event"
1334 );
1335 
1336 let second = a.stream(stream_id)?.read()?.unwrap();
1337 let mut payload = [0; 5];
1338 assert_eq!(second.read(&mut payload)?, payload.len());
1339 assert_eq!(&payload, b"gen-b");
1340 assert!(
1341 a.stream(stream_id).is_err(),
1342 "generation B was also reset and must retire after it is drained"
1343 );
1344 
1345 while let Some(event) = a.poll() {
1346 match event {
1347 Event::Stream(StreamEvent::Finished { id }) if id == stream_id => finished += 1,
1348 Event::Stream(StreamEvent::ResetComplete { id }) if id == stream_id => {
1349 reset_complete += 1;
1350 }
1351 _ => {}
1352 }
1353 }
1354 assert_eq!(finished, 2);
1355 assert_eq!(reset_complete, 2);
1356 Ok(())
1357}
1358 
1359fn association_with_retiring_boundary_data() -> Result<Association> {
1360 let stream_id = 1;
1361 let mut a = Association {
1362 state: AssociationState::Established,
1363 peer_last_tsn: 0,
1364 my_next_tsn: 1,
1365 max_receive_buffer_size: 1024,
1366 max_receive_message_size: 1024,
1367 max_payload_size: 1200,
1368 use_forward_tsn: true,
1369 ..Default::default()
1370 };
1371 assert!(
1372 a.create_stream(stream_id, false, PayloadProtocolIdentifier::Binary)
1373 .is_some()
1374 );
1375 let reset: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
1376 reconfig_request_sequence_number: 7,
1377 reconfig_response_sequence_number: u32::MAX,
1378 sender_last_tsn: 1,
1379 stream_identifiers: vec![stream_id],
1380 });
1381 a.handle_reconfig_param(&reset, &mut vec![])?;
1382 a.handle_data(&ChunkPayloadData {
1383 tsn: 1,
1384 stream_identifier: stream_id,
1385 stream_sequence_number: 0,
1386 beginning_fragment: true,
1387 ending_fragment: true,
1388 user_data: Bytes::from_static(b"old"),
1389 ..Default::default()
1390 })?;
1391 Ok(a)
1392}
1393 
1394fn assert_successor_ssn_one_is_readable(mut a: Association) -> Result<()> {
1395 let stream_id = 1;
1396 let old = a.stream(stream_id)?.read()?.unwrap();
1397 let mut payload = [0; 3];
1398 assert_eq!(old.read(&mut payload)?, payload.len());
1399 assert_eq!(&payload, b"old");
1400 
1401 a.handle_data(&ChunkPayloadData {
1402 tsn: 3,
1403 stream_identifier: stream_id,
1404 stream_sequence_number: 1,
1405 beginning_fragment: true,
1406 ending_fragment: true,
1407 user_data: Bytes::from_static(b"new"),
1408 ..Default::default()
1409 })?;
1410 let successor = a.stream(stream_id)?.read()?.unwrap();
1411 let mut payload = [0; 3];
1412 assert_eq!(successor.read(&mut payload)?, payload.len());
1413 assert_eq!(&payload, b"new");
1414 Ok(())
1415}
1416 
1417#[test]
1418fn test_forward_tsn_skip_is_applied_to_reset_successor() -> Result<()> {
1419 let mut a = association_with_retiring_boundary_data()?;
1420 a.handle_forward_tsn(&ChunkForwardTsn {
1421 new_cumulative_tsn: 2,
1422 streams: vec![ChunkForwardTsnStream {
1423 identifier: 1,
1424 sequence: 0,
1425 }],
1426 })?;
1427 assert_successor_ssn_one_is_readable(a)
1428}
1429 
1430#[test]
1431fn test_i_forward_tsn_skip_is_applied_to_reset_successor() -> Result<()> {
1432 let mut a = association_with_retiring_boundary_data()?;
1433 a.handle_i_forward_tsn(&ChunkIForwardTsn {
1434 new_cumulative_tsn: 2,
1435 streams: vec![ChunkIForwardTsnStream {
1436 identifier: 1,
1437 unordered: false,
1438 mid: 0,
1439 }],
1440 })?;
1441 assert_successor_ssn_one_is_readable(a)
1442}
1443 
1444#[test]
1445fn test_forward_tsn_during_in_progress_reset_applies_to_old_generation() -> Result<()> {
1446 let stream_id = 1;
1447 let mut a = Association {
1448 state: AssociationState::Established,
1449 peer_last_tsn: 0,
1450 my_next_tsn: 1,
1451 max_receive_buffer_size: 1024,
1452 max_receive_message_size: 1024,
1453 max_payload_size: 1200,
1454 use_forward_tsn: true,
1455 ..Default::default()
1456 };
1457 assert!(
1458 a.create_stream(stream_id, false, PayloadProtocolIdentifier::Binary)
1459 .is_some()
1460 );
1461 let reset: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
1462 reconfig_request_sequence_number: 7,
1463 reconfig_response_sequence_number: u32::MAX,
1464 sender_last_tsn: 1,
1465 stream_identifiers: vec![stream_id],
1466 });
1467 a.handle_reconfig_param(&reset, &mut vec![])?;
1468 
1469 // The reset has not reached its TSN boundary yet, so this skip belongs to
1470 // the old generation even though the association-level cumulative TSN
1471 // advances beyond that boundary.
1472 a.handle_forward_tsn(&ChunkForwardTsn {
1473 new_cumulative_tsn: 2,
1474 streams: vec![ChunkForwardTsnStream {
1475 identifier: stream_id,
1476 sequence: 0,
1477 }],
1478 })?;
1479 
1480 a.handle_data(&ChunkPayloadData {
1481 tsn: 3,
1482 stream_identifier: stream_id,
1483 stream_sequence_number: 0,
1484 beginning_fragment: true,
1485 ending_fragment: true,
1486 user_data: Bytes::from_static(b"new"),
1487 ..Default::default()
1488 })?;
1489 let successor = a.stream(stream_id)?.read()?.unwrap();
1490 let mut payload = [0; 3];
1491 assert_eq!(successor.read(&mut payload)?, payload.len());
1492 assert_eq!(&payload, b"new");
1493 Ok(())
1494}
1495 
1496fn association_with_queued_second_reset() -> Result<Association> {
1497 let mut a = association_with_retiring_boundary_data()?;
1498 let second_reset: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
1499 reconfig_request_sequence_number: 8,
1500 reconfig_response_sequence_number: u32::MAX,
1501 sender_last_tsn: 2,
1502 stream_identifiers: vec![1],
1503 });
1504 a.handle_reconfig_param(&second_reset, &mut vec![])?;
1505 assert!(a.reconfig_requests.contains_key(&8));
1506 Ok(a)
1507}
1508 
1509fn assert_generation_c_ssn_zero_is_readable(mut a: Association) -> Result<()> {
1510 let old = a.stream(1)?.read()?.unwrap();
1511 let mut payload = [0; 3];
1512 assert_eq!(old.read(&mut payload)?, payload.len());
1513 assert_eq!(&payload, b"old");
1514 
1515 a.handle_data(&ChunkPayloadData {
1516 tsn: 4,
1517 stream_identifier: 1,
1518 stream_sequence_number: 0,
1519 beginning_fragment: true,
1520 ending_fragment: true,
1521 user_data: Bytes::from_static(b"gen-c"),
1522 ..Default::default()
1523 })?;
1524 let generation_c = a.stream(1)?.read()?.unwrap();
1525 let mut payload = [0; 5];
1526 assert_eq!(generation_c.read(&mut payload)?, payload.len());
1527 assert_eq!(&payload, b"gen-c");
1528 Ok(())
1529}
1530 
1531#[test]
1532fn test_forward_tsn_skip_is_scoped_to_queued_reset_generation() -> Result<()> {
1533 let mut a = association_with_queued_second_reset()?;
1534 a.handle_forward_tsn(&ChunkForwardTsn {
1535 // TSN 2 is generation B's abandoned SSN 0. TSN 3 belongs to an
1536 // unrelated stream, so the aggregate cumulative point cannot identify
1537 // which stream generation the per-stream SSN describes.
1538 new_cumulative_tsn: 3,
1539 streams: vec![ChunkForwardTsnStream {
1540 identifier: 1,
1541 sequence: 0,
1542 }],
1543 })?;
1544 assert_generation_c_ssn_zero_is_readable(a)
1545}
1546 
1547#[test]
1548fn test_retransmitted_forward_tsn_stays_with_reset_generation() -> Result<()> {
1549 let mut a = association_with_queued_second_reset()?;
1550 
1551 // The first FORWARD-TSN advances through generation B's abandoned SSN 0
1552 // and lets its pending reset complete.
1553 a.handle_forward_tsn(&ChunkForwardTsn {
1554 new_cumulative_tsn: 2,
1555 streams: vec![ChunkForwardTsnStream {
1556 identifier: 1,
1557 sequence: 0,
1558 }],
1559 })?;
1560 assert!(!a.reconfig_requests.contains_key(&8));
1561 
1562 // If the resulting SACK is lost, the sender can legitimately repeat the
1563 // same stream skip while advancing over an unrelated abandoned TSN. The
1564 // repeated entry still belongs to generation B, not the future tail.
1565 a.handle_forward_tsn(&ChunkForwardTsn {
1566 new_cumulative_tsn: 3,
1567 streams: vec![ChunkForwardTsnStream {
1568 identifier: 1,
1569 sequence: 0,
1570 }],
1571 })?;
1572 
1573 assert_generation_c_ssn_zero_is_readable(a)
1574}
1575 
1576fn assert_successor_ssn_one_after_reused_forward_ssn(mut a: Association) -> Result<()> {
1577 let old = a.stream(1)?.read()?.unwrap();
1578 let mut payload = [0; 3];
1579 assert_eq!(old.read(&mut payload)?, payload.len());
1580 assert_eq!(&payload, b"old");
1581 
1582 a.handle_data(&ChunkPayloadData {
1583 tsn: 4,
1584 stream_identifier: 1,
1585 stream_sequence_number: 1,
1586 beginning_fragment: true,
1587 ending_fragment: true,
1588 user_data: Bytes::from_static(b"gen-c"),
1589 ..Default::default()
1590 })?;
1591 let generation_c = a
1592 .stream(1)?
1593 .read()?
1594 .expect("generation C SSN 1 should follow its abandoned SSN 0");
1595 let mut payload = [0; 5];
1596 assert_eq!(generation_c.read(&mut payload)?, payload.len());
1597 assert_eq!(&payload, b"gen-c");
1598 Ok(())
1599}
1600 
1601#[test]
1602fn test_successor_forward_tsn_can_reuse_reset_ssn() -> Result<()> {
1603 let mut a = association_with_queued_second_reset()?;
1604 for new_cumulative_tsn in [2, 3] {
1605 a.handle_forward_tsn(&ChunkForwardTsn {
1606 new_cumulative_tsn,
1607 streams: vec![ChunkForwardTsnStream {
1608 identifier: 1,
1609 sequence: 0,
1610 }],
1611 })?;
1612 }
1613 assert_successor_ssn_one_after_reused_forward_ssn(a)
1614}
1615 
1616#[test]
1617fn test_successor_i_forward_tsn_can_reuse_reset_mid() -> Result<()> {
1618 let mut a = association_with_queued_second_reset()?;
1619 for new_cumulative_tsn in [2, 3] {
1620 a.handle_i_forward_tsn(&ChunkIForwardTsn {
1621 new_cumulative_tsn,
1622 streams: vec![ChunkIForwardTsnStream {
1623 identifier: 1,
1624 unordered: false,
1625 mid: 0,
1626 }],
1627 })?;
1628 }
1629 assert_successor_ssn_one_after_reused_forward_ssn(a)
1630}
1631 
1632#[test]
1633fn test_repeated_old_forward_tsn_does_not_skip_later_successor_ssn() -> Result<()> {
1634 let mut a = association_with_queued_second_reset()?;
1635 for new_cumulative_tsn in [2, 3] {
1636 a.handle_forward_tsn(&ChunkForwardTsn {
1637 new_cumulative_tsn,
1638 streams: vec![ChunkForwardTsnStream {
1639 identifier: 1,
1640 sequence: 1,
1641 }],
1642 })?;
1643 }
1644 
1645 let old = a.stream(1)?.read()?.unwrap();
1646 let mut payload = [0; 3];
1647 assert_eq!(old.read(&mut payload)?, payload.len());
1648 
1649 for (tsn, ssn, bytes) in [(4, 0, b"zero".as_slice()), (5, 1, b"one".as_slice())] {
1650 a.handle_data(&ChunkPayloadData {
1651 tsn,
1652 stream_identifier: 1,
1653 stream_sequence_number: ssn,
1654 beginning_fragment: true,
1655 ending_fragment: true,
1656 user_data: Bytes::copy_from_slice(bytes),
1657 ..Default::default()
1658 })?;
1659 let message = a
1660 .stream(1)?
1661 .read()?
1662 .expect("an old-generation repeat must not skip a live successor SSN");
1663 let mut payload = [0; 4];
1664 assert_eq!(message.read(&mut payload)?, bytes.len());
1665 assert_eq!(&payload[..bytes.len()], bytes);
1666 }
1667 Ok(())
1668}
1669 
1670#[test]
1671fn test_repeated_old_forward_tsn_preserves_partial_successor_above_cumulative() -> Result<()> {
1672 let mut a = association_with_queued_second_reset()?;
1673 a.handle_forward_tsn(&ChunkForwardTsn {
1674 new_cumulative_tsn: 2,
1675 streams: vec![ChunkForwardTsnStream {
1676 identifier: 1,
1677 sequence: 0,
1678 }],
1679 })?;
1680 a.handle_data(&ChunkPayloadData {
1681 tsn: 4,
1682 stream_identifier: 1,
1683 stream_sequence_number: 0,
1684 beginning_fragment: true,
1685 ending_fragment: false,
1686 user_data: Bytes::from_static(b"live-"),
1687 ..Default::default()
1688 })?;
1689 
1690 // TSN 3 is unrelated to this stream. Repeating generation B's skip must
1691 // not discard generation C's fragment at TSN 4.
1692 a.handle_forward_tsn(&ChunkForwardTsn {
1693 new_cumulative_tsn: 3,
1694 streams: vec![ChunkForwardTsnStream {
1695 identifier: 1,
1696 sequence: 0,
1697 }],
1698 })?;
1699 
1700 let old = a.stream(1)?.read()?.unwrap();
1701 let mut payload = [0; 3];
1702 assert_eq!(old.read(&mut payload)?, payload.len());
1703 a.handle_data(&ChunkPayloadData {
1704 tsn: 5,
1705 stream_identifier: 1,
1706 stream_sequence_number: 0,
1707 beginning_fragment: false,
1708 ending_fragment: true,
1709 user_data: Bytes::from_static(b"data"),
1710 ..Default::default()
1711 })?;
1712 
1713 let message = a
1714 .stream(1)?
1715 .read()?
1716 .expect("the successor fragments should still reassemble");
1717 let mut payload = [0; 9];
1718 assert_eq!(message.read(&mut payload)?, payload.len());
1719 assert_eq!(&payload, b"live-data");
1720 Ok(())
1721}
1722 
1723#[test]
1724fn test_deferred_forward_tsn_discards_partial_missing_pre_cumulative_tsn() -> Result<()> {
1725 let mut a = association_with_retiring_boundary_data()?;
1726 let reset: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
1727 reconfig_request_sequence_number: 8,
1728 reconfig_response_sequence_number: u32::MAX,
1729 sender_last_tsn: 4,
1730 stream_identifiers: vec![1],
1731 });
1732 a.handle_reconfig_param(&reset, &mut vec![])?;
1733 
1734 // B/SSN 0 is missing its beginning at TSN 2, while B/SSN 1 is complete.
1735 for chunk in [
1736 ChunkPayloadData {
1737 tsn: 3,
1738 stream_identifier: 1,
1739 stream_sequence_number: 0,
1740 beginning_fragment: false,
1741 ending_fragment: true,
1742 user_data: Bytes::from_static(b"orphan"),
1743 ..Default::default()
1744 },
1745 ChunkPayloadData {
1746 tsn: 4,
1747 stream_identifier: 1,
1748 stream_sequence_number: 1,
1749 beginning_fragment: true,
1750 ending_fragment: true,
1751 user_data: Bytes::from_static(b"next"),
1752 ..Default::default()
1753 },
1754 ] {
1755 a.handle_data(&chunk)?;
1756 }
1757 a.handle_forward_tsn(&ChunkForwardTsn {
1758 new_cumulative_tsn: 2,
1759 streams: vec![ChunkForwardTsnStream {
1760 identifier: 1,
1761 sequence: 0,
1762 }],
1763 })?;
1764 
1765 assert_eq!(
1766 a.get_my_receiver_window_credit(),
1767 a.max_receive_buffer_size - 3 - 4,
1768 "the partial message missing TSN 2 must release its six bytes immediately"
1769 );
1770 let old = a.stream(1)?.read()?.unwrap();
1771 let mut payload = [0; 3];
1772 assert_eq!(old.read(&mut payload)?, payload.len());
1773 
1774 let next = a
1775 .stream(1)?
1776 .read()?
1777 .expect("discarding abandoned SSN 0 must unblock complete SSN 1");
1778 let mut payload = [0; 4];
1779 assert_eq!(next.read(&mut payload)?, payload.len());
1780 assert_eq!(&payload, b"next");
1781 Ok(())
1782}
1783 
1784#[test]
1785fn test_reordered_successor_waits_before_consuming_old_forward_tsn() -> Result<()> {
1786 let mut a = association_with_queued_second_reset()?;
1787 for new_cumulative_tsn in [2, 3] {
1788 a.handle_forward_tsn(&ChunkForwardTsn {
1789 new_cumulative_tsn,
1790 streams: vec![ChunkForwardTsnStream {
1791 identifier: 1,
1792 sequence: 0,
1793 }],
1794 })?;
1795 }
1796 a.handle_data(&ChunkPayloadData {
1797 tsn: 5,
1798 stream_identifier: 1,
1799 stream_sequence_number: 1,
1800 beginning_fragment: true,
1801 ending_fragment: true,
1802 user_data: Bytes::from_static(b"one"),
1803 ..Default::default()
1804 })?;
1805 
1806 let old = a.stream(1)?.read()?.unwrap();
1807 let mut payload = [0; 3];
1808 assert_eq!(old.read(&mut payload)?, payload.len());
1809 assert!(
1810 a.stream(1)?.read()?.is_none(),
1811 "SSN 1 must wait while SSN 0's post-FWD TSN is still missing"
1812 );
1813 
1814 a.handle_data(&ChunkPayloadData {
1815 tsn: 4,
1816 stream_identifier: 1,
1817 stream_sequence_number: 0,
1818 beginning_fragment: true,
1819 ending_fragment: true,
1820 user_data: Bytes::from_static(b"zero"),
1821 ..Default::default()
1822 })?;
1823 for expected in [b"zero".as_slice(), b"one".as_slice()] {
1824 let message = a
1825 .stream(1)?
1826 .read()?
1827 .expect("successor messages should become readable in SSN order");
1828 let mut payload = [0; 4];
1829 assert_eq!(message.read(&mut payload)?, expected.len());
1830 assert_eq!(&payload[..expected.len()], expected);
1831 }
1832 Ok(())
1833}
1834 
1835#[test]
1836fn test_reordered_successor_waits_before_consuming_repeated_old_forward_tsn() -> Result<()> {
1837 let mut a = association_with_queued_second_reset()?;
1838 for new_cumulative_tsn in [2, 3] {
1839 a.handle_forward_tsn(&ChunkForwardTsn {
1840 new_cumulative_tsn,
1841 streams: vec![ChunkForwardTsnStream {
1842 identifier: 1,
1843 sequence: 1,
1844 }],
1845 })?;
1846 }
1847 a.handle_data(&ChunkPayloadData {
1848 tsn: 5,
1849 stream_identifier: 1,
1850 stream_sequence_number: 1,
1851 beginning_fragment: true,
1852 ending_fragment: true,
1853 user_data: Bytes::from_static(b"one"),
1854 ..Default::default()
1855 })?;
1856 
1857 let old = a.stream(1)?.read()?.unwrap();
1858 let mut payload = [0; 3];
1859 assert_eq!(old.read(&mut payload)?, payload.len());
1860 assert!(
1861 a.stream(1)?.read()?.is_none(),
1862 "SSN 1 must wait while SSN 0's post-FWD TSN is still missing"
1863 );
1864 
1865 a.handle_data(&ChunkPayloadData {
1866 tsn: 4,
1867 stream_identifier: 1,
1868 stream_sequence_number: 0,
1869 beginning_fragment: true,
1870 ending_fragment: true,
1871 user_data: Bytes::from_static(b"zero"),
1872 ..Default::default()
1873 })?;
1874 for expected in [b"zero".as_slice(), b"one".as_slice()] {
1875 let message = a
1876 .stream(1)?
1877 .read()?
1878 .expect("successor messages should become readable in SSN order");
1879 let mut payload = [0; 4];
1880 assert_eq!(message.read(&mut payload)?, expected.len());
1881 assert_eq!(&payload[..expected.len()], expected);
1882 }
1883 Ok(())
1884}
1885 
1886#[test]
1887fn test_successor_forward_tsn_releases_covered_out_of_order_message() -> Result<()> {
1888 let mut a = association_with_queued_second_reset()?;
1889 a.handle_forward_tsn(&ChunkForwardTsn {
1890 new_cumulative_tsn: 2,
1891 streams: vec![ChunkForwardTsnStream {
1892 identifier: 1,
1893 sequence: 0,
1894 }],
1895 })?;
1896 a.handle_data(&ChunkPayloadData {
1897 tsn: 5,
1898 stream_identifier: 1,
1899 stream_sequence_number: 1,
1900 beginning_fragment: true,
1901 ending_fragment: true,
1902 user_data: Bytes::from_static(b"one"),
1903 ..Default::default()
1904 })?;
1905 
1906 // The sender did not receive the gap acknowledgement for TSN 5 and
1907 // abandons both missing SSN 0 and the already-complete SSN 1.
1908 a.handle_forward_tsn(&ChunkForwardTsn {
1909 new_cumulative_tsn: 5,
1910 streams: vec![ChunkForwardTsnStream {
1911 identifier: 1,
1912 sequence: 1,
1913 }],
1914 })?;
1915 
1916 let old = a.stream(1)?.read()?.unwrap();
1917 let mut payload = [0; 3];
1918 assert_eq!(old.read(&mut payload)?, payload.len());
1919 
1920 let successor = a
1921 .stream(1)?
1922 .read()?
1923 .expect("the successor skip must release the complete covered message");
1924 let mut payload = [0; 3];
1925 assert_eq!(successor.read(&mut payload)?, payload.len());
1926 assert_eq!(&payload, b"one");
1927 Ok(())
1928}
1929 
1930#[test]
1931fn test_cross_generation_forward_tsn_does_not_skip_uncovered_successor_ssn() -> Result<()> {
1932 let mut a = association_with_queued_second_reset()?;
1933 a.handle_forward_tsn(&ChunkForwardTsn {
1934 new_cumulative_tsn: 2,
1935 streams: vec![ChunkForwardTsnStream {
1936 identifier: 1,
1937 sequence: 1,
1938 }],
1939 })?;
1940 a.handle_data(&ChunkPayloadData {
1941 tsn: 4,
1942 stream_identifier: 1,
1943 stream_sequence_number: 0,
1944 beginning_fragment: true,
1945 ending_fragment: true,
1946 user_data: Bytes::from_static(b"zero"),
1947 ..Default::default()
1948 })?;
1949 
1950 // A lost SACK can make one FORWARD-TSN span old SSN 1 and successor SSN
1951 // 0. The old maximum SSN must not also skip successor SSN 1.
1952 a.handle_forward_tsn(&ChunkForwardTsn {
1953 new_cumulative_tsn: 4,
1954 streams: vec![ChunkForwardTsnStream {
1955 identifier: 1,
1956 sequence: 1,
1957 }],
1958 })?;
1959 
1960 let old = a.stream(1)?.read()?.unwrap();
1961 let mut payload = [0; 3];
1962 assert_eq!(old.read(&mut payload)?, payload.len());
1963 let zero = a.stream(1)?.read()?.unwrap();
1964 let mut payload = [0; 4];
1965 assert_eq!(zero.read(&mut payload)?, payload.len());
1966 assert_eq!(&payload, b"zero");
1967 
1968 a.handle_data(&ChunkPayloadData {
1969 tsn: 5,
1970 stream_identifier: 1,
1971 stream_sequence_number: 1,
1972 beginning_fragment: true,
1973 ending_fragment: true,
1974 user_data: Bytes::from_static(b"one"),
1975 ..Default::default()
1976 })?;
1977 let one = a
1978 .stream(1)?
1979 .read()?
1980 .expect("the successor SSN outside the covered TSN range must remain live");
1981 let mut payload = [0; 3];
1982 assert_eq!(one.read(&mut payload)?, payload.len());
1983 assert_eq!(&payload, b"one");
1984 Ok(())
1985}
1986 
1987#[test]
1988fn test_successor_forward_tsn_releases_covered_message_before_lost_tail() -> Result<()> {
1989 let mut a = association_with_queued_second_reset()?;
1990 a.handle_forward_tsn(&ChunkForwardTsn {
1991 new_cumulative_tsn: 2,
1992 streams: vec![ChunkForwardTsnStream {
1993 identifier: 1,
1994 sequence: 0,
1995 }],
1996 })?;
1997 a.handle_data(&ChunkPayloadData {
1998 tsn: 4,
1999 stream_identifier: 1,
2000 stream_sequence_number: 1,
2001 beginning_fragment: true,
2002 ending_fragment: true,
2003 user_data: Bytes::from_static(b"one"),
2004 ..Default::default()
2005 })?;
2006 
2007 // Successor SSN 0/TSN 3 and SSN 2/TSN 5 were both abandoned. The
2008 // complete covered SSN 1 can be released before SSN 2 is observed.
2009 a.handle_forward_tsn(&ChunkForwardTsn {
2010 new_cumulative_tsn: 5,
2011 streams: vec![ChunkForwardTsnStream {
2012 identifier: 1,
2013 sequence: 2,
2014 }],
2015 })?;
2016 
2017 let old = a.stream(1)?.read()?.unwrap();
2018 let mut payload = [0; 3];
2019 assert_eq!(old.read(&mut payload)?, payload.len());
2020 let one = a
2021 .stream(1)?
2022 .read()?
2023 .expect("covered successor data must not wait for the lost final SSN");
2024 let mut payload = [0; 3];
2025 assert_eq!(one.read(&mut payload)?, payload.len());
2026 assert_eq!(&payload, b"one");
2027 Ok(())
2028}
2029 
2030#[test]
2031fn test_resolved_successor_prefix_releases_message_before_stale_tail() -> Result<()> {
2032 let mut a = association_with_queued_second_reset()?;
2033 a.handle_forward_tsn(&ChunkForwardTsn {
2034 new_cumulative_tsn: 2,
2035 streams: vec![ChunkForwardTsnStream {
2036 identifier: 1,
2037 sequence: 2,
2038 }],
2039 })?;
2040 a.handle_forward_tsn(&ChunkForwardTsn {
2041 new_cumulative_tsn: 3,
2042 streams: vec![ChunkForwardTsnStream {
2043 identifier: 1,
2044 sequence: 2,
2045 }],
2046 })?;
2047 a.handle_data(&ChunkPayloadData {
2048 tsn: 4,
2049 stream_identifier: 1,
2050 stream_sequence_number: 1,
2051 beginning_fragment: true,
2052 ending_fragment: true,
2053 user_data: Bytes::from_static(b"one"),
2054 ..Default::default()
2055 })?;
2056 
2057 let old = a.stream(1)?.read()?.unwrap();
2058 let mut payload = [0; 3];
2059 assert_eq!(old.read(&mut payload)?, payload.len());
2060 let one = a
2061 .stream(1)?
2062 .read()?
2063 .expect("resolved successor data must not wait for the stale old tail");
2064 let mut payload = [0; 3];
2065 assert_eq!(one.read(&mut payload)?, payload.len());
2066 assert_eq!(&payload, b"one");
2067 Ok(())
2068}
2069 
2070fn assert_covered_partial_after_unaffected_anchor_is_discarded(last_ssn: u16) -> Result<()> {
2071 let mut a = association_with_queued_second_reset()?;
2072 for new_cumulative_tsn in [2, 3] {
2073 a.handle_forward_tsn(&ChunkForwardTsn {
2074 new_cumulative_tsn,
2075 streams: vec![ChunkForwardTsnStream {
2076 identifier: 1,
2077 sequence: last_ssn,
2078 }],
2079 })?;
2080 }
2081 for chunk in [
2082 ChunkPayloadData {
2083 tsn: 4,
2084 stream_identifier: 1,
2085 stream_sequence_number: 1,
2086 beginning_fragment: false,
2087 ending_fragment: true,
2088 user_data: Bytes::from_static(b"orphan"),
2089 ..Default::default()
2090 },
2091 ChunkPayloadData {
2092 tsn: 5,
2093 stream_identifier: 1,
2094 stream_sequence_number: 0,
2095 beginning_fragment: true,
2096 ending_fragment: true,
2097 user_data: Bytes::from_static(b"zero"),
2098 ..Default::default()
2099 },
2100 ] {
2101 a.handle_data(&chunk)?;
2102 }
2103 
2104 let old = a.stream(1)?.read()?.unwrap();
2105 let mut payload = [0; 3];
2106 assert_eq!(old.read(&mut payload)?, payload.len());
2107 let zero = a.stream(1)?.read()?.unwrap();
2108 let mut payload = [0; 4];
2109 assert_eq!(zero.read(&mut payload)?, payload.len());
2110 assert_eq!(&payload, b"zero");
2111 assert_eq!(
2112 a.get_my_receiver_window_credit(),
2113 a.max_receive_buffer_size,
2114 "the covered partial message must release its receive-window credit"
2115 );
2116 Ok(())
2117}
2118 
2119#[test]
2120fn test_full_forward_update_discards_covered_partial_after_anchor() -> Result<()> {
2121 assert_covered_partial_after_unaffected_anchor_is_discarded(1)
2122}
2123 
2124#[test]
2125fn test_partial_forward_update_discards_covered_partial_after_anchor() -> Result<()> {
2126 assert_covered_partial_after_unaffected_anchor_is_discarded(2)
2127}
2128 
2129#[test]
2130fn test_drained_generation_discards_retained_forward_suffix() -> Result<()> {
2131 let mut a = association_with_retiring_boundary_data()?;
2132 let reset: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
2133 reconfig_request_sequence_number: 8,
2134 reconfig_response_sequence_number: u32::MAX,
2135 sender_last_tsn: 5,
2136 stream_identifiers: vec![1],
2137 });
2138 a.handle_reconfig_param(&reset, &mut vec![])?;
2139 a.handle_data(&ChunkPayloadData {
2140 tsn: 4,
2141 stream_identifier: 1,
2142 stream_sequence_number: 1,
2143 beginning_fragment: true,
2144 ending_fragment: true,
2145 user_data: Bytes::from_static(b"one"),
2146 ..Default::default()
2147 })?;
2148 a.handle_forward_tsn(&ChunkForwardTsn {
2149 new_cumulative_tsn: 5,
2150 streams: vec![ChunkForwardTsnStream {
2151 identifier: 1,
2152 sequence: 2,
2153 }],
2154 })?;
2155 
2156 let old = a.stream(1)?.read()?.unwrap();
2157 let mut payload = [0; 3];
2158 assert_eq!(old.read(&mut payload)?, payload.len());
2159 let one = a.stream(1)?.read()?.unwrap();
2160 let mut payload = [0; 3];
2161 assert_eq!(one.read(&mut payload)?, payload.len());
2162 assert_eq!(&payload, b"one");
2163 assert!(
2164 a.deferred_forward_tsns.get(&1).is_none_or(|updates| {
2165 updates
2166 .iter()
2167 .all(|update| update.generation_boundary != Some(5))
2168 }),
2169 "a drained reset generation must discard its retained ambiguous suffix"
2170 );
2171 Ok(())
2172}
2173 
2174#[test]
2175fn test_first_reset_discards_retained_tail_forward_suffix() -> Result<()> {
2176 let mut a = association_with_queued_second_reset()?;
2177 a.handle_forward_tsn(&ChunkForwardTsn {
2178 new_cumulative_tsn: 2,
2179 streams: vec![ChunkForwardTsnStream {
2180 identifier: 1,
2181 sequence: 0,
2182 }],
2183 })?;
2184 a.handle_data(&ChunkPayloadData {
2185 tsn: 4,
2186 stream_identifier: 1,
2187 stream_sequence_number: 1,
2188 beginning_fragment: true,
2189 ending_fragment: true,
2190 user_data: Bytes::from_static(b"one"),
2191 ..Default::default()
2192 })?;
2193 a.handle_forward_tsn(&ChunkForwardTsn {
2194 new_cumulative_tsn: 5,
2195 streams: vec![ChunkForwardTsnStream {
2196 identifier: 1,
2197 sequence: 2,
2198 }],
2199 })?;
2200 
2201 let old = a.stream(1)?.read()?.unwrap();
2202 let mut payload = [0; 3];
2203 assert_eq!(old.read(&mut payload)?, payload.len());
2204 let one = a.stream(1)?.read()?.unwrap();
2205 let mut payload = [0; 3];
2206 assert_eq!(one.read(&mut payload)?, payload.len());
2207 assert_eq!(&payload, b"one");
2208 assert!(
2209 a.deferred_forward_tsns[&1]
2210 .iter()
2211 .any(|update| update.generation_boundary.is_none())
2212 );
2213 
2214 let reset: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
2215 reconfig_request_sequence_number: 9,
2216 reconfig_response_sequence_number: u32::MAX,
2217 sender_last_tsn: 5,
2218 stream_identifiers: vec![1],
2219 });
2220 a.handle_reconfig_param(&reset, &mut vec![])?;
2221 assert!(
2222 !a.deferred_forward_tsns.contains_key(&1),
2223 "resetting a drained tail must discard its retained ambiguous suffix"
2224 );
2225 Ok(())
2226}
2227 
2228#[test]
2229fn test_i_forward_tsn_skip_is_scoped_to_queued_reset_generation() -> Result<()> {
2230 let mut a = association_with_queued_second_reset()?;
2231 a.handle_i_forward_tsn(&ChunkIForwardTsn {
2232 new_cumulative_tsn: 3,
2233 streams: vec![ChunkIForwardTsnStream {
2234 identifier: 1,
2235 unordered: false,
2236 mid: 0,
2237 }],
2238 })?;
2239 assert_generation_c_ssn_zero_is_readable(a)
2240}
2241 
2242fn assert_later_reset_claims_tail_forward_tsn(mut a: Association) -> Result<()> {
2243 let second_reset: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
2244 reconfig_request_sequence_number: 8,
2245 reconfig_response_sequence_number: u32::MAX,
2246 sender_last_tsn: 2,
2247 stream_identifiers: vec![1],
2248 });
2249 a.handle_reconfig_param(&second_reset, &mut vec![])?;
2250 
2251 let old = a.stream(1)?.read()?.unwrap();
2252 let mut payload = [0; 3];
2253 assert_eq!(old.read(&mut payload)?, payload.len());
2254 assert_eq!(&payload, b"old");
2255 
2256 a.handle_data(&ChunkPayloadData {
2257 tsn: 3,
2258 stream_identifier: 1,
2259 stream_sequence_number: 0,
2260 beginning_fragment: true,
2261 ending_fragment: true,
2262 user_data: Bytes::from_static(b"gen-c"),
2263 ..Default::default()
2264 })?;
2265 let generation_c = a.stream(1)?.read()?.unwrap();
2266 let mut payload = [0; 5];
2267 assert_eq!(generation_c.read(&mut payload)?, payload.len());
2268 assert_eq!(&payload, b"gen-c");
2269 Ok(())
2270}
2271 
2272#[test]
2273fn test_later_reset_claims_tail_forward_tsn_generation() -> Result<()> {
2274 let mut a = association_with_retiring_boundary_data()?;
2275 a.handle_forward_tsn(&ChunkForwardTsn {
2276 new_cumulative_tsn: 2,
2277 streams: vec![ChunkForwardTsnStream {
2278 identifier: 1,
2279 sequence: 0,
2280 }],
2281 })?;
2282 assert_later_reset_claims_tail_forward_tsn(a)
2283}
2284 
2285#[test]
2286fn test_later_reset_claims_tail_i_forward_tsn_generation() -> Result<()> {
2287 let mut a = association_with_retiring_boundary_data()?;
2288 a.handle_i_forward_tsn(&ChunkIForwardTsn {
2289 new_cumulative_tsn: 2,
2290 streams: vec![ChunkIForwardTsnStream {
2291 identifier: 1,
2292 unordered: false,
2293 mid: 0,
2294 }],
2295 })?;
2296 assert_later_reset_claims_tail_forward_tsn(a)
2297}
2298 
2299fn association_with_complete_deferred_successor() -> Result<Association> {
2300 let mut a = association_with_retiring_boundary_data()?;
2301 a.handle_data(&ChunkPayloadData {
2302 tsn: 2,
2303 stream_identifier: 1,
2304 stream_sequence_number: 0,
2305 beginning_fragment: true,
2306 ending_fragment: true,
2307 user_data: Bytes::from_static(b"received"),
2308 ..Default::default()
2309 })?;
2310 Ok(a)
2311}
2312 
2313fn assert_received_message_survives_forward_tsn(mut a: Association) -> Result<()> {
2314 let old = a.stream(1)?.read()?.unwrap();
2315 let mut payload = [0; 3];
2316 assert_eq!(old.read(&mut payload)?, payload.len());
2317 assert_eq!(&payload, b"old");
2318 
2319 let received = a.stream(1)?.read()?.unwrap();
2320 let mut payload = [0; 8];
2321 assert_eq!(received.read(&mut payload)?, payload.len());
2322 assert_eq!(&payload, b"received");
2323 Ok(())
2324}
2325 
2326#[test]
2327fn test_forward_tsn_preserves_complete_deferred_message() -> Result<()> {
2328 let mut a = association_with_complete_deferred_successor()?;
2329 a.handle_forward_tsn(&ChunkForwardTsn {
2330 // SSN 0 was received completely; only the following SSN 1 was
2331 // abandoned. RFC 3758 requires the stranded complete message to be
2332 // made available when the skip advances the ordered stream.
2333 new_cumulative_tsn: 3,
2334 streams: vec![ChunkForwardTsnStream {
2335 identifier: 1,
2336 sequence: 1,
2337 }],
2338 })?;
2339 assert_received_message_survives_forward_tsn(a)
2340}
2341 
2342#[test]
2343fn test_forward_tsn_preserves_out_of_order_complete_deferred_message() -> Result<()> {
2344 let mut a = association_with_retiring_boundary_data()?;
2345 
2346 // Generation B's complete SSN 0 arrives above a TSN gap, so it remains in
2347 // both the association payload queue and the reset-generation holding map.
2348 a.handle_data(&ChunkPayloadData {
2349 tsn: 3,
2350 stream_identifier: 1,
2351 stream_sequence_number: 0,
2352 beginning_fragment: true,
2353 ending_fragment: true,
2354 user_data: Bytes::from_static(b"received"),
2355 ..Default::default()
2356 })?;
2357 
2358 // The sender may not have received the gap SACK before abandoning TSNs 2
2359 // and 3. Advancing the cumulative TSN must not erase a complete message
2360 // that this receiver already holds for the successor generation.
2361 a.handle_forward_tsn(&ChunkForwardTsn {
2362 new_cumulative_tsn: 3,
2363 streams: vec![ChunkForwardTsnStream {
2364 identifier: 1,
2365 sequence: 0,
2366 }],
2367 })?;
2368 
2369 assert_received_message_survives_forward_tsn(a)
2370}
2371 
2372#[test]
2373fn test_i_forward_tsn_preserves_complete_deferred_message() -> Result<()> {
2374 let mut a = association_with_complete_deferred_successor()?;
2375 a.handle_i_forward_tsn(&ChunkIForwardTsn {
2376 new_cumulative_tsn: 3,
2377 streams: vec![ChunkIForwardTsnStream {
2378 identifier: 1,
2379 unordered: false,
2380 mid: 1,
2381 }],
2382 })?;
2383 assert_received_message_survives_forward_tsn(a)
2384}
2385 
2386fn association_with_deferred_unordered_fragment() -> Result<Association> {
2387 let mut a = association_with_retiring_boundary_data()?;
2388 a.streams.get_mut(&1).unwrap().unordered = true;
2389 a.handle_data(&ChunkPayloadData {
2390 tsn: 2,
2391 stream_identifier: 1,
2392 unordered: true,
2393 beginning_fragment: true,
2394 ending_fragment: false,
2395 user_data: Bytes::from_static(b"orphan"),
2396 ..Default::default()
2397 })?;
2398 Ok(a)
2399}
2400 
2401fn assert_abandoned_unordered_fragment_is_discarded(mut a: Association) -> Result<()> {
2402 let old = a.stream(1)?.read()?.unwrap();
2403 let mut payload = [0; 3];
2404 assert_eq!(old.read(&mut payload)?, payload.len());
2405 assert_eq!(&payload, b"old");
2406 
2407 assert_eq!(
2408 a.streams
2409 .get(&1)
2410 .map(StreamState::get_num_bytes_in_reassembly_queue)
2411 .unwrap_or_default(),
2412 0,
2413 "an abandoned partial unordered message must not survive replay"
2414 );
2415 assert_eq!(a.get_my_receiver_window_credit(), a.max_receive_buffer_size);
2416 Ok(())
2417}
2418 
2419#[test]
2420fn test_forward_tsn_discards_deferred_unordered_fragment() -> Result<()> {
2421 let mut a = association_with_deferred_unordered_fragment()?;
2422 a.handle_forward_tsn(&ChunkForwardTsn {
2423 new_cumulative_tsn: 3,
2424 streams: vec![],
2425 })?;
2426 assert_abandoned_unordered_fragment_is_discarded(a)
2427}
2428 
2429#[test]
2430fn test_i_forward_tsn_discards_deferred_unordered_fragment() -> Result<()> {
2431 let mut a = association_with_deferred_unordered_fragment()?;
2432 a.handle_i_forward_tsn(&ChunkIForwardTsn {
2433 new_cumulative_tsn: 3,
2434 streams: vec![ChunkIForwardTsnStream {
2435 identifier: 1,
2436 unordered: true,
2437 mid: 0,
2438 }],
2439 })?;
2440 assert_abandoned_unordered_fragment_is_discarded(a)
2441}
2442 
2443#[test]
2444fn test_forward_tsn_releases_abandoned_deferred_fragment_credit_immediately() -> Result<()> {
2445 let mut a = association_with_deferred_unordered_fragment()?;
2446 let old_generation_bytes = a
2447 .streams
2448 .get(&1)
2449 .unwrap()
2450 .get_num_bytes_in_reassembly_queue() as u32;
2451 
2452 a.handle_forward_tsn(&ChunkForwardTsn {
2453 new_cumulative_tsn: 3,
2454 streams: vec![],
2455 })?;
2456 
2457 assert_eq!(
2458 a.get_my_receiver_window_credit(),
2459 a.max_receive_buffer_size - old_generation_bytes,
2460 "abandoned successor fragments must not hold receive-window credit until the old generation drains"
2461 );
2462 Ok(())
2463}
2464 
2465#[test]
2466fn test_deferred_unordered_forward_tsn_updates_are_coalesced() -> Result<()> {
2467 let mut a = association_with_retiring_boundary_data()?;
2468 a.streams.get_mut(&1).unwrap().unordered = true;
2469 
2470 for new_cumulative_tsn in 2..=9 {
2471 a.handle_forward_tsn(&ChunkForwardTsn {
2472 new_cumulative_tsn,
2473 streams: vec![],
2474 })?;
2475 }
2476 
2477 let updates = a.deferred_forward_tsns.get(&1).unwrap();
2478 assert_eq!(
2479 updates.len(),
2480 1,
2481 "newer unordered cumulative TSNs supersede older updates for the same generation"
2482 );
2483 assert!(matches!(
2484 updates.front(),
2485 Some(DeferredForwardTsn {
2486 kind: DeferredForwardTsnKind::Unordered {
2487 new_cumulative_tsn: 9
2488 },
2489 ..
2490 })
2491 ));
2492 Ok(())
2493}
2494 
2495#[test]
2496fn test_deferred_tail_ordered_forward_tsn_duplicates_are_coalesced() -> Result<()> {
2497 let mut a = association_with_retiring_boundary_data()?;
2498 for new_cumulative_tsn in 2..=9 {
2499 a.handle_forward_tsn(&ChunkForwardTsn {
2500 new_cumulative_tsn,
2501 streams: vec![ChunkForwardTsnStream {
2502 identifier: 1,
2503 sequence: 0,
2504 }],
2505 })?;
2506 }
2507 
2508 let ordered: Vec<_> = a.deferred_forward_tsns[&1]
2509 .iter()
2510 .filter(|update| {
2511 update.generation_boundary.is_none()
2512 && matches!(update.kind, DeferredForwardTsnKind::Ordered { .. })
2513 })
2514 .collect();
2515 assert_eq!(
2516 ordered.len(),
2517 1,
2518 "lost-SACK repeats for an ambiguous tail SSN should coalesce"
2519 );
2520 Ok(())
2521}
2522 
2523#[test]
2524fn test_deferred_bounded_ordered_forward_tsn_updates_are_coalesced() -> Result<()> {
2525 let mut a = association_with_retiring_boundary_data()?;
2526 let reset: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
2527 reconfig_request_sequence_number: 8,
2528 reconfig_response_sequence_number: u32::MAX,
2529 sender_last_tsn: 10,
2530 stream_identifiers: vec![1],
2531 });
2532 a.handle_reconfig_param(&reset, &mut vec![])?;
2533 
2534 for new_cumulative_tsn in 2..=9 {
2535 a.handle_forward_tsn(&ChunkForwardTsn {
2536 new_cumulative_tsn,
2537 streams: vec![ChunkForwardTsnStream {
2538 identifier: 1,
2539 sequence: (new_cumulative_tsn - 2) as u16,
2540 }],
2541 })?;
2542 }
2543 
2544 let ordered: Vec<_> = a.deferred_forward_tsns[&1]
2545 .iter()
2546 .filter(|update| {
2547 update.generation_boundary == Some(10)
2548 && matches!(update.kind, DeferredForwardTsnKind::Ordered { .. })
2549 })
2550 .collect();
2551 assert_eq!(
2552 ordered.len(),
2553 1,
2554 "a newer ordered skip supersedes older skips for the same reset boundary"
2555 );
2556 assert!(matches!(
2557 ordered[0].kind,
2558 DeferredForwardTsnKind::Ordered { last_ssn: 7, .. }
2559 ));
2560 Ok(())
2561}
2562 
2563#[test]
2564fn test_unread_retiring_stream_does_not_block_unrelated_events() -> Result<()> {
2565 let mut a = association_with_retiring_boundary_data()?;
2566 assert!(
2567 a.create_stream(2, false, PayloadProtocolIdentifier::Binary)
2568 .is_some()
2569 );
2570 
2571 assert!(
2572 core::iter::from_fn(|| a.poll())
2573 .any(|event| matches!(event, Event::Stream(StreamEvent::Readable { id: 1 })))
2574 );
2575 assert!(a.poll().is_none(), "stream 1 Finished should still wait");
2576 
2577 a.handle_data(&ChunkPayloadData {
2578 tsn: 2,
2579 stream_identifier: 2,
2580 stream_sequence_number: 0,
2581 beginning_fragment: true,
2582 ending_fragment: true,
2583 user_data: Bytes::from_static(b"other"),
2584 ..Default::default()
2585 })?;
2586 
2587 assert!(matches!(a.poll(), Some(Event::DatagramReceived)));
2588 assert!(matches!(
2589 a.poll(),
2590 Some(Event::Stream(StreamEvent::Readable { id: 2 }))
2591 ));
2592 Ok(())
2593}
2594 
2595#[test]
2596fn test_stop_discards_unread_retiring_data_and_unblocks_finished() -> Result<()> {
2597 let mut a = association_with_retiring_boundary_data()?;
2598 assert!(
2599 core::iter::from_fn(|| a.poll())
2600 .any(|event| matches!(event, Event::Stream(StreamEvent::Readable { id: 1 })))
2601 );
2602 assert!(a.poll().is_none(), "Finished should wait for old DATA");
2603 
2604 a.stream(1)?.stop()?;
2605 assert!(matches!(
2606 a.poll(),
2607 Some(Event::Stream(StreamEvent::Finished { id: 1 }))
2608 ));
2609 Ok(())
2610}
2611 
2612#[test]
2613fn test_stop_during_retirement_does_not_queue_duplicate_reset() -> Result<()> {
2614 let mut a = association_with_retiring_boundary_data()?;
2615 assert_eq!(
2616 a.reconfigs.len(),
2617 1,
2618 "peer reset should queue one reciprocal"
2619 );
2620 assert!(a.pending_reset_streams.is_empty());
2621 
2622 a.stream(1)?.stop()?;
2623 
2624 assert!(
2625 a.pending_reset_streams.is_empty(),
2626 "the reciprocal already resets this outgoing stream"
2627 );
2628 assert_eq!(a.reconfigs.len(), 1);
2629 Ok(())
2630}
2631 
2632fn assert_shutdown_during_retirement_allows_reuse(close: bool) -> Result<()> {
2633 let mut a = association_with_retiring_boundary_data()?;
2634 if close {
2635 a.stream(1)?.close()?;
2636 } else {
2637 a.stream(1)?.stop()?;
2638 }
2639 
2640 let rsn = *a.reconfigs.keys().next().unwrap();
2641 a.active_reconfig = Some(rsn);
2642 let response: Box<dyn Param + Send + Sync> = Box::new(ParamReconfigResponse {
2643 reconfig_response_sequence_number: rsn,
2644 result: ReconfigResult::SuccessPerformed,
2645 });
2646 a.handle_reconfig_param(&response, &mut vec![])?;
2647 
2648 assert!(
2649 !a.streams.contains_key(&1),
2650 "a completed reciprocal reset must remove the locally shut down stream"
2651 );
2652 assert!(
2653 a.open_stream(1, PayloadProtocolIdentifier::Binary).is_ok(),
2654 "ResetComplete must make the retired stream id reusable"
2655 );
2656 Ok(())
2657}
2658 
2659#[test]
2660fn test_stop_during_retirement_allows_reuse_after_reset_complete() -> Result<()> {
2661 assert_shutdown_during_retirement_allows_reuse(false)
2662}
2663 
2664#[test]
2665fn test_close_during_retirement_allows_reuse_after_reset_complete() -> Result<()> {
2666 assert_shutdown_during_retirement_allows_reuse(true)
2667}
2668 
2669fn assert_reset_complete_before_shutdown_allows_reuse(close: bool) -> Result<()> {
2670 let mut a = association_with_retiring_boundary_data()?;
2671 let rsn = *a.reconfigs.keys().next().unwrap();
2672 a.active_reconfig = Some(rsn);
2673 let response: Box<dyn Param + Send + Sync> = Box::new(ParamReconfigResponse {
2674 reconfig_response_sequence_number: rsn,
2675 result: ReconfigResult::SuccessPerformed,
2676 });
2677 a.handle_reconfig_param(&response, &mut vec![])?;
2678 
2679 if close {
2680 a.stream(1)?.close()?;
2681 } else {
2682 a.stream(1)?.stop()?;
2683 }
2684 
2685 assert!(
2686 a.pending_reset_streams.is_empty() && a.reconfigs.is_empty(),
2687 "shutting down a completed retiring stream must not queue a duplicate reset"
2688 );
2689 assert!(
2690 !a.streams.contains_key(&1),
2691 "discarding the retired generation after ResetComplete must remove the stream"
2692 );
2693 assert!(
2694 a.open_stream(1, PayloadProtocolIdentifier::Binary).is_ok(),
2695 "the completed stream id should be immediately reusable"
2696 );
2697 Ok(())
2698}
2699 
2700#[test]
2701fn test_reset_complete_before_stop_allows_reuse() -> Result<()> {
2702 assert_reset_complete_before_shutdown_allows_reuse(false)
2703}
2704 
2705#[test]
2706fn test_reset_complete_before_close_allows_reuse() -> Result<()> {
2707 assert_reset_complete_before_shutdown_allows_reuse(true)
2708}
2709 
2710#[test]
2711fn test_stop_does_not_reopen_read_half_on_deferred_successor() -> Result<()> {
2712 let mut a = association_with_retiring_boundary_data()?;
2713 a.handle_data(&ChunkPayloadData {
2714 tsn: 2,
2715 stream_identifier: 1,
2716 stream_sequence_number: 0,
2717 beginning_fragment: true,
2718 ending_fragment: true,
2719 user_data: Bytes::from_static(b"successor"),
2720 ..Default::default()
2721 })?;
2722 
2723 let mut stream = a.stream(1)?;
2724 stream.stop()?;
2725 assert_eq!(stream.read().unwrap_err(), Error::ErrStreamClosed);
2726 Ok(())
2727}
2728 
2729#[test]
2730fn test_stop_discards_deferred_successor_bytes_and_notifications() -> Result<()> {
2731 let mut a = association_with_retiring_boundary_data()?;
2732 while a.poll().is_some() {}
2733 a.handle_data(&ChunkPayloadData {
2734 tsn: 2,
2735 stream_identifier: 1,
2736 stream_sequence_number: 0,
2737 beginning_fragment: true,
2738 ending_fragment: true,
2739 user_data: Bytes::from_static(b"successor"),
2740 ..Default::default()
2741 })?;
2742 
2743 a.stream(1)?.stop()?;
2744 
2745 assert_eq!(
2746 a.streams
2747 .get(&1)
2748 .unwrap()
2749 .get_num_bytes_in_reassembly_queue(),
2750 0,
2751 "stop must not strand bytes in a write-only successor"
2752 );
2753 assert_eq!(a.get_my_receiver_window_credit(), a.max_receive_buffer_size);
2754 assert!(
2755 !core::iter::from_fn(|| a.poll()).any(|event| matches!(
2756 event,
2757 Event::Stream(StreamEvent::Opened { id } | StreamEvent::Readable { id }) if id == 1
2758 )),
2759 "stop must not advertise an unreadable successor"
2760 );
2761 Ok(())
2762}
2763 
2764#[test]
2765fn test_finish_remains_closed_across_deferred_successor() -> Result<()> {
2766 let mut a = association_with_retiring_boundary_data()?;
2767 a.handle_data(&ChunkPayloadData {
2768 tsn: 2,
2769 stream_identifier: 1,
2770 stream_sequence_number: 0,
2771 beginning_fragment: true,
2772 ending_fragment: true,
2773 user_data: Bytes::from_static(b"successor"),
2774 ..Default::default()
2775 })?;
2776 a.stream(1)?.finish()?;
2777 
2778 let old = a.stream(1)?.read()?.unwrap();
2779 let mut payload = [0; 3];
2780 assert_eq!(old.read(&mut payload)?, payload.len());
2781 assert_eq!(&payload, b"old");
2782 assert_eq!(a.streams.get(&1).unwrap().state, RecvSendState::Readable);
2783 
2784 let successor = a.stream(1)?.read()?.unwrap();
2785 let mut payload = [0; 9];
2786 assert_eq!(successor.read(&mut payload)?, payload.len());
2787 assert_eq!(&payload, b"successor");
2788 Ok(())
2789}
2790 
2791#[test]
2792fn test_close_remains_closed_across_deferred_successor() -> Result<()> {
2793 let mut a = association_with_retiring_boundary_data()?;
2794 a.handle_data(&ChunkPayloadData {
2795 tsn: 2,
2796 stream_identifier: 1,
2797 stream_sequence_number: 0,
2798 beginning_fragment: true,
2799 ending_fragment: true,
2800 user_data: Bytes::from_static(b"successor"),
2801 ..Default::default()
2802 })?;
2803 
2804 a.stream(1)?.close()?;
2805 
2806 let successor = a.streams.get(&1).unwrap();
2807 assert_eq!(successor.state, RecvSendState::Closed);
2808 assert_eq!(successor.get_num_bytes_in_reassembly_queue(), 0);
2809 assert_eq!(a.get_my_receiver_window_credit(), a.max_receive_buffer_size);
2810 Ok(())
2811}
2812 
2813#[test]
2814fn test_duplicate_stream_ids_retire_only_one_generation() -> Result<()> {
2815 let stream_id = 1;
2816 let mut a = Association {
2817 state: AssociationState::Established,
2818 peer_last_tsn: 0,
2819 my_next_tsn: 1,
2820 max_receive_buffer_size: 1024,
2821 max_receive_message_size: 1024,
2822 max_payload_size: 1200,
2823 ..Default::default()
2824 };
2825 assert!(
2826 a.create_stream(stream_id, false, PayloadProtocolIdentifier::Binary)
2827 .is_some()
2828 );
2829 a.handle_data(&ChunkPayloadData {
2830 tsn: 1,
2831 stream_identifier: stream_id,
2832 stream_sequence_number: 0,
2833 beginning_fragment: true,
2834 ending_fragment: true,
2835 user_data: Bytes::from_static(b"old"),
2836 ..Default::default()
2837 })?;
2838 
2839 let reset: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
2840 reconfig_request_sequence_number: 7,
2841 reconfig_response_sequence_number: u32::MAX,
2842 sender_last_tsn: 1,
2843 stream_identifiers: vec![stream_id, stream_id],
2844 });
2845 a.handle_reconfig_param(&reset, &mut vec![])?;
2846 
2847 assert_eq!(a.retiring_streams.get(&stream_id).unwrap().len(), 1);
2848 assert_eq!(a.pending_reset_completions.0.get(&stream_id), Some(&1));
2849 assert_eq!(
2850 a.reconfigs
2851 .values()
2852 .flat_map(Association::reconfig_stream_ids)
2853 .collect::<Vec<_>>(),
2854 vec![stream_id]
2855 );
2856 Ok(())
2857}
2858 
2859#[test]
2860fn test_reconfig_rejects_two_outgoing_resets_without_mutation() -> Result<()> {
2861 let mut a = association_with_retiring_boundary_data()?;
2862 let boundaries_before = a.retiring_streams.get(&1).unwrap().len();
2863 let completions_before = a.pending_reset_completions.0.get(&1).copied();
2864 let reconfigs_before = a.reconfigs.len();
2865 
2866 let reconfig = ChunkReconfig {
2867 param_a: Some(Box::new(ParamOutgoingResetRequest {
2868 reconfig_request_sequence_number: 8,
2869 reconfig_response_sequence_number: u32::MAX,
2870 sender_last_tsn: a.peer_last_tsn,
2871 stream_identifiers: vec![1],
2872 })),
2873 param_b: Some(Box::new(ParamOutgoingResetRequest {
2874 reconfig_request_sequence_number: 9,
2875 reconfig_response_sequence_number: u32::MAX,
2876 sender_last_tsn: a.peer_last_tsn,
2877 stream_identifiers: vec![1],
2878 })),
2879 };
2880 let packet = a.create_packet(vec![Box::new(reconfig)]);
2881 let result = a.handle_chunk(&packet, &packet.chunks[0], Instant::now());
2882 
2883 assert!(result.is_err(), "two Outgoing Reset parameters are invalid");
2884 assert_eq!(a.retiring_streams.get(&1).unwrap().len(), boundaries_before);
2885 assert_eq!(
2886 a.pending_reset_completions.0.get(&1).copied(),
2887 completions_before
2888 );
2889 assert_eq!(a.reconfigs.len(), reconfigs_before);
2890 Ok(())
2891}
2892 
2893#[test]
2894fn test_stop_before_reset_boundary_does_not_leave_finished_blocked() -> Result<()> {
2895 let stream_id = 1;
2896 let mut a = Association {
2897 state: AssociationState::Established,
2898 peer_last_tsn: 0,
2899 my_next_tsn: 1,
2900 max_receive_buffer_size: 1024,
2901 max_receive_message_size: 1024,
2902 max_payload_size: 1200,
2903 ..Default::default()
2904 };
2905 for id in [stream_id, 2] {
2906 assert!(
2907 a.create_stream(id, false, PayloadProtocolIdentifier::Binary)
2908 .is_some()
2909 );
2910 }
2911 
2912 // TSN 2 is readable but cannot advance the cumulative point past missing
2913 // TSN 1, leaving time for the application to stop the read half.
2914 a.handle_data(&ChunkPayloadData {
2915 tsn: 2,
2916 stream_identifier: stream_id,
2917 stream_sequence_number: 0,
2918 beginning_fragment: true,
2919 ending_fragment: true,
2920 user_data: Bytes::from_static(b"discard"),
2921 ..Default::default()
2922 })?;
2923 let reset: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
2924 reconfig_request_sequence_number: 7,
2925 reconfig_response_sequence_number: u32::MAX,
2926 sender_last_tsn: 2,
2927 stream_identifiers: vec![stream_id],
2928 });
2929 a.handle_reconfig_param(&reset, &mut vec![])?;
2930 a.stream(stream_id)?.stop()?;
2931 
2932 a.handle_data(&ChunkPayloadData {
2933 tsn: 1,
2934 stream_identifier: 2,
2935 stream_sequence_number: 0,
2936 beginning_fragment: true,
2937 ending_fragment: true,
2938 user_data: Bytes::from_static(b"gap"),
2939 ..Default::default()
2940 })?;
2941 
2942 assert!(a.stream(stream_id).is_err());
2943 assert!(
2944 core::iter::from_fn(|| a.poll()).any(|event| matches!(
2945 event,
2946 Event::Stream(StreamEvent::Finished { id }) if id == stream_id
2947 )),
2948 "a stopped read half cannot be left waiting for an impossible drain"
2949 );
2950 Ok(())
2951}
2952 
2953#[test]
2954fn test_oversized_deferred_successor_is_rejected_before_old_read() -> Result<()> {
2955 let mut a = association_with_retiring_boundary_data()?;
2956 a.max_receive_message_size = 3;
2957 
2958 let error = a
2959 .handle_data(&ChunkPayloadData {
2960 tsn: 2,
2961 stream_identifier: 1,
2962 stream_sequence_number: 0,
2963 beginning_fragment: true,
2964 ending_fragment: true,
2965 user_data: Bytes::from_static(b"too-big"),
2966 ..Default::default()
2967 })
2968 .unwrap_err();
2969 assert_eq!(error, Error::ErrInboundPacketTooLarge);
2970 
2971 let old = a.stream(1)?.read()?.unwrap();
2972 let mut payload = [0; 3];
2973 assert_eq!(old.read(&mut payload)?, payload.len());
2974 assert_eq!(&payload, b"old");
2975 Ok(())
2976}
2977 
2978#[test]
2979fn test_future_peer_reconfig_sequence_is_rejected() -> Result<()> {
2980 let stream_id = 1;
2981 let mut a = Association {
2982 peer_last_tsn: 0,
2983 peer_last_reconfig_rsn: 40,
2984 peer_reconfig_rsn_initialized: true,
2985 my_next_tsn: 1,
2986 ..Default::default()
2987 };
2988 assert!(
2989 a.create_stream(stream_id, false, PayloadProtocolIdentifier::Binary)
2990 .is_some()
2991 );
2992 
2993 let future: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
2994 reconfig_request_sequence_number: 42,
2995 reconfig_response_sequence_number: u32::MAX,
2996 sender_last_tsn: a.peer_last_tsn,
2997 stream_identifiers: vec![stream_id],
2998 });
2999 let mut reply = vec![];
3000 a.handle_reconfig_param(&future, &mut reply)?;
3001 
3002 assert_eq!(
3003 reconfig_response_result(&reply, 42),
3004 Some(ReconfigResult::ErrorBadSequenceNumber)
3005 );
3006 assert_eq!(a.peer_last_reconfig_rsn, 40);
3007 assert!(a.max_completed_reconfig_rsn.is_none());
3008 assert!(a.streams.contains_key(&stream_id));
3009 
3010 let expected: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
3011 reconfig_request_sequence_number: 41,
3012 reconfig_response_sequence_number: u32::MAX,
3013 sender_last_tsn: a.peer_last_tsn,
3014 stream_identifiers: vec![stream_id],
3015 });
3016 a.handle_reconfig_param(&expected, &mut vec![])?;
3017 assert!(!a.streams.contains_key(&stream_id));
3018 assert_eq!(a.max_completed_reconfig_rsn, Some(41));
3019 Ok(())
3020}
3021 
3022#[test]
3023fn test_peer_reconfig_sequence_accepts_wrap_to_zero() -> Result<()> {
3024 let stream_id = 1;
3025 let mut a = Association {
3026 peer_last_tsn: 0,
3027 peer_last_reconfig_rsn: u32::MAX,
3028 peer_reconfig_rsn_initialized: true,
3029 my_next_tsn: 1,
3030 ..Default::default()
3031 };
3032 assert!(
3033 a.create_stream(stream_id, false, PayloadProtocolIdentifier::Binary)
3034 .is_some()
3035 );
3036 
3037 let request: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
3038 reconfig_request_sequence_number: 0,
3039 reconfig_response_sequence_number: u32::MAX,
3040 sender_last_tsn: a.peer_last_tsn,
3041 stream_identifiers: vec![stream_id],
3042 });
3043 let mut reply = vec![];
3044 a.handle_reconfig_param(&request, &mut reply)?;
3045 
3046 assert_ne!(
3047 reconfig_response_result(&reply, 0),
3048 Some(ReconfigResult::ErrorBadSequenceNumber)
3049 );
3050 assert_eq!(a.peer_last_reconfig_rsn, 0);
3051 assert_eq!(a.max_completed_reconfig_rsn, Some(0));
3052 assert!(!a.streams.contains_key(&stream_id));
3053 Ok(())
3054}
3055 
3056#[test]
3057fn test_empty_reset_stream_list_resets_all_streams() -> Result<()> {
3058 let mut a = Association {
3059 my_next_tsn: 1,
3060 ..Default::default()
3061 };
3062 for stream_id in [1, 2] {
3063 assert!(
3064 a.create_stream(stream_id, false, PayloadProtocolIdentifier::Binary)
3065 .is_some()
3066 );
3067 }
3068 
3069 let request: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
3070 reconfig_request_sequence_number: 7,
3071 reconfig_response_sequence_number: u32::MAX,
3072 sender_last_tsn: a.peer_last_tsn,
3073 stream_identifiers: vec![],
3074 });
3075 let mut reply = vec![];
3076 a.handle_reconfig_param(&request, &mut reply)?;
3077 
3078 assert!(
3079 a.streams.is_empty(),
3080 "an omitted stream list means all streams"
3081 );
3082 let reciprocal_ids = a
3083 .reconfigs
3084 .values()
3085 .next()
3086 .map(Association::reconfig_stream_ids)
3087 .unwrap_or_default();
3088 assert!(
3089 reciprocal_ids.is_empty(),
3090 "reset-all must remain compact on the wire"
3091 );
3092 let mut completion_ids = a
3093 .reconfig_reset_streams
3094 .values()
3095 .next()
3096 .cloned()
3097 .unwrap_or_default();
3098 completion_ids.sort_unstable();
3099 assert_eq!(completion_ids, vec![1, 2]);
3100 Ok(())
3101}
3102 
3103#[test]
3104fn test_reset_all_reciprocal_fits_reconfig_chunk_length() -> Result<()> {
3105 let mut a = Association {
3106 my_next_tsn: 1,
3107 ..Default::default()
3108 };
3109 for stream_id in 0..32_758u32 {
3110 assert!(
3111 a.create_stream(
3112 stream_id as StreamId,
3113 false,
3114 PayloadProtocolIdentifier::Binary,
3115 )
3116 .is_some()
3117 );
3118 }
3119 
3120 let request: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest {
3121 reconfig_request_sequence_number: 7,
3122 reconfig_response_sequence_number: u32::MAX,
3123 sender_last_tsn: a.peer_last_tsn,
3124 stream_identifiers: vec![],
3125 });
3126 let mut reply = vec![];
3127 a.handle_reconfig_param(&request, &mut reply)?;
3128 
3129 for packet in reply {
3130 packet.marshal()?;
3131 }
3132 Ok(())
3133}
3134 
3135#[test]
3136fn test_blocked_reconfig_does_not_allow_later_rsn_to_overtake() {
3137 let mut a = Association {
3138 state: AssociationState::Established,
3139 my_next_tsn: 1,
3140 cwnd: 0,
3141 rwnd: 0,
3142 mtu: 1400,
3143 ..Default::default()
3144 };
3145 a.pending_queue.push(ChunkPayloadData {
3146 stream_identifier: 1,
3147 beginning_fragment: true,
3148 ending_fragment: true,
3149 user_data: Bytes::from_static(b"pending"),
3150 ..Default::default()
3151 });
3152 a.inflight_queue.push_no_check(ChunkPayloadData {
3153 tsn: 0,
3154 stream_identifier: 2,
3155 user_data: Bytes::from_static(b"inflight"),
3156 ..Default::default()
3157 });
3158 insert_queued_reset(&mut a, 7, 1);
3159 insert_queued_reset(&mut a, 8, 2);
3160 
3161 let _ = a.gather_outbound(Instant::now());
3162 assert!(
3163 a.active_reconfig.is_none(),
3164 "RSN 8 must not overtake RSN 7 while RSN 7 waits for its DATA boundary"
3165 );
3166 assert_eq!(a.control_queue.len(), 2);
3167 
3168 a.cwnd = 1400;
3169 a.rwnd = 1400;
3170 let _ = a.gather_outbound(Instant::now());
3171 assert_eq!(a.active_reconfig, Some(7));
3172 assert_eq!(a.control_queue.len(), 1);
3173}
3174 
3175#[test]
3176fn test_retransmission_does_not_send_buffered_reconfigs() {
3177 let mut a = Association {
3178 state: AssociationState::Established,
3179 ..Default::default()
3180 };
3181 
3182 insert_active_reset(&mut a, 7, 1);
3183 insert_queued_reset(&mut a, 8, 2);
3184 a.timers
3185 .start(Timer::Reconfig, Instant::now(), a.rto_mgr.get_rto());
3186 a.will_retransmit_reconfig = true;
3187 
3188 let (packets, _) = a.gather_outbound(Instant::now());
3189 assert_eq!(
3190 packets.len(),
3191 1,
3192 "retransmission must include only the request associated with the running timer"
3193 );
3194}
3195 
3196#[test]
3197fn test_older_completion_does_not_make_newer_generation_writable() -> Result<()> {
3198 let stream_id = 1;
3199 let mut a = Association::default();
3200 
3201 a.pending_reset_completions.insert(stream_id);
3202 a.pending_reset_completions.insert(stream_id);
3203 insert_active_reset(&mut a, 7, stream_id);
3204 insert_queued_reset(&mut a, 8, stream_id);
3205 a.timers
3206 .start(Timer::Reconfig, Instant::now(), a.rto_mgr.get_rto());
3207 
3208 assert!(a.get_or_create_stream(stream_id).is_some());
3209 assert!(!a.stream(stream_id)?.is_writable());
3210 
3211 let response: Box<dyn Param + Send + Sync> = Box::new(ParamReconfigResponse {
3212 reconfig_response_sequence_number: 7,
3213 result: ReconfigResult::SuccessPerformed,
3214 });
3215 a.handle_reconfig_param(&response, &mut vec![])?;
3216 
3217 assert!(a.reconfigs.contains_key(&8));
3218 assert!(
3219 !a.stream(stream_id)?.is_writable(),
3220 "an older success cannot release the newer generation's outgoing quarantine"
3221 );
3222 Ok(())
3223}
3224 
3225#[test]
3226fn test_stale_completion_does_not_override_newer_denial() -> Result<()> {
3227 let stream_id = 1;
3228 let mut a = Association::default();
3229 
3230 a.pending_reset_completions.insert(stream_id);
3231 insert_active_reset(&mut a, 7, stream_id);
3232 
3233 let success: Box<dyn Param + Send + Sync> = Box::new(ParamReconfigResponse {
3234 reconfig_response_sequence_number: 7,
3235 result: ReconfigResult::SuccessPerformed,
3236 });
3237 a.handle_reconfig_param(&success, &mut vec![])?;
3238 
3239 // Before the application polls generation 7's completion, the peer opens
3240 // and uses generation 8, then resets it. Its nonzero outgoing sequence
3241 // means a denied reciprocal reset must quarantine the id.
3242 assert!(a.get_or_create_stream(stream_id).is_some());
3243 a.streams.get_mut(&stream_id).unwrap().sequence_number = 1;
3244 a.unregister_stream(stream_id, true);
3245 insert_active_reset(&mut a, 8, stream_id);
3246 a.timers
3247 .start(Timer::Reconfig, Instant::now(), a.rto_mgr.get_rto());
3248 
3249 let denied: Box<dyn Param + Send + Sync> = Box::new(ParamReconfigResponse {
3250 reconfig_response_sequence_number: 8,
3251 result: ReconfigResult::Denied,
3252 });
3253 a.handle_reconfig_param(&denied, &mut vec![])?;
3254 
3255 // Polling the stale completion must not resurrect reuse permission after
3256 // generation 8 has already failed.
3257 assert!(matches!(
3258 a.poll(),
3259 Some(Event::Stream(StreamEvent::ResetComplete { id })) if id == stream_id
3260 ));
3261 
3262 assert!(matches!(
3263 a.open_stream(stream_id, PayloadProtocolIdentifier::Binary),
3264 Err(Error::ErrStreamResetPending)
3265 ));
3266 Ok(())
3267}
3268 
3269#[test]
3270fn test_create_forward_tsn_forward_one_abandoned() -> Result<()> {
3271 let mut a = Association {
3272 cumulative_tsn_ack_point: 9,
3273 advanced_peer_tsn_ack_point: 10,
3274 ..Default::default()
3275 };
3276 
3277 a.inflight_queue.push_no_check(ChunkPayloadData {
3278 beginning_fragment: true,
3279 ending_fragment: true,
3280 tsn: 10,
3281 stream_identifier: 1,
3282 stream_sequence_number: 2,
3283 user_data: Bytes::from_static(b"ABC"),
3284 nsent: 1,
3285 abandoned: true,
3286 ..Default::default()
3287 });
3288 
3289 let fwdtsn = a.create_forward_tsn();
3290 
3291 assert_eq!(10, fwdtsn.new_cumulative_tsn, "should be able to serialize");
3292 assert_eq!(1, fwdtsn.streams.len(), "there should be one stream");
3293 assert_eq!(1, fwdtsn.streams[0].identifier, "si should be 1");
3294 assert_eq!(2, fwdtsn.streams[0].sequence, "ssn should be 2");
3295 
3296 Ok(())
3297}
3298 
3299#[test]
3300fn test_create_forward_tsn_forward_two_abandoned_with_the_same_si() -> Result<()> {
3301 let mut a = Association {
3302 cumulative_tsn_ack_point: 9,
3303 advanced_peer_tsn_ack_point: 12,
3304 ..Default::default()
3305 };
3306 
3307 a.inflight_queue.push_no_check(ChunkPayloadData {
3308 beginning_fragment: true,
3309 ending_fragment: true,
3310 tsn: 10,
3311 stream_identifier: 1,
3312 stream_sequence_number: 2,
3313 user_data: Bytes::from_static(b"ABC"),
3314 nsent: 1,
3315 abandoned: true,
3316 ..Default::default()
3317 });
3318 a.inflight_queue.push_no_check(ChunkPayloadData {
3319 beginning_fragment: true,
3320 ending_fragment: true,
3321 tsn: 11,
3322 stream_identifier: 1,
3323 stream_sequence_number: 3,
3324 user_data: Bytes::from_static(b"DEF"),
3325 nsent: 1,
3326 abandoned: true,
3327 ..Default::default()
3328 });
3329 a.inflight_queue.push_no_check(ChunkPayloadData {
3330 beginning_fragment: true,
3331 ending_fragment: true,
3332 tsn: 12,
3333 stream_identifier: 2,
3334 stream_sequence_number: 1,
3335 user_data: Bytes::from_static(b"123"),
3336 nsent: 1,
3337 abandoned: true,
3338 ..Default::default()
3339 });
3340 
3341 let fwdtsn = a.create_forward_tsn();
3342 
3343 assert_eq!(12, fwdtsn.new_cumulative_tsn, "should be able to serialize");
3344 assert_eq!(2, fwdtsn.streams.len(), "there should be two stream");
3345 
3346 let mut si1ok = false;
3347 let mut si2ok = false;
3348 for s in &fwdtsn.streams {
3349 match s.identifier {
3350 1 => {
3351 assert_eq!(3, s.sequence, "ssn should be 3");
3352 si1ok = true;
3353 }
3354 2 => {
3355 assert_eq!(1, s.sequence, "ssn should be 1");
3356 si2ok = true;
3357 }
3358 _ => panic!("unexpected stream indentifier"),
3359 }
3360 }
3361 assert!(si1ok, "si=1 should be present");
3362 assert!(si2ok, "si=2 should be present");
3363 
3364 Ok(())
3365}
3366 
3367#[test]
3368fn test_handle_forward_tsn_forward_3unreceived_chunks() -> Result<()> {
3369 let mut a = Association {
3370 use_forward_tsn: true,
3371 ..Default::default()
3372 };
3373 
3374 let prev_tsn = a.peer_last_tsn;
3375 
3376 let fwdtsn = ChunkForwardTsn {
3377 new_cumulative_tsn: a.peer_last_tsn + 3,
3378 streams: vec![ChunkForwardTsnStream {
3379 identifier: 0,
3380 sequence: 0,
3381 }],
3382 };
3383 
3384 let p = a.handle_forward_tsn(&fwdtsn)?;
3385 
3386 let delayed_ack_triggered = a.delayed_ack_triggered;
3387 let immediate_ack_triggered = a.immediate_ack_triggered;
3388 assert_eq!(
3389 a.peer_last_tsn,
3390 prev_tsn + 3,
3391 "peerLastTSN should advance by 3 "
3392 );
3393 assert!(delayed_ack_triggered, "delayed sack should be triggered");
3394 assert!(
3395 !immediate_ack_triggered,
3396 "immediate sack should NOT be triggered"
3397 );
3398 assert!(p.is_empty(), "should return empty");
3399 
3400 Ok(())
3401}
3402 
3403#[test]
3404fn test_handle_forward_tsn_forward_1for1_missing() -> Result<()> {
3405 let mut a = Association {
3406 use_forward_tsn: true,
3407 ..Default::default()
3408 };
3409 
3410 let prev_tsn = a.peer_last_tsn;
3411 
3412 // this chunk is blocked by the missing chunk at tsn=1
3413 a.payload_queue.push(
3414 ChunkPayloadData {
3415 beginning_fragment: true,
3416 ending_fragment: true,
3417 tsn: a.peer_last_tsn + 2,
3418 stream_identifier: 0,
3419 stream_sequence_number: 1,
3420 user_data: Bytes::from_static(b"ABC"),
3421 ..Default::default()
3422 },
3423 a.peer_last_tsn,
3424 );
3425 
3426 let fwdtsn = ChunkForwardTsn {
3427 new_cumulative_tsn: a.peer_last_tsn + 1,
3428 streams: vec![ChunkForwardTsnStream {
3429 identifier: 0,
3430 sequence: 1,
3431 }],
3432 };
3433 
3434 let p = a.handle_forward_tsn(&fwdtsn)?;
3435 
3436 let delayed_ack_triggered = a.delayed_ack_triggered;
3437 let immediate_ack_triggered = a.immediate_ack_triggered;
3438 assert_eq!(
3439 a.peer_last_tsn,
3440 prev_tsn + 2,
3441 "peerLastTSN should advance by 2"
3442 );
3443 assert!(delayed_ack_triggered, "delayed sack should be triggered");
3444 assert!(
3445 !immediate_ack_triggered,
3446 "immediate sack should NOT be triggered"
3447 );
3448 assert!(p.is_empty(), "should return empty");
3449 
3450 Ok(())
3451}
3452 
3453#[test]
3454fn test_handle_forward_tsn_forward_1for2_missing() -> Result<()> {
3455 let mut a = Association {
3456 use_forward_tsn: true,
3457 ..Default::default()
3458 };
3459 
3460 a.use_forward_tsn = true;
3461 let prev_tsn = a.peer_last_tsn;
3462 
3463 // this chunk is blocked by the missing chunk at tsn=1
3464 a.payload_queue.push(
3465 ChunkPayloadData {
3466 beginning_fragment: true,
3467 ending_fragment: true,
3468 tsn: a.peer_last_tsn + 3,
3469 stream_identifier: 0,
3470 stream_sequence_number: 1,
3471 user_data: Bytes::from_static(b"ABC"),
3472 ..Default::default()
3473 },
3474 a.peer_last_tsn,
3475 );
3476 
3477 let fwdtsn = ChunkForwardTsn {
3478 new_cumulative_tsn: a.peer_last_tsn + 1,
3479 streams: vec![ChunkForwardTsnStream {
3480 identifier: 0,
3481 sequence: 1,
3482 }],
3483 };
3484 
3485 let p = a.handle_forward_tsn(&fwdtsn)?;
3486 
3487 let immediate_ack_triggered = a.immediate_ack_triggered;
3488 assert_eq!(
3489 a.peer_last_tsn,
3490 prev_tsn + 1,
3491 "peerLastTSN should advance by 1"
3492 );
3493 assert!(
3494 immediate_ack_triggered,
3495 "immediate sack should be triggered"
3496 );
3497 assert!(p.is_empty(), "should return empty");
3498 
3499 Ok(())
3500}
3501 
3502#[test]
3503fn test_handle_forward_tsn_dup_forward_tsn_chunk_should_generate_sack() -> Result<()> {
3504 let mut a = Association {
3505 use_forward_tsn: true,
3506 ..Default::default()
3507 };
3508 
3509 let prev_tsn = a.peer_last_tsn;
3510 
3511 let fwdtsn = ChunkForwardTsn {
3512 new_cumulative_tsn: a.peer_last_tsn,
3513 streams: vec![ChunkForwardTsnStream {
3514 identifier: 0,
3515 sequence: 1,
3516 }],
3517 };
3518 
3519 let p = a.handle_forward_tsn(&fwdtsn)?;
3520 
3521 let ack_state = a.ack_state;
3522 assert_eq!(a.peer_last_tsn, prev_tsn, "peerLastTSN should not advance");
3523 assert_eq!(AckState::Immediate, ack_state, "sack should be requested");
3524 assert!(p.is_empty(), "should return empty");
3525 
3526 Ok(())
3527}
3528 
3529fn queued_chunk(tsn: u32) -> ChunkPayloadData {
3530 ChunkPayloadData {
3531 beginning_fragment: true,
3532 ending_fragment: true,
3533 tsn,
3534 user_data: Bytes::from_static(b"ABC"),
3535 ..Default::default()
3536 }
3537}
3538 
3539/// A peer names `new_cumulative_tsn` freely, up to ~2^31 ahead in serial-number
3540/// arithmetic. Advancing one TSN per iteration turned that into a multi-second
3541/// hang per chunk; the cost must follow the queue, not the size of the jump.
3542#[test]
3543fn test_handle_forward_tsn_huge_gap_is_bounded() -> Result<()> {
3544 let mut a = Association {
3545 use_forward_tsn: true,
3546 ..Default::default()
3547 };
3548 let gap: u32 = 1 << 30;
3549 
3550 // One chunk the jump abandons, one beyond it that must survive.
3551 assert!(a.payload_queue.push(queued_chunk(5), 0));
3552 assert!(a.payload_queue.push(queued_chunk(gap + 5), 0));
3553 
3554 let start = Instant::now();
3555 a.handle_forward_tsn(&ChunkForwardTsn {
3556 new_cumulative_tsn: gap,
3557 streams: vec![],
3558 })?;
3559 let elapsed = start.elapsed();
3560 
3561 assert!(
3562 elapsed < Duration::from_secs(2),
3563 "advance must be bounded by queued data, took {elapsed:?}"
3564 );
3565 assert_eq!(a.peer_last_tsn, gap, "should reach the new cumulative TSN");
3566 assert!(a.payload_queue.get(5).is_none(), "abandoned chunk dropped");
3567 assert!(a.payload_queue.get(gap + 5).is_some(), "later chunk kept");
3568 assert_eq!(a.payload_queue.get_num_bytes(), 3, "bytes follow the drop");
3569 
3570 Ok(())
3571}
3572 
3573/// The I-FORWARD-TSN handler (RFC 8260) carries its own copy of the advance.
3574#[test]
3575fn test_handle_i_forward_tsn_huge_gap_is_bounded() -> Result<()> {
3576 let mut a = Association {
3577 use_forward_tsn: true,
3578 ..Default::default()
3579 };
3580 let gap: u32 = 1 << 30;
3581 
3582 let start = Instant::now();
3583 a.handle_i_forward_tsn(&ChunkIForwardTsn {
3584 new_cumulative_tsn: gap,
3585 streams: vec![],
3586 })?;
3587 let elapsed = start.elapsed();
3588 
3589 assert!(
3590 elapsed < Duration::from_secs(2),
3591 "advance must be bounded by queued data, took {elapsed:?}"
3592 );
3593 assert_eq!(a.peer_last_tsn, gap, "should reach the new cumulative TSN");
3594 
3595 Ok(())
3596}
3597 
3598#[test]
3599fn test_assoc_create_new_stream() -> Result<()> {
3600 let mut a = Association::default();
3601 
3602 for i in 0..ACCEPT_CH_SIZE {
3603 let stream_identifier =
3604 if let Some(s) = a.create_stream(i as u16, true, PayloadProtocolIdentifier::Unknown) {
3605 s.stream_identifier
3606 } else {
3607 panic!("{} should success", i);
3608 };
3609 let result = a.streams.get(&stream_identifier);
3610 assert!(result.is_some(), "should be in a.streams map");
3611 }
3612 
3613 let new_si = ACCEPT_CH_SIZE as u16;
3614 let result = a.streams.get(&new_si);
3615 assert!(result.is_none(), "should NOT be in a.streams map");
3616 
3617 let to_be_ignored = ChunkPayloadData {
3618 beginning_fragment: true,
3619 ending_fragment: true,
3620 tsn: a.peer_last_tsn + 1,
3621 stream_identifier: new_si,
3622 user_data: Bytes::from_static(b"ABC"),
3623 ..Default::default()
3624 };
3625 
3626 let p = a.handle_data(&to_be_ignored)?;
3627 assert!(p.is_empty(), "should return empty");
3628 
3629 Ok(())
3630}
3631 
3632fn handle_init_test(name: &str, initial_state: AssociationState, expect_err: bool) {
3633 let mut a = create_association(TransportConfig::default());
3634 a.set_state(initial_state);
3635 let pkt = Packet {
3636 common_header: CommonHeader {
3637 source_port: 5001,
3638 destination_port: 5002,
3639 ..Default::default()
3640 },
3641 ..Default::default()
3642 };
3643 let mut init = ChunkInit {
3644 initial_tsn: 1234,
3645 num_outbound_streams: 1001,
3646 num_inbound_streams: 1002,
3647 initiate_tag: 5678,
3648 advertised_receiver_window_credit: 512 * 1024,
3649 ..Default::default()
3650 };
3651 init.set_supported_extensions();
3652 
3653 let result = a.handle_init(&pkt, &init);
3654 if expect_err {
3655 assert!(result.is_err(), "{} should fail", name);
3656 return;
3657 } else {
3658 assert!(result.is_ok(), "{} should be ok", name);
3659 }
3660 assert_eq!(
3661 if init.initial_tsn == 0 {
3662 u32::MAX
3663 } else {
3664 init.initial_tsn - 1
3665 },
3666 a.peer_last_tsn,
3667 "{} should match",
3668 name
3669 );
3670 assert_eq!(1001, a.my_max_num_outbound_streams, "{} should match", name);
3671 assert_eq!(1002, a.my_max_num_inbound_streams, "{} should match", name);
3672 assert_eq!(5678, a.peer_verification_tag, "{} should match", name);
3673 assert_eq!(
3674 pkt.common_header.source_port, a.destination_port,
3675 "{} should match",
3676 name
3677 );
3678 assert_eq!(
3679 pkt.common_header.destination_port, a.source_port,
3680 "{} should match",
3681 name
3682 );
3683 assert!(a.use_forward_tsn, "{} should be set to true", name);
3684 assert_eq!(
3685 512 * 1024,
3686 a.rwnd,
3687 "{} rwnd should be initialized from peer's advertised_receiver_window_credit",
3688 name
3689 );
3690 assert_eq!(
3691 a.rwnd, a.ssthresh,
3692 "{} ssthresh should be initialized to rwnd",
3693 name
3694 );
3695}
3696 
3697#[test]
3698fn test_assoc_handle_init() -> Result<()> {
3699 handle_init_test("normal", AssociationState::Closed, false);
3700 
3701 handle_init_test(
3702 "unexpected state established",
3703 AssociationState::Established,
3704 true,
3705 );
3706 
3707 handle_init_test(
3708 "unexpected state shutdownAckSent",
3709 AssociationState::ShutdownAckSent,
3710 true,
3711 );
3712 
3713 handle_init_test(
3714 "unexpected state shutdownPending",
3715 AssociationState::ShutdownPending,
3716 true,
3717 );
3718 
3719 handle_init_test(
3720 "unexpected state shutdownReceived",
3721 AssociationState::ShutdownReceived,
3722 true,
3723 );
3724 
3725 handle_init_test(
3726 "unexpected state shutdownSent",
3727 AssociationState::ShutdownSent,
3728 true,
3729 );
3730 
3731 Ok(())
3732}
3733 
3734#[test]
3735fn test_assoc_max_send_message_size_default() -> Result<()> {
3736 let mut a = create_association(TransportConfig::default());
3737 assert_eq!(65536, a.max_send_message_size, "should match");
3738 
3739 let ppi = PayloadProtocolIdentifier::Unknown;
3740 let stream = a.create_stream(1, false, ppi);
3741 assert!(stream.is_some(), "should succeed");
3742 
3743 if let Some(mut s) = stream {
3744 let p = Bytes::from(vec![0u8; 65537]);
3745 
3746 if let Err(err) = s.write_sctp(&p.slice(..65536), ppi) {
3747 assert_ne!(
3748 Error::ErrOutboundPacketTooLarge,
3749 err,
3750 "should be not Error::ErrOutboundPacketTooLarge"
3751 );
3752 } else {
3753 panic!("should be error");
3754 }
3755 
3756 if let Err(err) = s.write_sctp(&p.slice(..65537), ppi) {
3757 assert_eq!(
3758 Error::ErrOutboundPacketTooLarge,
3759 err,
3760 "should be Error::ErrOutboundPacketTooLarge"
3761 );
3762 } else {
3763 panic!("should be error");
3764 }
3765 }
3766 
3767 Ok(())
3768}
3769 
3770#[test]
3771fn test_assoc_max_send_message_size_explicit() -> Result<()> {
3772 let mut a = create_association(TransportConfig::default().with_max_send_message_size(30000));
3773 assert_eq!(30000, a.max_send_message_size, "should match");
3774 
3775 let ppi = PayloadProtocolIdentifier::Unknown;
3776 let stream = a.create_stream(1, false, ppi);
3777 assert!(stream.is_some(), "should succeed");
3778 
3779 if let Some(mut s) = stream {
3780 let p = Bytes::from(vec![0u8; 30001]);
3781 
3782 if let Err(err) = s.write_sctp(&p.slice(..30000), ppi) {
3783 assert_ne!(
3784 Error::ErrOutboundPacketTooLarge,
3785 err,
3786 "should be not Error::ErrOutboundPacketTooLarge"
3787 );
3788 } else {
3789 panic!("should be error");
3790 }
3791 
3792 if let Err(err) = s.write_sctp(&p.slice(..30001), ppi) {
3793 assert_eq!(
3794 Error::ErrOutboundPacketTooLarge,
3795 err,
3796 "should be Error::ErrOutboundPacketTooLarge"
3797 );
3798 } else {
3799 panic!("should be error");
3800 }
3801 }
3802 
3803 Ok(())
3804}
3805 
3806#[test]
3807fn test_assoc_max_receive_message_size_default() -> Result<()> {
3808 let mut a = create_association(TransportConfig::default());
3809 assert_eq!(65536, a.max_receive_message_size, "should match");
3810 
3811 let p = Bytes::from(vec![0u8; 65537]);
3812 
3813 let size_ok = ChunkPayloadData {
3814 beginning_fragment: true,
3815 ending_fragment: true,
3816 tsn: a.peer_last_tsn + 1,
3817 stream_identifier: 1,
3818 user_data: p.slice(..65536),
3819 ..Default::default()
3820 };
3821 
3822 assert!(a.handle_data(&size_ok).is_ok(), "should succeed");
3823 
3824 let too_large = ChunkPayloadData {
3825 beginning_fragment: true,
3826 ending_fragment: true,
3827 tsn: a.peer_last_tsn + 1,
3828 stream_identifier: 1,
3829 user_data: p,
3830 ..Default::default()
3831 };
3832 
3833 if let Err(err) = a.handle_data(&too_large) {
3834 assert_eq!(
3835 Error::ErrInboundPacketTooLarge,
3836 err,
3837 "should be Error::ErrInboundPacketTooLarge"
3838 );
3839 } else {
3840 panic!("should be error");
3841 }
3842 
3843 Ok(())
3844}
3845 
3846#[test]
3847fn test_assoc_max_receive_message_size_explicit() -> Result<()> {
3848 let mut a = create_association(TransportConfig::default().with_max_receive_message_size(1024));
3849 assert_eq!(1024, a.max_receive_message_size, "should match");
3850 
3851 let first_chunk = ChunkPayloadData {
3852 beginning_fragment: true,
3853 ending_fragment: true,
3854 tsn: a.peer_last_tsn + 1,
3855 stream_identifier: 1,
3856 user_data: Bytes::from(vec![0u8; 512]),
3857 ..Default::default()
3858 };
3859 
3860 assert!(a.handle_data(&first_chunk).is_ok(), "should succeed");
3861 
3862 let second_chunk = ChunkPayloadData {
3863 beginning_fragment: true,
3864 ending_fragment: true,
3865 tsn: a.peer_last_tsn + 1,
3866 stream_identifier: 1,
3867 user_data: Bytes::from(vec![0u8; 513]),
3868 ..Default::default()
3869 };
3870 
3871 if let Err(err) = a.handle_data(&second_chunk) {
3872 assert_eq!(
3873 Error::ErrInboundPacketTooLarge,
3874 err,
3875 "should be Error::ErrInboundPacketTooLarge"
3876 );
3877 } else {
3878 panic!("should be error");
3879 }
3880 
3881 Ok(())
3882}
3883 
3884#[test]
3885fn test_assoc_max_message_size_asymmetric() -> Result<()> {
3886 let config = TransportConfig::default()
3887 .with_max_send_message_size(1024)
3888 .with_max_receive_message_size(30000);
3889 
3890 let mut a = create_association(config);
3891 assert_eq!(1024, a.max_send_message_size, "should match");
3892 assert_eq!(30000, a.max_receive_message_size, "should match");
3893 
3894 let ppi = PayloadProtocolIdentifier::Unknown;
3895 let stream = a.create_stream(1, false, ppi);
3896 assert!(stream.is_some(), "should succeed");
3897 
3898 if let Some(mut s) = stream {
3899 let p = Bytes::from(vec![0u8; 1025]);
3900 
3901 if let Err(err) = s.write_sctp(&p.slice(..1024), ppi) {
3902 assert_ne!(
3903 Error::ErrOutboundPacketTooLarge,
3904 err,
3905 "should be not Error::ErrOutboundPacketTooLarge"
3906 );
3907 } else {
3908 panic!("should be error");
3909 }
3910 
3911 if let Err(err) = s.write_sctp(&p.slice(..1025), ppi) {
3912 assert_eq!(
3913 Error::ErrOutboundPacketTooLarge,
3914 err,
3915 "should be Error::ErrOutboundPacketTooLarge"
3916 );
3917 } else {
3918 panic!("should be error");
3919 }
3920 }
3921 
3922 let p = Bytes::from(vec![0u8; 30001]);
3923 
3924 let size_ok = ChunkPayloadData {
3925 beginning_fragment: true,
3926 ending_fragment: true,
3927 tsn: a.peer_last_tsn + 1,
3928 stream_identifier: 1,
3929 user_data: p.slice(..30000),
3930 ..Default::default()
3931 };
3932 
3933 assert!(a.handle_data(&size_ok).is_ok(), "should succeed");
3934 
3935 let too_large = ChunkPayloadData {
3936 beginning_fragment: true,
3937 ending_fragment: true,
3938 tsn: a.peer_last_tsn + 1,
3939 stream_identifier: 1,
3940 user_data: p,
3941 ..Default::default()
3942 };
3943 
3944 if let Err(err) = a.handle_data(&too_large) {
3945 assert_eq!(
3946 Error::ErrInboundPacketTooLarge,
3947 err,
3948 "should be Error::ErrInboundPacketTooLarge"
3949 );
3950 } else {
3951 panic!("should be error");
3952 }
3953 
3954 Ok(())
3955}
3956 
3957#[test]
3958fn test_generate_out_of_band_init() {
3959 let config = TransportConfig::default();
3960 let init_bytes = generate_snap_token(&config).unwrap();
3961 
3962 // Parse it back to validate
3963 let parsed = ChunkInit::unmarshal(&init_bytes).unwrap();
3964 
3965 assert!(!parsed.is_ack, "Should be INIT, not INIT ACK");
3966 assert!(parsed.initiate_tag != 0, "Initiate tag should not be zero");
3967 // Token always advertises u16::MAX for stream counts;
3968 // actual limits are applied from TransportConfig during negotiation.
3969 assert_eq!(
3970 parsed.num_outbound_streams,
3971 u16::MAX,
3972 "Outbound streams should always be u16::MAX in token"
3973 );
3974 assert_eq!(
3975 parsed.num_inbound_streams,
3976 u16::MAX,
3977 "Inbound streams should always be u16::MAX in token"
3978 );
3979 assert_eq!(
3980 parsed.advertised_receiver_window_credit,
3981 config.max_receive_buffer_size(),
3982 "ARWND should match config"
3983 );
3984}
3985 
3986#[test]
3987fn test_generate_out_of_band_init_with_custom_config() {
3988 let config = TransportConfig::default()
3989 .with_max_receive_buffer_size(2_000_000)
3990 .with_max_num_outbound_streams(256)
3991 .with_max_num_inbound_streams(512);
3992 
3993 let init_bytes = generate_snap_token(&config).unwrap();
3994 let parsed = ChunkInit::unmarshal(&init_bytes).unwrap();
3995 
3996 // Token always advertises u16::MAX for stream counts;
3997 // actual limits are applied from TransportConfig during negotiation.
3998 assert_eq!(parsed.num_outbound_streams, u16::MAX);
3999 assert_eq!(parsed.num_inbound_streams, u16::MAX);
4000 assert_eq!(parsed.advertised_receiver_window_credit, 2_000_000);
4001}
4002 
4003#[test]
4004fn test_generate_out_of_band_init_uniqueness() {
4005 // Each call to generate_snap_token creates a new INIT with unique random tags.
4006 let config1 = TransportConfig::default();
4007 let config2 = TransportConfig::default();
4008 
4009 let init1 = generate_snap_token(&config1).unwrap();
4010 let init2 = generate_snap_token(&config2).unwrap();
4011 
4012 let parsed1 = ChunkInit::unmarshal(&init1).unwrap();
4013 let parsed2 = ChunkInit::unmarshal(&init2).unwrap();
4014 
4015 // Initiate tags should be different (random)
4016 assert_ne!(
4017 parsed1.initiate_tag, parsed2.initiate_tag,
4018 "Initiate tags should be unique across calls"
4019 );
4020 
4021 // Initial TSNs should be different (random)
4022 assert_ne!(
4023 parsed1.initial_tsn, parsed2.initial_tsn,
4024 "Initial TSNs should be unique across calls"
4025 );
4026}
4027 
4028#[test]
4029fn test_out_of_band_association_creation() {
4030 let local_config = Arc::new(TransportConfig::default());
4031 let remote_config = TransportConfig::default();
4032 let max_payload_size = 1200;
4033 
4034 let local_init_bytes = generate_snap_token(&local_config).unwrap();
4035 let remote_init_bytes = generate_snap_token(&remote_config).unwrap();
4036 
4037 let local_init = ChunkInit::unmarshal(&local_init_bytes).unwrap();
4038 let remote_init = ChunkInit::unmarshal(&remote_init_bytes).unwrap();
4039 
4040 let remote_addr: SocketAddr = "192.168.1.1:5000".parse().unwrap();
4041 
4042 let assoc = Association::new_with_out_of_band_init(
4043 local_config.clone(),
4044 max_payload_size,
4045 remote_addr,
4046 None,
4047 local_init.clone(),
4048 remote_init.clone(),
4049 )
4050 .expect("Should create out-of-band init association");
4051 
4052 // Verify the association is in ESTABLISHED state
4053 assert_eq!(
4054 assoc.state(),
4055 AssociationState::Established,
4056 "Out-of-band init association should be in ESTABLISHED state"
4057 );
4058 
4059 // Verify handshake is marked complete
4060 assert!(
4061 assoc.handshake_completed,
4062 "Out-of-band init association should have handshake completed"
4063 );
4064 
4065 // Verify verification tags
4066 assert_eq!(
4067 assoc.my_verification_tag, local_init.initiate_tag,
4068 "My verification tag should match local init"
4069 );
4070 assert_eq!(
4071 assoc.peer_verification_tag, remote_init.initiate_tag,
4072 "Peer verification tag should match remote init"
4073 );
4074 
4075 // Verify TSN setup
4076 assert_eq!(
4077 assoc.my_next_tsn, local_init.initial_tsn,
4078 "My next TSN should match local init"
4079 );
4080 assert_eq!(
4081 assoc.peer_last_tsn,
4082 remote_init.initial_tsn.wrapping_sub(1),
4083 "Peer last TSN should be remote init TSN - 1"
4084 );
4085 
4086 // Verify rwnd
4087 assert_eq!(
4088 assoc.rwnd, remote_init.advertised_receiver_window_credit,
4089 "rwnd should match remote advertised credit"
4090 );
4091}
4092 
4093#[test]
4094fn test_out_of_band_association_stream_negotiation() {
4095 let config = Arc::new(
4096 TransportConfig::default()
4097 .with_max_num_outbound_streams(100)
4098 .with_max_num_inbound_streams(200),
4099 );
4100 
4101 // Remote uses default config — token always advertises u16::MAX
4102 let remote_config = TransportConfig::default();
4103 
4104 let local_init_bytes = generate_snap_token(&config).unwrap();
4105 let remote_init_bytes = generate_snap_token(&remote_config).unwrap();
4106 
4107 let local_init = ChunkInit::unmarshal(&local_init_bytes).unwrap();
4108 let remote_init = ChunkInit::unmarshal(&remote_init_bytes).unwrap();
4109 
4110 // Token always advertises u16::MAX for stream counts
4111 assert_eq!(local_init.num_outbound_streams, u16::MAX);
4112 assert_eq!(local_init.num_inbound_streams, u16::MAX);
4113 assert_eq!(remote_init.num_outbound_streams, u16::MAX);
4114 assert_eq!(remote_init.num_inbound_streams, u16::MAX);
4115 
4116 let remote_addr: SocketAddr = "192.168.1.1:5000".parse().unwrap();
4117 
4118 let assoc = Association::new_with_out_of_band_init(
4119 config.clone(),
4120 1200,
4121 remote_addr,
4122 None,
4123 local_init,
4124 remote_init,
4125 )
4126 .expect("Should create out-of-band init association");
4127 
4128 // Stream limits should be clamped by the local config, since the
4129 // remote token always offers u16::MAX.
4130 // my_max_num_outbound_streams = min(config.max_out=100, remote_in=MAX) = 100
4131 assert_eq!(
4132 assoc.my_max_num_outbound_streams, 100,
4133 "Outbound streams should be clamped by local config"
4134 );
4135 
4136 // my_max_num_inbound_streams = min(config.max_in=200, remote_out=MAX) = 200
4137 assert_eq!(
4138 assoc.my_max_num_inbound_streams, 200,
4139 "Inbound streams should be clamped by local config"
4140 );
4141}
4142 
4143#[test]
4144fn test_out_of_band_connected_event() {
4145 let local_config = Arc::new(TransportConfig::default());
4146 let remote_config = TransportConfig::default();
4147 
4148 let local_init_bytes = generate_snap_token(&local_config).unwrap();
4149 let remote_init_bytes = generate_snap_token(&remote_config).unwrap();
4150 
4151 let local_init = ChunkInit::unmarshal(&local_init_bytes).unwrap();
4152 let remote_init = ChunkInit::unmarshal(&remote_init_bytes).unwrap();
4153 
4154 let remote_addr: SocketAddr = "192.168.1.1:5000".parse().unwrap();
4155 
4156 let mut assoc = Association::new_with_out_of_band_init(
4157 local_config.clone(),
4158 1200,
4159 remote_addr,
4160 None,
4161 local_init,
4162 remote_init,
4163 )
4164 .expect("Should create out-of-band init association");
4165 
4166 // Poll should return a Connected event
4167 let event = assoc.poll();
4168 assert!(
4169 matches!(event, Some(Event::Connected)),
4170 "Should emit Connected event, got {:?}",
4171 event
4172 );
4173}
4174 
4175#[test]
4176fn test_out_of_band_symmetric_setup() {
4177 // Test that both sides of an out-of-band init association work correctly
4178 let config_a = Arc::new(TransportConfig::default());
4179 let config_b = Arc::new(TransportConfig::default());
4180 
4181 let init_a_bytes = generate_snap_token(&config_a).unwrap();
4182 let init_b_bytes = generate_snap_token(&config_b).unwrap();
4183 
4184 let init_a = ChunkInit::unmarshal(&init_a_bytes).unwrap();
4185 let init_b = ChunkInit::unmarshal(&init_b_bytes).unwrap();
4186 
4187 let addr_a: SocketAddr = "192.168.1.1:5000".parse().unwrap();
4188 let addr_b: SocketAddr = "192.168.1.2:5000".parse().unwrap();
4189 
4190 // Create association A (local=A, remote=B)
4191 let assoc_a = Association::new_with_out_of_band_init(
4192 config_a.clone(),
4193 1200,
4194 addr_b,
4195 None,
4196 init_a.clone(),
4197 init_b.clone(),
4198 )
4199 .expect("Should create association A");
4200 
4201 // Create association B (local=B, remote=A)
4202 let assoc_b = Association::new_with_out_of_band_init(
4203 config_b.clone(),
4204 1200,
4205 addr_a,
4206 None,
4207 init_b.clone(),
4208 init_a.clone(),
4209 )
4210 .expect("Should create association B");
4211 
4212 // Verify both are in ESTABLISHED state
4213 assert_eq!(assoc_a.state(), AssociationState::Established);
4214 assert_eq!(assoc_b.state(), AssociationState::Established);
4215 
4216 // Verify verification tags are cross-matched
4217 assert_eq!(assoc_a.my_verification_tag, assoc_b.peer_verification_tag);
4218 assert_eq!(assoc_b.my_verification_tag, assoc_a.peer_verification_tag);
4219}
4220 
4221#[test]
4222fn test_out_of_band_with_forward_tsn_support() {
4223 let local_config = Arc::new(TransportConfig::default());
4224 let remote_config = TransportConfig::default();
4225 
4226 let local_init_bytes = generate_snap_token(&local_config).unwrap();
4227 let remote_init_bytes = generate_snap_token(&remote_config).unwrap();
4228 
4229 let local_init = ChunkInit::unmarshal(&local_init_bytes).unwrap();
4230 let remote_init = ChunkInit::unmarshal(&remote_init_bytes).unwrap();
4231 
4232 // Verify supported extensions are present
4233 let mut has_forward_tsn = false;
4234 for param in &local_init.params {
4235 if let Some(ext) = param
4236 .as_any()
4237 .downcast_ref::<crate::param::param_supported_extensions::ParamSupportedExtensions>(
4238 ) {
4239 for ct in &ext.chunk_types {
4240 if *ct == crate::chunk::chunk_type::CT_FORWARD_TSN {
4241 has_forward_tsn = true;
4242 }
4243 }
4244 }
4245 }
4246 assert!(
4247 has_forward_tsn,
4248 "Generated INIT should include ForwardTSN support"
4249 );
4250 
4251 let remote_addr: SocketAddr = "192.168.1.1:5000".parse().unwrap();
4252 
4253 let assoc = Association::new_with_out_of_band_init(
4254 local_config.clone(),
4255 1200,
4256 remote_addr,
4257 None,
4258 local_init,
4259 remote_init,
4260 )
4261 .expect("Should create out-of-band init association");
4262 
4263 assert!(
4264 assoc.use_forward_tsn,
4265 "Out-of-band init association should have ForwardTSN enabled"
4266 );
4267}
4268 
4269#[test]
4270fn test_out_of_band_initial_tsn_zero_wrap() {
4271 // Test edge case where initial TSN is 0 (wraps to MAX)
4272 let config = Arc::new(TransportConfig::default());
4273 
4274 let local_init_bytes = generate_snap_token(&config).unwrap();
4275 let mut remote_init = ChunkInit::unmarshal(&local_init_bytes).unwrap();
4276 
4277 // Set initial TSN to 0 to test the edge case
4278 remote_init.initial_tsn = 0;
4279 remote_init.initiate_tag = 12345;
4280 
4281 // Generate a fresh local init for the association
4282 let actual_local_init_bytes = generate_snap_token(&config).unwrap();
4283 let local_init = ChunkInit::unmarshal(&actual_local_init_bytes).unwrap();
4284 let remote_addr: SocketAddr = "192.168.1.1:5000".parse().unwrap();
4285 
4286 let assoc = Association::new_with_out_of_band_init(
4287 config.clone(),
4288 1200,
4289 remote_addr,
4290 None,
4291 local_init,
4292 remote_init,
4293 )
4294 .expect("Should create out-of-band init association");
4295 
4296 // peer_last_tsn should be u32::MAX when initial_tsn is 0
4297 assert_eq!(
4298 assoc.peer_last_tsn,
4299 u32::MAX,
4300 "peer_last_tsn should wrap to MAX when initial_tsn is 0"
4301 );
4302}
4303 
4304#[test]
4305fn test_out_of_band_rwnd_negotiation() {
4306 let local_config = Arc::new(TransportConfig::default().with_max_receive_buffer_size(500_000));
4307 
4308 let remote_config = TransportConfig::default().with_max_receive_buffer_size(300_000);
4309 
4310 let local_init_bytes = generate_snap_token(&local_config).unwrap();
4311 let remote_init_bytes = generate_snap_token(&remote_config).unwrap();
4312 
4313 let local_init = ChunkInit::unmarshal(&local_init_bytes).unwrap();
4314 let remote_init = ChunkInit::unmarshal(&remote_init_bytes).unwrap();
4315 
4316 let remote_addr: SocketAddr = "192.168.1.1:5000".parse().unwrap();
4317 
4318 let assoc = Association::new_with_out_of_band_init(
4319 local_config.clone(),
4320 1200,
4321 remote_addr,
4322 None,
4323 local_init,
4324 remote_init,
4325 )
4326 .expect("Should create out-of-band init association");
4327 
4328 // rwnd should be set to remote's advertised receiver window credit
4329 assert_eq!(
4330 assoc.rwnd, 300_000,
4331 "rwnd should be remote's advertised window"
4332 );
4333}
4334 
4335#[test]
4336fn test_initial_cwnd_small_mtu() {
4337 // Regression test for #47: a small MTU made 4*MTU < 4380, which previously
4338 // panicked because clamp(4380, 4*MTU) was called with min > max.
4339 let max_payload_size = 100;
4340 let mtu = max_payload_size + COMMON_HEADER_SIZE + DATA_CHUNK_HEADER_SIZE;
4341 assert!(4 * mtu < 4380, "test must exercise the small-MTU branch");
4342 
4343 let assoc = Association::new(
4344 None,
4345 Arc::new(TransportConfig::default()),
4346 max_payload_size,
4347 0,
4348 SocketAddr::from_str("0.0.0.0:0").unwrap(),
4349 None,
4350 Instant::now(),
4351 );
4352 
4353 // RFC 4960 Sec 7.2.1: min(4*MTU, max(2*MTU, 4380)).
4354 assert_eq!(assoc.cwnd, (4 * mtu).min((2 * mtu).max(4380)));
4355 assert_eq!(assoc.cwnd, 4 * mtu);
4356}