Skip to main content

o_sfu_protocol/core/
connection_lifecycle.rs

1//! socket lifecycle transitions for [`ProtocolCore`]
2//!
3//! [`ProtocolCore::connect`] starts a user attempt and clears replayable intent
4//! [`ProtocolCore::disconnect`] ends that attempt and suppresses recovery
5//! [`ProtocolCore::on_ws_close`] maps terminal codes to [`ConnectionState::Closed`]
6//! while transient closes preserve the connect context for [`handle_recovery_timer`]
7//!
8//! welcome messages enter through [`ProtocolCore::on_ws_message`]
9//! transport readiness enters through [`ProtocolCore::on_transport_ready`]
10//! each transition returns ordered [`Command`] values for the host
11
12use 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    /// Starts a fresh connection attempt when the current state permits one.
20    ///
21    /// Accepts [`ConnectionState::Disconnected`], [`ConnectionState::Closed`] and
22    /// [`ConnectionState::Recovering`]. Calls from [`ConnectionState::Connecting`],
23    /// [`ConnectionState::Authenticated`] and [`ConnectionState::Connected`] return
24    /// no commands without replacing the saved admission context.
25    ///
26    /// Clears sticky replay and runtime state so a caller switching rooms or credentials cannot
27    /// accidentally leak the previous user intent into the new connection.
28    /// A call from [`ConnectionState::Recovering`] cancels its recovery timer
29    /// before the new socket attempt starts.
30    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    /// ends the current user attempt on purpose
58    ///
59    /// unlike [`ProtocolCore::on_ws_close`], this is not a recovery path
60    /// it clears the saved connect context, runtime state and sticky replay state,
61    /// then closes the websocket and peer connection
62    /// any later recovery-timer delivery becomes a no-op because the caller
63    /// explicitly asked to stop
64    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    /// handles websocket closure after a user was already in flight
89    ///
90    /// there are three different cases here and mixing them up is the main way to
91    /// break reconnect behavior:
92    ///
93    /// - terminal close codes move to [`ConnectionState::Closed`], clear the saved connect context,
94    ///   and suppress recovery
95    /// - non-terminal closes with saved connect context move to [`ConnectionState::Recovering`] and
96    ///   schedule the recovery timer
97    /// - non-terminal closes without saved connect context fall back to
98    ///   [`ConnectionState::Disconnected`], because there is nothing safe to reconnect to
99    ///
100    /// Queue overflow uses [`WebSocketCloseCode::Overloaded`] so recovery retains
101    /// publication and subscription intent after a temporary backlog.
102    ///
103    /// example:
104    ///
105    /// ```text
106    /// Connected --on_ws_close(AuthFailed)--> Closed
107    /// Connected --on_ws_close(1011)--> Recovering
108    /// ```
109    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
159/// retries the saved websocket connection after a recovery delay
160///
161/// this is narrow
162/// only [`ConnectionState::Recovering`] may consume the recovery timer
163/// a stale timer firing after a successful reconnect or explicit
164/// disconnect must do nothing, otherwise old scheduled work can restart an
165/// inactive attempt
166///
167/// example:
168///
169/// ```text
170/// Connected --on_ws_close(1011)--> Recovering
171/// Recovering --handle_recovery_timer()--> Connecting
172/// Connected --handle_recovery_timer()--> no-op
173/// ```
174pub(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
206/// Returns the compatibility cause label for a terminal WebSocket close code.
207fn 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}