Skip to main content

o_sfu/application/user_session/
lifecycle.rs

1use std::sync::Arc;
2
3use o_sfu_protocol::wire::{ServerEnvelope, ServerMessage, UserId, WelcomePayload};
4use tracing::{debug, error};
5
6use super::{ServerMediaNegotiation, User, UserError, UserOutput};
7use crate::{
8    core::prelude::{MediaSession, TransportSessionHealth},
9    runtime::ConnectionId,
10};
11
12impl User {
13    #[must_use]
14    pub fn new(session: MediaSession, remote_address: Arc<str>) -> Self {
15        Self {
16            remote_address,
17            media: ServerMediaNegotiation::new(session),
18            cleanup_finished: false,
19        }
20    }
21
22    pub(crate) fn room_id(&self) -> &str {
23        self.media.session().room_id()
24    }
25
26    pub(crate) const fn connection_id(&self) -> ConnectionId {
27        self.media.session().connection_id()
28    }
29
30    pub(crate) fn user_id(&self) -> &UserId {
31        self.media.session().user_id()
32    }
33
34    pub(crate) fn remote_address(&self) -> &str {
35        self.remote_address.as_ref()
36    }
37
38    pub(crate) async fn is_current_connection(&self) -> bool {
39        self.media.session().is_current_connection().await
40    }
41
42    #[must_use]
43    pub fn transport_disconnected(&self) -> bool {
44        self.media.session().endpoint_health() == Some(TransportSessionHealth::Disconnected)
45    }
46
47    pub async fn start(&mut self) -> Result<UserOutput, UserError> {
48        let welcome = WelcomePayload {
49            features: self.media.session().available_features(),
50            recording: self.media.session().recording_state().await,
51            peers: self.media.session().peer_snapshots().await,
52        };
53        self.reject_stale_connection().await?;
54        let mut output = vec![ServerEnvelope::Message(ServerMessage::Welcome(welcome))];
55        output.extend(self.run_initial_offer().await?);
56        Ok(output)
57    }
58
59    pub async fn close(&mut self) {
60        self.media.close().await;
61        self.cleanup_finished = true;
62    }
63
64    pub(super) async fn reject_stale_connection(&self) -> Result<(), UserError> {
65        if self.is_current_connection().await {
66            return Ok(());
67        }
68        debug!(
69            user_id = ?self.user_id(),
70            connection_id = ?self.connection_id(),
71            "rejecting intent from a stale user connection"
72        );
73        Err(UserError::Kicked)
74    }
75}
76
77impl Drop for User {
78    fn drop(&mut self) {
79        if self.cleanup_finished {
80            return;
81        }
82        error!(
83            user_id = ?self.user_id(),
84            connection_id = ?self.connection_id(),
85            "dropped websocket user without completing explicit cleanup"
86        );
87        debug_assert!(
88            self.cleanup_finished,
89            "websocket user dropped before explicit cleanup completed"
90        );
91    }
92}