Skip to main content

sccp_protocol/
qos.rs

1//! Bounded reservation ownership for service-node QoS messages.
2//!
3//! This state machine has no station-session dependency. Callers own the
4//! service-node transport, encode messages returned in [`QosTransition`], and
5//! feed decoded reservation notifications and errors back through
6//! [`QosReservationController::handle_message`].
7
8use 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/// Resource and policy bounds applied before a reservation enters live state.
19#[derive(Clone, Copy, Debug, Eq, PartialEq)]
20pub struct QosReservationLimits {
21    pub maximum_live_reservations: usize,
22    /// Lifetime limit for unique flow/direction generations. Retired wire
23    /// identities remain reserved so a late notification cannot settle a
24    /// replacement generation.
25    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    /// Returns the unchanged limits after proving that live state and its
33    /// non-evicting correlation history can both remain bounded.
34    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/// Wire setup form selected for one reservation generation.
60#[derive(Clone, Copy, Debug, Eq, PartialEq)]
61pub enum QosReservationSetup {
62    Listen { confirmation_required: bool },
63    Path,
64}
65
66/// Shared admission and traffic policy for listen and path setup messages.
67#[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/// Complete request for one exactly correlated service-node reservation.
79///
80/// `direction` identifies the notification/error key expected from the
81/// service node. `response_timeout` is an application-selected local deadline;
82/// it does not assign units to the separate wire `retry_timer` quantity.
83#[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/// Monotonic controller identity for one non-reusable wire generation.
93#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
94pub struct QosReservationId(NonZeroU64);
95
96impl QosReservationId {
97    /// Returns the wire-independent monotonic generation number.
98    pub const fn get(self) -> u64 {
99        self.0.get()
100    }
101}
102
103/// Externally observable reservation phase.
104#[derive(Clone, Copy, Debug, Eq, PartialEq)]
105pub enum QosReservationState {
106    Establishing,
107    Active,
108    Modifying,
109}
110
111/// Exact service-node failure fields retained for policy decisions.
112#[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/// State change emitted after exact flow/direction correlation.
122#[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    /// The local observation window elapsed without an explicit failure. The
146    /// wire contract does not define a modification-success response.
147    ModificationOutcomeUnknown {
148        id: QosReservationId,
149    },
150    TornDown {
151        id: QosReservationId,
152    },
153}
154
155/// Ordered service-node writes and state events produced by one transition.
156#[derive(Clone, Debug, Default, Eq, PartialEq)]
157pub struct QosTransition {
158    messages: Vec<ControlMessage>,
159    events: Vec<QosReservationEvent>,
160}
161
162impl QosTransition {
163    /// Borrows service-node writes in their required emission order.
164    pub fn messages(&self) -> &[ControlMessage] {
165        &self.messages
166    }
167
168    /// Borrows events in the order produced by the corresponding state change.
169    pub fn events(&self) -> &[QosReservationEvent] {
170        &self.events
171    }
172
173    /// Transfers both ordered output queues to a transport/event owner.
174    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/// Failure returned before a state transition or service-node write occurs.
206#[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/// Owns bounded QoS reservation generations independently of station calls.
231#[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    /// Creates an empty controller after validating every resource bound.
242    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    /// Reserves a fresh lifetime identity and returns the exact setup message.
253    /// A flow/direction pair can never be reused by this controller.
254    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    /// Returns `None` for unknown and already retired generations.
308    pub fn state(&self, id: QosReservationId) -> Option<QosReservationState> {
309        self.by_id
310            .get(&id)
311            .map(|reservation| state(&reservation.phase))
312    }
313
314    /// Counts live generations without discarding retained correlation history.
315    pub fn len(&self) -> usize {
316        self.by_id.len()
317    }
318
319    /// Reports whether live state is empty; retired wire identities remain reserved.
320    pub fn is_empty(&self) -> bool {
321        self.by_id.is_empty()
322    }
323
324    /// Sends the sole correlatable modification for an active generation.
325    ///
326    /// The message family has no transaction identity or explicit success
327    /// response, so repeated modifications would make late failures ambiguous.
328    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    /// Builds a six-bit DSCP update for an active generation.
371    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    /// Retires the generation before returning its teardown message.
402    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    /// Handles only service-node reservation notifications and errors. Other
411    /// typed control messages leave state unchanged.
412    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    /// Expires every elapsed setup or modification observation in ID order.
445    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    /// Retires all live generations and returns teardowns in ID order.
492    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}