moqtap_client/draft18/subscription.rs
1/// Subscription lifecycle states.
2#[derive(Debug, Clone, Copy, PartialEq, Eq)]
3pub enum SubscriptionState {
4 /// Initial state before any SUBSCRIBE message is sent.
5 Idle,
6 /// SUBSCRIBE has been sent; awaiting OK or ERROR.
7 Subscribing,
8 /// Subscription is accepted and data may be flowing.
9 Active,
10 /// Subscription has ended (error, cancellation, or PUBLISH_DONE).
11 Done,
12}
13
14/// Errors that can occur during subscription state transitions.
15#[derive(Debug, thiserror::Error, PartialEq, Eq)]
16pub enum SubscriptionError {
17 /// An event was received that is not valid for the current state.
18 #[error("invalid transition from {from:?} on event {event}")]
19 InvalidTransition {
20 /// The state the machine was in when the invalid event arrived.
21 from: SubscriptionState,
22 /// The name of the event that was rejected.
23 event: String,
24 },
25}
26
27/// Pure state machine for a MoQT subscription.
28/// Transitions: Idle -> Subscribing -> Active -> Done.
29pub struct SubscriptionStateMachine {
30 state: SubscriptionState,
31}
32
33impl Default for SubscriptionStateMachine {
34 fn default() -> Self {
35 Self::new()
36 }
37}
38
39impl SubscriptionStateMachine {
40 /// Creates a new state machine in the [`SubscriptionState::Idle`] state.
41 pub fn new() -> Self {
42 Self { state: SubscriptionState::Idle }
43 }
44
45 /// Returns the current state of the subscription.
46 pub fn state(&self) -> SubscriptionState {
47 self.state
48 }
49
50 /// Idle -> Subscribing (SUBSCRIBE sent).
51 pub fn on_subscribe_sent(&mut self) -> Result<(), SubscriptionError> {
52 if self.state == SubscriptionState::Idle {
53 self.state = SubscriptionState::Subscribing;
54 Ok(())
55 } else {
56 Err(SubscriptionError::InvalidTransition {
57 from: self.state,
58 event: "on_subscribe_sent".to_string(),
59 })
60 }
61 }
62
63 /// Subscribing -> Active (SUBSCRIBE_OK received).
64 pub fn on_subscribe_ok(&mut self) -> Result<(), SubscriptionError> {
65 if self.state == SubscriptionState::Subscribing {
66 self.state = SubscriptionState::Active;
67 Ok(())
68 } else {
69 Err(SubscriptionError::InvalidTransition {
70 from: self.state,
71 event: "on_subscribe_ok".to_string(),
72 })
73 }
74 }
75
76 /// Subscribing -> Done (REQUEST_ERROR received).
77 pub fn on_subscribe_error(&mut self) -> Result<(), SubscriptionError> {
78 if self.state == SubscriptionState::Subscribing {
79 self.state = SubscriptionState::Done;
80 Ok(())
81 } else {
82 Err(SubscriptionError::InvalidTransition {
83 from: self.state,
84 event: "on_subscribe_error".to_string(),
85 })
86 }
87 }
88
89 /// Subscribing | Active -> Done, Done -> Done (this subscription's request
90 /// stream was cancelled).
91 ///
92 /// This draft has no UNSUBSCRIBE message. Section 3.3.2: "Once a request
93 /// stream has been opened, the request MAY be cancelled by either endpoint."
94 ///
95 /// `Subscribing` is accepted because the precondition is the stream being
96 /// open, and it is open from the SUBSCRIBE that opened it: a subscription
97 /// can be withdrawn before it is ever answered.
98 ///
99 /// `Idle` is refused, on the other half of the same sentence: nothing has
100 /// been written, so there is no stream to terminate. `Done` stays `Done` —
101 /// nothing finishes a request stream's send half on the ordinary path, so a
102 /// caller that walks away from a request that has already ended still
103 /// resets the stream, and that reset is an ordinary end rather than a
104 /// fault.
105 pub fn on_request_cancelled(&mut self) -> Result<(), SubscriptionError> {
106 match self.state {
107 SubscriptionState::Subscribing | SubscriptionState::Active => {
108 self.state = SubscriptionState::Done;
109 Ok(())
110 }
111 SubscriptionState::Done => Ok(()),
112 SubscriptionState::Idle => Err(SubscriptionError::InvalidTransition {
113 from: self.state,
114 event: "on_request_cancelled".to_string(),
115 }),
116 }
117 }
118
119 /// REQUEST_UPDATE received -- a self-transition, from Subscribing as well as
120 /// from Active.
121 ///
122 /// Section 10.9 orders an update against the request rather than against
123 /// the request's answer: the sender of a SUBSCRIBE "can later send a
124 /// REQUEST_UPDATE on the same bidi stream as the request to modify it",
125 /// where later is later than the SUBSCRIBE. The stream is open from the
126 /// moment the SUBSCRIBE opens it.
127 ///
128 /// So a peer that sends SUBSCRIBE and REQUEST_UPDATE back to back breaks no
129 /// rule this draft states, and an update arriving before the answer leaves
130 /// the subscription where it found it. `Idle` and `Done` are still refused:
131 /// in neither does the subscription an update names exist.
132 pub fn on_subscribe_update(&mut self) -> Result<(), SubscriptionError> {
133 if matches!(self.state, SubscriptionState::Subscribing | SubscriptionState::Active) {
134 Ok(())
135 } else {
136 Err(SubscriptionError::InvalidTransition {
137 from: self.state,
138 event: "on_subscribe_update".to_string(),
139 })
140 }
141 }
142
143 /// Active -> Done (PUBLISH_DONE received).
144 pub fn on_publish_done(&mut self) -> Result<(), SubscriptionError> {
145 if self.state == SubscriptionState::Active {
146 self.state = SubscriptionState::Done;
147 Ok(())
148 } else {
149 Err(SubscriptionError::InvalidTransition {
150 from: self.state,
151 event: "on_publish_done".to_string(),
152 })
153 }
154 }
155}