o_sfu/application/user_session/
lifecycle.rs1use 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}