o_sfu_protocol/core/
connection_lifecycle.rs1use super::{
13 Command, Commands, ConnectContext, ConnectionState, INITIAL_RECOVERY_DELAY_MS, ProtocolCore,
14 ProtocolPhase, RECOVERY_TIMER_ID, empty_features, next_recovery_delay,
15};
16use crate::{shared::RecordingState, signaling::WebSocketCloseCode};
17
18impl ProtocolCore {
19 pub fn connect(
31 &mut self,
32 url: impl Into<String>,
33 jwt: impl Into<String>,
34 room: Option<String>,
35 ) -> Vec<Command> {
36 let url = url.into();
37 let jwt = jwt.into();
38 let mut commands = match self.state() {
39 ConnectionState::Disconnected | ConnectionState::Closed => Vec::new(),
40 ConnectionState::Recovering => vec![Command::CancelTimer {
41 id: RECOVERY_TIMER_ID,
42 }],
43 _ => return Vec::new(),
44 };
45 let connect_url = url.clone();
46 self.connect_context = Some(ConnectContext { url, jwt, room });
47 self.recovery_delay_ms = INITIAL_RECOVERY_DELAY_MS;
48 self.phase = ProtocolPhase::Connecting;
49 self.clear_runtime_state();
50 self.sticky_replay.clear();
51 reset_public_state(&mut commands);
52 commands.push(state_change(self.state(), None));
53 commands.push(Command::Connect { url: connect_url });
54 commands
55 }
56
57 pub fn disconnect(&mut self) -> Vec<Command> {
65 if matches!(
66 self.state(),
67 ConnectionState::Disconnected | ConnectionState::Closed
68 ) {
69 return Vec::new();
70 }
71 self.phase = ProtocolPhase::Disconnected;
72 self.connect_context = None;
73 self.recovery_delay_ms = INITIAL_RECOVERY_DELAY_MS;
74 let mut commands = vec![Command::CancelTimer {
75 id: RECOVERY_TIMER_ID,
76 }];
77 commands.extend(self.teardown_runtime_state());
78 self.sticky_replay.clear();
79 commands.push(Command::CloseWebSocket {
80 code: u16::from(WebSocketCloseCode::Clean),
81 });
82 commands.push(Command::ClosePeerConnection);
83 reset_public_state(&mut commands);
84 commands.push(state_change(self.state(), None));
85 commands
86 }
87
88 pub fn on_ws_close(&mut self, close_code: u16) -> Vec<Command> {
110 if matches!(
111 self.state(),
112 ConnectionState::Disconnected | ConnectionState::Closed
113 ) {
114 return Vec::new();
115 }
116 if let Some(
117 terminal_code @ (WebSocketCloseCode::ProtocolError
118 | WebSocketCloseCode::AuthFailed
119 | WebSocketCloseCode::Kicked
120 | WebSocketCloseCode::RoomFull),
121 ) = WebSocketCloseCode::from_u16(close_code)
122 {
123 self.phase = ProtocolPhase::Closed;
124 self.connect_context = None;
125 self.recovery_delay_ms = INITIAL_RECOVERY_DELAY_MS;
126 let mut commands = self.teardown_runtime_state();
127 commands.push(Command::CancelTimer {
128 id: RECOVERY_TIMER_ID,
129 });
130 commands.push(Command::ClosePeerConnection);
131 reset_public_state(&mut commands);
132 commands.push(state_change(
133 self.state(),
134 terminal_close_cause(terminal_code),
135 ));
136 return commands;
137 }
138 if self.connect_context.is_none() {
139 self.phase = ProtocolPhase::Disconnected;
140 let mut commands = self.teardown_runtime_state();
141 reset_public_state(&mut commands);
142 commands.push(state_change(self.state(), None));
143 return commands;
144 }
145 let scheduled_delay_ms = self.recovery_delay_ms;
146 self.recovery_delay_ms = next_recovery_delay(scheduled_delay_ms);
147 self.phase = ProtocolPhase::Recovering;
148 let mut commands = self.teardown_runtime_state();
149 commands.push(Command::ClosePeerConnection);
150 commands.push(state_change(self.state(), None));
151 commands.push(Command::ScheduleTimer {
152 id: RECOVERY_TIMER_ID,
153 ms: scheduled_delay_ms,
154 });
155 commands
156 }
157}
158
159pub(super) fn handle_recovery_timer(core: &mut ProtocolCore) -> Commands {
175 if core.state() != ConnectionState::Recovering {
176 return Vec::new();
177 }
178 let Some(connect_context) = core.connect_context.as_ref() else {
179 return Vec::new();
180 };
181 let connect_url = connect_context.url.clone();
182 core.phase = ProtocolPhase::Connecting;
183 let mut commands = vec![state_change(core.state(), None)];
184 commands.push(Command::Connect { url: connect_url });
185 commands
186}
187
188fn reset_public_state(commands: &mut Commands) {
189 commands.extend([
190 Command::SetAvailableFeatures {
191 features: empty_features(),
192 },
193 Command::SetRecordingState {
194 state: RecordingState::default(),
195 },
196 ]);
197}
198
199fn state_change(state: ConnectionState, cause: Option<&'static str>) -> Command {
200 Command::EmitStateChange {
201 state,
202 cause: cause.map(str::to_owned),
203 }
204}
205
206fn terminal_close_cause(close_code: WebSocketCloseCode) -> Option<&'static str> {
208 match close_code {
209 WebSocketCloseCode::AuthFailed => Some("auth_failed"),
210 WebSocketCloseCode::Kicked => Some("kicked"),
211 WebSocketCloseCode::RoomFull => Some("full"),
212 _ => None,
213 }
214}