1use std::collections::{HashMap, HashSet};
9use std::net::Ipv4Addr;
10use std::num::NonZeroU64;
11use std::time::{Duration, Instant};
12
13use thiserror::Error;
14
15use crate::message::values::{QosDirection, QosErrorCode, QosReservationStyle, RsvpErrorCode};
16use crate::message::{ControlMessage, QosApplicationIdentifier, QosFlow, QosTrafficSpecification};
17
18#[derive(Clone, Copy, Debug, Eq, PartialEq)]
20pub struct QosReservationLimits {
21 pub maximum_live_reservations: usize,
22 pub maximum_generations: usize,
26 pub maximum_retries: u32,
27 pub maximum_retry_timer: u32,
28 pub maximum_response_timeout: Duration,
29}
30
31impl QosReservationLimits {
32 pub fn validate(self) -> Result<Self, QosReservationError> {
35 if self.maximum_live_reservations == 0 {
36 return Err(QosReservationError::InvalidLimits(
37 "maximum live reservations must be nonzero",
38 ));
39 }
40 if self.maximum_generations < self.maximum_live_reservations {
41 return Err(QosReservationError::InvalidLimits(
42 "generation capacity must cover every live reservation",
43 ));
44 }
45 if self.maximum_retry_timer == 0 {
46 return Err(QosReservationError::InvalidLimits(
47 "maximum retry timer must be nonzero",
48 ));
49 }
50 if self.maximum_response_timeout.is_zero() {
51 return Err(QosReservationError::InvalidLimits(
52 "maximum response timeout must be nonzero",
53 ));
54 }
55 Ok(self)
56 }
57}
58
59#[derive(Clone, Copy, Debug, Eq, PartialEq)]
61pub enum QosReservationSetup {
62 Listen { confirmation_required: bool },
63 Path,
64}
65
66#[derive(Clone, Debug, Eq, PartialEq)]
68pub struct QosReservationPolicy {
69 pub reservation_style: QosReservationStyle,
70 pub maximum_retries: u32,
71 pub retry_timer: u32,
72 pub preemption_priority: u32,
73 pub defending_priority: u32,
74 pub traffic: QosTrafficSpecification,
75 pub application: QosApplicationIdentifier,
76}
77
78#[derive(Clone, Debug, Eq, PartialEq)]
84pub struct QosReservationRequest {
85 pub flow: QosFlow,
86 pub direction: QosDirection,
87 pub setup: QosReservationSetup,
88 pub policy: QosReservationPolicy,
89 pub response_timeout: Duration,
90}
91
92#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
94pub struct QosReservationId(NonZeroU64);
95
96impl QosReservationId {
97 pub const fn get(self) -> u64 {
99 self.0.get()
100 }
101}
102
103#[derive(Clone, Copy, Debug, Eq, PartialEq)]
105pub enum QosReservationState {
106 Establishing,
107 Active,
108 Modifying,
109}
110
111#[derive(Clone, Copy, Debug, Eq, PartialEq)]
113pub struct QosReservationFailure {
114 pub error_code: QosErrorCode,
115 pub failure_node: Ipv4Addr,
116 pub rsvp_error_code: RsvpErrorCode,
117 pub rsvp_error_subcode: u32,
118 pub rsvp_error_flags: u32,
119}
120
121#[derive(Clone, Debug, Eq, PartialEq)]
123pub enum QosReservationEvent {
124 Established {
125 id: QosReservationId,
126 },
127 ActiveWithoutConfirmation {
128 id: QosReservationId,
129 },
130 EstablishmentTimedOut {
131 id: QosReservationId,
132 },
133 Failed {
134 id: QosReservationId,
135 failure: QosReservationFailure,
136 },
137 Preempted {
138 id: QosReservationId,
139 failure: QosReservationFailure,
140 },
141 ModificationFailed {
142 id: QosReservationId,
143 failure: QosReservationFailure,
144 },
145 ModificationOutcomeUnknown {
148 id: QosReservationId,
149 },
150 TornDown {
151 id: QosReservationId,
152 },
153}
154
155#[derive(Clone, Debug, Default, Eq, PartialEq)]
157pub struct QosTransition {
158 messages: Vec<ControlMessage>,
159 events: Vec<QosReservationEvent>,
160}
161
162impl QosTransition {
163 pub fn messages(&self) -> &[ControlMessage] {
165 &self.messages
166 }
167
168 pub fn events(&self) -> &[QosReservationEvent] {
170 &self.events
171 }
172
173 pub fn into_parts(self) -> (Vec<ControlMessage>, Vec<QosReservationEvent>) {
175 (self.messages, self.events)
176 }
177}
178
179#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
180struct ReservationKey {
181 flow: QosFlow,
182 direction: QosDirection,
183}
184
185#[derive(Clone, Debug)]
186enum ReservationPhase {
187 Establishing { deadline: Instant },
188 Active(ModificationCorrelation),
189}
190
191#[derive(Clone, Debug)]
192enum ModificationCorrelation {
193 Available,
194 Pending { deadline: Instant },
195 OutcomeUnknown,
196 FailureReported,
197}
198
199#[derive(Clone, Debug)]
200struct Reservation {
201 key: ReservationKey,
202 phase: ReservationPhase,
203}
204
205#[derive(Clone, Debug, Error, Eq, PartialEq)]
207pub enum QosReservationError {
208 #[error("invalid QoS reservation limits: {0}")]
209 InvalidLimits(&'static str),
210 #[error("invalid QoS reservation: {0}")]
211 InvalidRequest(&'static str),
212 #[error("QoS live-reservation capacity is exhausted")]
213 LiveCapacityExhausted,
214 #[error("QoS reservation generation capacity is exhausted")]
215 GenerationCapacityExhausted,
216 #[error("QoS wire flow identity was already used")]
217 FlowIdentityReused,
218 #[error("QoS reservation {0:?} does not exist")]
219 UnknownReservation(QosReservationId),
220 #[error("QoS reservation {0:?} already issued its one correlatable modification")]
221 ModificationAlreadyIssued(QosReservationId),
222 #[error("QoS reservation {id:?} cannot {operation} while {state:?}")]
223 InvalidState {
224 id: QosReservationId,
225 operation: &'static str,
226 state: QosReservationState,
227 },
228}
229
230#[derive(Debug)]
232pub struct QosReservationController {
233 limits: QosReservationLimits,
234 next_id: Option<NonZeroU64>,
235 by_id: HashMap<QosReservationId, Reservation>,
236 by_key: HashMap<ReservationKey, QosReservationId>,
237 retired: HashSet<ReservationKey>,
238}
239
240impl QosReservationController {
241 pub fn new(limits: QosReservationLimits) -> Result<Self, QosReservationError> {
243 Ok(Self {
244 limits: limits.validate()?,
245 next_id: NonZeroU64::new(1),
246 by_id: HashMap::new(),
247 by_key: HashMap::new(),
248 retired: HashSet::new(),
249 })
250 }
251
252 pub fn start(
255 &mut self,
256 request: QosReservationRequest,
257 now: Instant,
258 ) -> Result<(QosReservationId, QosTransition), QosReservationError> {
259 self.validate_request(&request)?;
260 if self.by_id.len() >= self.limits.maximum_live_reservations {
261 return Err(QosReservationError::LiveCapacityExhausted);
262 }
263 if self.retired.len()
264 >= self
265 .limits
266 .maximum_generations
267 .saturating_sub(self.by_id.len())
268 {
269 return Err(QosReservationError::GenerationCapacityExhausted);
270 }
271 let key = ReservationKey {
272 flow: request.flow,
273 direction: request.direction,
274 };
275 if self.by_key.contains_key(&key) || self.retired.contains(&key) {
276 return Err(QosReservationError::FlowIdentityReused);
277 }
278 let id = self.allocate_id()?;
279 let message = setup_message(&request);
280 let confirmed = !matches!(
281 request.setup,
282 QosReservationSetup::Listen {
283 confirmation_required: false
284 }
285 );
286 let phase = if confirmed {
287 ReservationPhase::Establishing {
288 deadline: deadline(now, request.response_timeout)?,
289 }
290 } else {
291 ReservationPhase::Active(ModificationCorrelation::Available)
292 };
293 let mut transition = QosTransition {
294 messages: vec![message],
295 events: Vec::new(),
296 };
297 if !confirmed {
298 transition
299 .events
300 .push(QosReservationEvent::ActiveWithoutConfirmation { id });
301 }
302 self.by_key.insert(key, id);
303 self.by_id.insert(id, Reservation { key, phase });
304 Ok((id, transition))
305 }
306
307 pub fn state(&self, id: QosReservationId) -> Option<QosReservationState> {
309 self.by_id
310 .get(&id)
311 .map(|reservation| state(&reservation.phase))
312 }
313
314 pub fn len(&self) -> usize {
316 self.by_id.len()
317 }
318
319 pub fn is_empty(&self) -> bool {
321 self.by_id.is_empty()
322 }
323
324 pub fn modify(
329 &mut self,
330 id: QosReservationId,
331 traffic: QosTrafficSpecification,
332 application: QosApplicationIdentifier,
333 response_timeout: Duration,
334 now: Instant,
335 ) -> Result<QosTransition, QosReservationError> {
336 validate_application(&application)?;
337 self.validate_response_timeout(response_timeout)?;
338 let reservation = self
339 .by_id
340 .get_mut(&id)
341 .ok_or(QosReservationError::UnknownReservation(id))?;
342 match reservation.phase {
343 ReservationPhase::Active(ModificationCorrelation::Available) => {}
344 ReservationPhase::Active(_) => {
345 return Err(QosReservationError::ModificationAlreadyIssued(id));
346 }
347 ReservationPhase::Establishing { .. } => {
348 return Err(QosReservationError::InvalidState {
349 id,
350 operation: "modify",
351 state: state(&reservation.phase),
352 });
353 }
354 }
355 let message = ControlMessage::QosModify {
356 flow: reservation.key.flow,
357 direction: reservation.key.direction,
358 traffic,
359 application,
360 };
361 reservation.phase = ReservationPhase::Active(ModificationCorrelation::Pending {
362 deadline: deadline(now, response_timeout)?,
363 });
364 Ok(QosTransition {
365 messages: vec![message],
366 events: Vec::new(),
367 })
368 }
369
370 pub fn update_dscp(
372 &self,
373 id: QosReservationId,
374 dscp: u8,
375 ) -> Result<QosTransition, QosReservationError> {
376 if dscp > 63 {
377 return Err(QosReservationError::InvalidRequest(
378 "DSCP must fit six bits",
379 ));
380 }
381 let reservation = self
382 .by_id
383 .get(&id)
384 .ok_or(QosReservationError::UnknownReservation(id))?;
385 if state(&reservation.phase) != QosReservationState::Active {
386 return Err(QosReservationError::InvalidState {
387 id,
388 operation: "update DSCP",
389 state: state(&reservation.phase),
390 });
391 }
392 Ok(QosTransition {
393 messages: vec![ControlMessage::UpdateDscp {
394 flow: reservation.key.flow,
395 dscp,
396 }],
397 events: Vec::new(),
398 })
399 }
400
401 pub fn teardown(&mut self, id: QosReservationId) -> Result<QosTransition, QosReservationError> {
403 let reservation = self.retire(id)?;
404 Ok(QosTransition {
405 messages: vec![teardown_message(reservation.key)],
406 events: vec![QosReservationEvent::TornDown { id }],
407 })
408 }
409
410 pub fn handle_message(&mut self, message: &ControlMessage) -> QosTransition {
413 match message {
414 ControlMessage::QosReservationNotify { flow, direction } => self
415 .handle_reservation_notify(ReservationKey {
416 flow: *flow,
417 direction: *direction,
418 }),
419 ControlMessage::QosErrorNotify {
420 flow,
421 direction,
422 error_code,
423 failure_node,
424 rsvp_error_code,
425 rsvp_error_subcode,
426 rsvp_error_flags,
427 } => self.handle_error(
428 ReservationKey {
429 flow: *flow,
430 direction: *direction,
431 },
432 QosReservationFailure {
433 error_code: *error_code,
434 failure_node: *failure_node,
435 rsvp_error_code: *rsvp_error_code,
436 rsvp_error_subcode: *rsvp_error_subcode,
437 rsvp_error_flags: *rsvp_error_flags,
438 },
439 ),
440 _ => QosTransition::default(),
441 }
442 }
443
444 pub fn poll(&mut self, now: Instant) -> QosTransition {
446 let mut ids = self
447 .by_id
448 .iter()
449 .filter_map(|(&id, reservation)| match reservation.phase {
450 ReservationPhase::Establishing { deadline }
451 | ReservationPhase::Active(ModificationCorrelation::Pending { deadline })
452 if deadline <= now =>
453 {
454 Some(id)
455 }
456 _ => None,
457 })
458 .collect::<Vec<_>>();
459 ids.sort_unstable();
460
461 let mut transition = QosTransition::default();
462 for id in ids {
463 match self
464 .by_id
465 .get(&id)
466 .map(|reservation| state(&reservation.phase))
467 {
468 Some(QosReservationState::Establishing) => {
469 if let Ok(reservation) = self.retire(id) {
470 transition.messages.push(teardown_message(reservation.key));
471 transition
472 .events
473 .push(QosReservationEvent::EstablishmentTimedOut { id });
474 }
475 }
476 Some(QosReservationState::Modifying) => {
477 if let Some(reservation) = self.by_id.get_mut(&id) {
478 reservation.phase =
479 ReservationPhase::Active(ModificationCorrelation::OutcomeUnknown);
480 transition
481 .events
482 .push(QosReservationEvent::ModificationOutcomeUnknown { id });
483 }
484 }
485 _ => {}
486 }
487 }
488 transition
489 }
490
491 pub fn drain(&mut self) -> QosTransition {
493 let mut ids = self.by_id.keys().copied().collect::<Vec<_>>();
494 ids.sort_unstable();
495 let mut transition = QosTransition::default();
496 for id in ids {
497 if let Ok(reservation) = self.retire(id) {
498 transition.messages.push(teardown_message(reservation.key));
499 transition.events.push(QosReservationEvent::TornDown { id });
500 }
501 }
502 transition
503 }
504
505 fn handle_reservation_notify(&mut self, key: ReservationKey) -> QosTransition {
506 let Some(id) = self.by_key.get(&key).copied() else {
507 return QosTransition::default();
508 };
509 let Some(reservation) = self.by_id.get_mut(&id) else {
510 return QosTransition::default();
511 };
512 let event = match &reservation.phase {
513 ReservationPhase::Establishing { .. } => QosReservationEvent::Established { id },
514 ReservationPhase::Active(_) => {
515 return QosTransition::default();
516 }
517 };
518 reservation.phase = ReservationPhase::Active(ModificationCorrelation::Available);
519 QosTransition {
520 messages: Vec::new(),
521 events: vec![event],
522 }
523 }
524
525 fn handle_error(
526 &mut self,
527 key: ReservationKey,
528 failure: QosReservationFailure,
529 ) -> QosTransition {
530 let Some(id) = self.by_key.get(&key).copied() else {
531 return QosTransition::default();
532 };
533 if is_modify_failure(failure.error_code) {
534 let correlates_modification = self.by_id.get(&id).is_some_and(|reservation| {
535 matches!(
536 reservation.phase,
537 ReservationPhase::Active(ModificationCorrelation::Pending { .. })
538 )
539 });
540 if !correlates_modification {
541 return QosTransition::default();
542 }
543 if let Some(reservation) = self.by_id.get_mut(&id) {
544 reservation.phase =
545 ReservationPhase::Active(ModificationCorrelation::FailureReported);
546 }
547 return QosTransition {
548 messages: Vec::new(),
549 events: vec![QosReservationEvent::ModificationFailed { id, failure }],
550 };
551 }
552
553 let Ok(reservation) = self.retire(id) else {
554 return QosTransition::default();
555 };
556 if failure.error_code == QosErrorCode::ReservationTornDown {
557 return QosTransition {
558 messages: Vec::new(),
559 events: vec![QosReservationEvent::TornDown { id }],
560 };
561 }
562 let event = if is_preemption(failure.error_code) {
563 QosReservationEvent::Preempted { id, failure }
564 } else {
565 QosReservationEvent::Failed { id, failure }
566 };
567 QosTransition {
568 messages: vec![teardown_message(reservation.key)],
569 events: vec![event],
570 }
571 }
572
573 fn retire(&mut self, id: QosReservationId) -> Result<Reservation, QosReservationError> {
574 let reservation = self
575 .by_id
576 .remove(&id)
577 .ok_or(QosReservationError::UnknownReservation(id))?;
578 self.by_key.remove(&reservation.key);
579 self.retired.insert(reservation.key);
580 Ok(reservation)
581 }
582
583 fn allocate_id(&mut self) -> Result<QosReservationId, QosReservationError> {
584 let value = self
585 .next_id
586 .take()
587 .ok_or(QosReservationError::GenerationCapacityExhausted)?;
588 self.next_id = value.get().checked_add(1).and_then(NonZeroU64::new);
589 Ok(QosReservationId(value))
590 }
591
592 fn validate_request(&self, request: &QosReservationRequest) -> Result<(), QosReservationError> {
593 validate_flow(request.flow)?;
594 if !request.direction.is_known() {
595 return Err(QosReservationError::InvalidRequest(
596 "reservation direction is unknown",
597 ));
598 }
599 if !request.policy.reservation_style.is_known() {
600 return Err(QosReservationError::InvalidRequest(
601 "reservation style is unknown",
602 ));
603 }
604 if request.policy.maximum_retries > self.limits.maximum_retries {
605 return Err(QosReservationError::InvalidRequest(
606 "retry count exceeds controller policy",
607 ));
608 }
609 if request.policy.retry_timer == 0
610 || request.policy.retry_timer > self.limits.maximum_retry_timer
611 {
612 return Err(QosReservationError::InvalidRequest(
613 "retry timer is outside controller policy",
614 ));
615 }
616 self.validate_response_timeout(request.response_timeout)?;
617 validate_application(&request.policy.application)
618 }
619
620 fn validate_response_timeout(&self, timeout: Duration) -> Result<(), QosReservationError> {
621 if timeout.is_zero() || timeout > self.limits.maximum_response_timeout {
622 return Err(QosReservationError::InvalidRequest(
623 "response timeout is outside controller policy",
624 ));
625 }
626 Ok(())
627 }
628}
629
630fn setup_message(request: &QosReservationRequest) -> ControlMessage {
631 let policy = &request.policy;
632 match request.setup {
633 QosReservationSetup::Listen {
634 confirmation_required,
635 } => ControlMessage::QosListen {
636 flow: request.flow,
637 reservation_style: policy.reservation_style,
638 maximum_retries: policy.maximum_retries,
639 retry_timer: policy.retry_timer,
640 confirmation_required,
641 preemption_priority: policy.preemption_priority,
642 defending_priority: policy.defending_priority,
643 traffic: policy.traffic,
644 application: policy.application.clone(),
645 },
646 QosReservationSetup::Path => ControlMessage::QosPath {
647 flow: request.flow,
648 reservation_style: policy.reservation_style,
649 maximum_retries: policy.maximum_retries,
650 retry_timer: policy.retry_timer,
651 preemption_priority: policy.preemption_priority,
652 defending_priority: policy.defending_priority,
653 traffic: policy.traffic,
654 application: policy.application.clone(),
655 },
656 }
657}
658
659fn validate_flow(flow: QosFlow) -> Result<(), QosReservationError> {
660 if flow.conference_id.get() == 0
661 || flow.call_reference.get() == 0
662 || flow.passthrough_party_id.get() == 0
663 {
664 return Err(QosReservationError::InvalidRequest(
665 "flow identities must be nonzero",
666 ));
667 }
668 if flow.address.is_unspecified() || flow.address.is_multicast() || flow.port == 0 {
669 return Err(QosReservationError::InvalidRequest(
670 "flow endpoint must be usable unicast",
671 ));
672 }
673 Ok(())
674}
675
676fn validate_application(application: &QosApplicationIdentifier) -> Result<(), QosReservationError> {
677 for (value, maximum) in [
678 (&application.vendor_id, 31),
679 (&application.version, 15),
680 (&application.application_name, 31),
681 (&application.sub_application_id, 31),
682 ] {
683 if value.len() > maximum || value.contains('\0') {
684 return Err(QosReservationError::InvalidRequest(
685 "application identity exceeds its fixed text field",
686 ));
687 }
688 }
689 Ok(())
690}
691
692fn deadline(now: Instant, timeout: Duration) -> Result<Instant, QosReservationError> {
693 now.checked_add(timeout)
694 .ok_or(QosReservationError::InvalidRequest(
695 "response deadline overflows the monotonic clock",
696 ))
697}
698
699fn teardown_message(key: ReservationKey) -> ControlMessage {
700 ControlMessage::QosTeardown {
701 flow: key.flow,
702 direction: key.direction,
703 }
704}
705
706fn state(phase: &ReservationPhase) -> QosReservationState {
707 match phase {
708 ReservationPhase::Establishing { .. } => QosReservationState::Establishing,
709 ReservationPhase::Active(ModificationCorrelation::Pending { .. }) => {
710 QosReservationState::Modifying
711 }
712 ReservationPhase::Active(_) => QosReservationState::Active,
713 }
714}
715
716fn is_modify_failure(error: QosErrorCode) -> bool {
717 matches!(
718 error,
719 QosErrorCode::ReservationModifyFailed | QosErrorCode::PathModifyFailed
720 )
721}
722
723fn is_preemption(error: QosErrorCode) -> bool {
724 matches!(
725 error,
726 QosErrorCode::ReservationPreempted | QosErrorCode::PathPreempted
727 )
728}
729
730#[cfg(test)]
731mod tests {
732 use super::*;
733 use crate::message::catalog::{MessageId, MessageRoute, RuntimeUse};
734 use crate::message::values::{Codec, ProtocolVersion};
735 use crate::message::wire::FrameDecoder;
736 use crate::types::{CallReference, ConferenceId, PassthroughPartyId};
737
738 fn limits() -> QosReservationLimits {
739 QosReservationLimits {
740 maximum_live_reservations: 4,
741 maximum_generations: 8,
742 maximum_retries: 3,
743 maximum_retry_timer: 10,
744 maximum_response_timeout: Duration::from_secs(30),
745 }
746 }
747
748 fn flow(token: u32) -> QosFlow {
749 QosFlow {
750 conference_id: ConferenceId::new(40),
751 call_reference: CallReference::new(41),
752 passthrough_party_id: PassthroughPartyId::new(token),
753 address: "192.0.2.40".parse().unwrap(),
754 port: 16_000,
755 }
756 }
757
758 fn application() -> QosApplicationIdentifier {
759 QosApplicationIdentifier {
760 vendor_id: "vendor".into(),
761 version: "1".into(),
762 application_name: "audio".into(),
763 sub_application_id: "primary".into(),
764 }
765 }
766
767 fn traffic(average_bit_rate: u32) -> QosTrafficSpecification {
768 QosTrafficSpecification {
769 codec: Codec::Pcmu,
770 average_bit_rate,
771 burst_size: 1_200,
772 peak_rate: average_bit_rate * 2,
773 }
774 }
775
776 fn request(token: u32, setup: QosReservationSetup) -> QosReservationRequest {
777 QosReservationRequest {
778 flow: flow(token),
779 direction: QosDirection::Receive,
780 setup,
781 policy: QosReservationPolicy {
782 reservation_style: QosReservationStyle::SharedExplicit,
783 maximum_retries: 3,
784 retry_timer: 4,
785 preemption_priority: 5,
786 defending_priority: 6,
787 traffic: traffic(64_000),
788 application: application(),
789 },
790 response_timeout: Duration::from_secs(12),
791 }
792 }
793
794 fn reservation_notify(request: &QosReservationRequest) -> ControlMessage {
795 ControlMessage::QosReservationNotify {
796 flow: request.flow,
797 direction: request.direction,
798 }
799 }
800
801 fn error_notify(request: &QosReservationRequest, error_code: QosErrorCode) -> ControlMessage {
802 let failure = failure(error_code);
803 ControlMessage::QosErrorNotify {
804 flow: request.flow,
805 direction: request.direction,
806 error_code: failure.error_code,
807 failure_node: failure.failure_node,
808 rsvp_error_code: failure.rsvp_error_code,
809 rsvp_error_subcode: failure.rsvp_error_subcode,
810 rsvp_error_flags: failure.rsvp_error_flags,
811 }
812 }
813
814 fn failure(error_code: QosErrorCode) -> QosReservationFailure {
815 QosReservationFailure {
816 error_code,
817 failure_node: "198.51.100.4".parse().unwrap(),
818 rsvp_error_code: RsvpErrorCode::ServicePreempted,
819 rsvp_error_subcode: 7,
820 rsvp_error_flags: 8,
821 }
822 }
823
824 #[test]
825 fn setup_routes_exact_policy_and_correlates_only_the_live_key() {
826 assert_eq!(
827 MessageId::QosListen.contract().unwrap().route,
828 MessageRoute::ControlToServiceNode
829 );
830 assert_eq!(
831 MessageId::QosReservationNotify.contract().unwrap().route,
832 MessageRoute::ServiceNodeToControl
833 );
834 assert_eq!(
835 MessageId::QosListen.contract().unwrap().runtime_use,
836 RuntimeUse::ConditionalServiceNodeOutput
837 );
838 assert_eq!(
839 MessageId::QosReservationNotify
840 .contract()
841 .unwrap()
842 .runtime_use,
843 RuntimeUse::ServiceNodeInput
844 );
845 let now = Instant::now();
846 let request = request(
847 42,
848 QosReservationSetup::Listen {
849 confirmation_required: true,
850 },
851 );
852 let mut controller = QosReservationController::new(limits()).unwrap();
853 let (id, transition) = controller.start(request.clone(), now).unwrap();
854 assert_eq!(
855 controller.state(id),
856 Some(QosReservationState::Establishing)
857 );
858 assert_eq!(
859 transition.messages(),
860 &[ControlMessage::QosListen {
861 flow: request.flow,
862 reservation_style: request.policy.reservation_style,
863 maximum_retries: 3,
864 retry_timer: 4,
865 confirmation_required: true,
866 preemption_priority: 5,
867 defending_priority: 6,
868 traffic: request.policy.traffic,
869 application: request.policy.application.clone(),
870 }]
871 );
872
873 let mut wrong = request.clone();
874 wrong.flow.passthrough_party_id = PassthroughPartyId::new(43);
875 assert_eq!(
876 controller.handle_message(&reservation_notify(&wrong)),
877 QosTransition::default()
878 );
879 assert_eq!(
880 controller
881 .handle_message(&reservation_notify(&request))
882 .events(),
883 &[QosReservationEvent::Established { id }]
884 );
885 assert_eq!(controller.state(id), Some(QosReservationState::Active));
886 assert_eq!(
887 controller.handle_message(&reservation_notify(&request)),
888 QosTransition::default()
889 );
890 }
891
892 #[test]
893 fn retry_policy_and_deadlines_are_bounded_without_assigning_wire_timer_units() {
894 let now = Instant::now();
895 let mut controller = QosReservationController::new(limits()).unwrap();
896 let mut invalid = request(44, QosReservationSetup::Path);
897 invalid.policy.maximum_retries = 4;
898 assert!(matches!(
899 controller.start(invalid, now),
900 Err(QosReservationError::InvalidRequest(_))
901 ));
902 let mut invalid = request(44, QosReservationSetup::Path);
903 invalid.policy.retry_timer = 11;
904 assert!(matches!(
905 controller.start(invalid, now),
906 Err(QosReservationError::InvalidRequest(_))
907 ));
908
909 let request = request(44, QosReservationSetup::Path);
910 let (id, _) = controller.start(request.clone(), now).unwrap();
911 assert!(
912 controller
913 .poll(now + request.response_timeout - Duration::from_nanos(1))
914 .events()
915 .is_empty()
916 );
917 let expired = controller.poll(now + request.response_timeout);
918 assert_eq!(
919 expired.messages(),
920 &[ControlMessage::QosTeardown {
921 flow: request.flow,
922 direction: request.direction,
923 }]
924 );
925 assert_eq!(
926 expired.events(),
927 &[QosReservationEvent::EstablishmentTimedOut { id }]
928 );
929 assert!(controller.is_empty());
930 assert!(matches!(
931 controller.start(request, now),
932 Err(QosReservationError::FlowIdentityReused)
933 ));
934 }
935
936 #[test]
937 fn modification_failure_preserves_active_state_and_preemption_retires_it() {
938 let now = Instant::now();
939 let request = request(
940 45,
941 QosReservationSetup::Listen {
942 confirmation_required: true,
943 },
944 );
945 let mut controller = QosReservationController::new(limits()).unwrap();
946 let (id, _) = controller.start(request.clone(), now).unwrap();
947 controller.handle_message(&reservation_notify(&request));
948
949 let modified_traffic = traffic(96_000);
950 let modify = controller
951 .modify(
952 id,
953 modified_traffic,
954 application(),
955 Duration::from_secs(5),
956 now,
957 )
958 .unwrap();
959 assert_eq!(
960 modify.messages(),
961 &[ControlMessage::QosModify {
962 flow: request.flow,
963 direction: request.direction,
964 traffic: modified_traffic,
965 application: application(),
966 }]
967 );
968 assert_eq!(
969 controller.handle_message(&reservation_notify(&request)),
970 QosTransition::default()
971 );
972 assert_eq!(controller.state(id), Some(QosReservationState::Modifying));
973 let failed = controller.handle_message(&error_notify(
974 &request,
975 QosErrorCode::ReservationModifyFailed,
976 ));
977 assert_eq!(
978 failed.events(),
979 &[QosReservationEvent::ModificationFailed {
980 id,
981 failure: failure(QosErrorCode::ReservationModifyFailed),
982 }]
983 );
984 assert_eq!(controller.state(id), Some(QosReservationState::Active));
985 assert_eq!(
986 controller.handle_message(&error_notify(
987 &request,
988 QosErrorCode::ReservationModifyFailed,
989 )),
990 QosTransition::default()
991 );
992 assert!(matches!(
993 controller.modify(
994 id,
995 traffic(112_000),
996 application(),
997 Duration::from_secs(5),
998 now,
999 ),
1000 Err(QosReservationError::ModificationAlreadyIssued(actual)) if actual == id
1001 ));
1002
1003 let preempted =
1004 controller.handle_message(&error_notify(&request, QosErrorCode::ReservationPreempted));
1005 assert_eq!(
1006 preempted.messages(),
1007 &[ControlMessage::QosTeardown {
1008 flow: request.flow,
1009 direction: request.direction,
1010 }]
1011 );
1012 assert_eq!(
1013 preempted.events(),
1014 &[QosReservationEvent::Preempted {
1015 id,
1016 failure: failure(QosErrorCode::ReservationPreempted),
1017 }]
1018 );
1019 assert!(controller.is_empty());
1020 assert_eq!(
1021 controller.handle_message(&reservation_notify(&request)),
1022 QosTransition::default()
1023 );
1024 }
1025
1026 #[test]
1027 fn terminal_errors_teardown_once_and_remote_teardown_does_not_echo() {
1028 let now = Instant::now();
1029 let mut controller = QosReservationController::new(limits()).unwrap();
1030 let failed_request = request(
1031 48,
1032 QosReservationSetup::Listen {
1033 confirmation_required: true,
1034 },
1035 );
1036 let (failed_id, _) = controller.start(failed_request.clone(), now).unwrap();
1037 let failed = controller.handle_message(&error_notify(
1038 &failed_request,
1039 QosErrorCode::ResourceUnavailable,
1040 ));
1041 assert_eq!(
1042 failed.messages(),
1043 &[ControlMessage::QosTeardown {
1044 flow: failed_request.flow,
1045 direction: failed_request.direction,
1046 }]
1047 );
1048 assert_eq!(
1049 failed.events(),
1050 &[QosReservationEvent::Failed {
1051 id: failed_id,
1052 failure: failure(QosErrorCode::ResourceUnavailable),
1053 }]
1054 );
1055 assert_eq!(
1056 controller.handle_message(&error_notify(
1057 &failed_request,
1058 QosErrorCode::ResourceUnavailable,
1059 )),
1060 QosTransition::default()
1061 );
1062
1063 let torn_down_request = request(
1064 49,
1065 QosReservationSetup::Listen {
1066 confirmation_required: true,
1067 },
1068 );
1069 let (torn_down_id, _) = controller.start(torn_down_request.clone(), now).unwrap();
1070 let torn_down = controller.handle_message(&error_notify(
1071 &torn_down_request,
1072 QosErrorCode::ReservationTornDown,
1073 ));
1074 assert!(torn_down.messages().is_empty());
1075 assert_eq!(
1076 torn_down.events(),
1077 &[QosReservationEvent::TornDown { id: torn_down_id }]
1078 );
1079 assert!(controller.is_empty());
1080 }
1081
1082 #[test]
1083 fn fragmented_and_coalesced_service_frames_keep_exact_reservation_ownership() {
1084 let now = Instant::now();
1085 let protocol = ProtocolVersion::V22;
1086 let mut controller = QosReservationController::new(limits()).unwrap();
1087 let first = request(
1088 50,
1089 QosReservationSetup::Listen {
1090 confirmation_required: true,
1091 },
1092 );
1093 let second = request(
1094 51,
1095 QosReservationSetup::Listen {
1096 confirmation_required: true,
1097 },
1098 );
1099 let (first_id, _) = controller.start(first.clone(), now).unwrap();
1100 let (second_id, _) = controller.start(second.clone(), now).unwrap();
1101 let mut wrong = first.clone();
1102 wrong.direction = QosDirection::Send;
1103
1104 let bytes = [
1105 reservation_notify(&wrong).encode(protocol).unwrap(),
1106 reservation_notify(&first).encode(protocol).unwrap(),
1107 reservation_notify(&second).encode(protocol).unwrap(),
1108 ]
1109 .concat();
1110 let mut decoder = FrameDecoder::new();
1111 let mut events = Vec::new();
1112 for fragment in bytes.chunks(7) {
1113 for frame in decoder.push(fragment).unwrap() {
1114 let message = ControlMessage::decode(frame, protocol).unwrap();
1115 events.extend(controller.handle_message(&message).into_parts().1);
1116 }
1117 }
1118 assert_eq!(
1119 events,
1120 [
1121 QosReservationEvent::Established { id: first_id },
1122 QosReservationEvent::Established { id: second_id },
1123 ]
1124 );
1125 assert_eq!(
1126 controller.state(first_id),
1127 Some(QosReservationState::Active)
1128 );
1129 assert_eq!(
1130 controller.state(second_id),
1131 Some(QosReservationState::Active)
1132 );
1133 }
1134
1135 #[test]
1136 fn generation_capacity_never_evicts_a_stale_correlation_tombstone() {
1137 let now = Instant::now();
1138 let mut bounded_limits = limits();
1139 bounded_limits.maximum_live_reservations = 1;
1140 bounded_limits.maximum_generations = 2;
1141 let mut controller = QosReservationController::new(bounded_limits).unwrap();
1142 let first = request(
1143 60,
1144 QosReservationSetup::Listen {
1145 confirmation_required: false,
1146 },
1147 );
1148 let (first_id, _) = controller.start(first.clone(), now).unwrap();
1149 controller.teardown(first_id).unwrap();
1150 assert!(matches!(
1151 controller.start(first, now),
1152 Err(QosReservationError::FlowIdentityReused)
1153 ));
1154
1155 let second = request(
1156 61,
1157 QosReservationSetup::Listen {
1158 confirmation_required: false,
1159 },
1160 );
1161 let (second_id, _) = controller.start(second, now).unwrap();
1162 controller.teardown(second_id).unwrap();
1163 assert!(matches!(
1164 controller.start(
1165 request(
1166 62,
1167 QosReservationSetup::Listen {
1168 confirmation_required: false,
1169 },
1170 ),
1171 now,
1172 ),
1173 Err(QosReservationError::GenerationCapacityExhausted)
1174 ));
1175 }
1176
1177 #[test]
1178 fn modification_timeout_and_drain_are_deterministic_and_idempotent() {
1179 let now = Instant::now();
1180 let mut controller = QosReservationController::new(limits()).unwrap();
1181 let first = request(
1182 46,
1183 QosReservationSetup::Listen {
1184 confirmation_required: false,
1185 },
1186 );
1187 let second = request(
1188 47,
1189 QosReservationSetup::Listen {
1190 confirmation_required: false,
1191 },
1192 );
1193 let (first_id, first_start) = controller.start(first.clone(), now).unwrap();
1194 let (second_id, _) = controller.start(second.clone(), now).unwrap();
1195 assert_eq!(
1196 first_start.events(),
1197 &[QosReservationEvent::ActiveWithoutConfirmation { id: first_id }]
1198 );
1199 controller
1200 .modify(
1201 first_id,
1202 traffic(80_000),
1203 application(),
1204 Duration::from_secs(5),
1205 now,
1206 )
1207 .unwrap();
1208 let timed_out = controller.poll(now + Duration::from_secs(5));
1209 assert_eq!(
1210 timed_out.events(),
1211 &[QosReservationEvent::ModificationOutcomeUnknown { id: first_id }]
1212 );
1213 assert!(timed_out.messages().is_empty());
1214 assert_eq!(
1215 controller.state(first_id),
1216 Some(QosReservationState::Active)
1217 );
1218
1219 let drained = controller.drain();
1220 assert_eq!(
1221 drained.messages(),
1222 &[
1223 ControlMessage::QosTeardown {
1224 flow: first.flow,
1225 direction: first.direction,
1226 },
1227 ControlMessage::QosTeardown {
1228 flow: second.flow,
1229 direction: second.direction,
1230 },
1231 ]
1232 );
1233 assert_eq!(
1234 drained.events(),
1235 &[
1236 QosReservationEvent::TornDown { id: first_id },
1237 QosReservationEvent::TornDown { id: second_id },
1238 ]
1239 );
1240 assert_eq!(controller.drain(), QosTransition::default());
1241 }
1242}