Skip to main content

o_sfu/application/user_session/
media.rs

1use std::collections::BTreeMap;
2
3use o_sfu_protocol::{
4    host::NegotiationKind,
5    wire::{
6        ClientResponse, DownloadStates, NegotiationUploadEncoding, NegotiationUploadSlot,
7        RequestId, ServerEnvelope, ServerRequest, SessionDescriptionPayload, StreamType, UserId,
8    },
9};
10use tracing::{Span, field, instrument, warn};
11
12use super::{User, UserError, UserOutput};
13use crate::{
14    application::stream_catalog::{DiscussStream, source_publish_intent_for_stream_type},
15    core::prelude::{
16        Bitrate, MediaSession, NegotiationOffer, SessionError, SfuCoreError,
17        SourceDeactivateIntent, SourcePublishIntent,
18    },
19    runtime::telemetry::schema::event as telemetry_event,
20};
21
22impl User {
23    pub(super) async fn complete_negotiation(
24        &mut self,
25        response_to: RequestId,
26        response: ClientResponse,
27    ) -> Result<UserOutput, UserError> {
28        let result = self.media.answer(response_to, response).await;
29        result.map_err(|e| self.answer_error(e))
30    }
31
32    #[instrument(
33        name = "transport.renegotiate",
34        skip_all,
35        fields(
36            room_id = %self.room_id(),
37            user_id = %self.user_id().path_segment(),
38            connection_id = self.connection_id().as_u64()
39        )
40    )]
41    pub(super) async fn renegotiate(&mut self) -> Result<UserOutput, UserError> {
42        self.reject_stale_connection().await?;
43        let result = self.media.renegotiate().await;
44        result.map_err(|e| self.negotiation_error(NegotiationKind::Renegotiate, None, e))
45    }
46
47    #[instrument(
48        name = "publish.intent",
49        skip_all,
50        fields(
51            room_id = %self.room_id(),
52            user_id = %self.user_id().path_segment(),
53            connection_id = self.connection_id().as_u64(),
54            ?stream_type,
55            active
56        )
57    )]
58    pub(super) async fn set_publication_active(
59        &mut self,
60        stream_type: StreamType,
61        active: bool,
62    ) -> Result<UserOutput, UserError> {
63        if active {
64            let intent = source_publish_intent_for_stream_type(stream_type);
65            return self
66                .media
67                .publish(intent)
68                .await
69                .map_err(|error| self.publish_error(stream_type, error));
70        }
71        let intent = DiscussStream::for_type(stream_type).deactivate_intent();
72        Ok(self.media.deactivate_publication(intent).await)
73    }
74
75    #[instrument(
76        name = "subscribe.intent",
77        skip_all,
78        fields(
79            room_id = %self.room_id(),
80            user_id = %self.user_id().path_segment(),
81            connection_id = self.connection_id().as_u64(),
82            target_session_id = field::Empty,
83            source_count = field::Empty
84        )
85    )]
86    pub(super) async fn subscribe(
87        &self,
88        target_user_id: UserId,
89        states: DownloadStates,
90    ) -> Result<UserOutput, UserError> {
91        let target_user_id = target_user_id.normalized_for_runtime();
92        let span = Span::current();
93        span.record(
94            "target_session_id",
95            field::display(target_user_id.path_segment()),
96        );
97        let source_intents = DiscussStream::all()
98            .into_iter()
99            .filter_map(|stream| stream.subscription_intent_if_requested(&states))
100            .collect::<BTreeMap<_, _>>();
101        span.record("source_count", source_intents.len());
102        let session = self.media.session();
103        let result = session.subscribe(&target_user_id, &source_intents).await;
104        result.map_err(|error| self.subscribe_error(&target_user_id, error))?;
105        Ok(UserOutput::new())
106    }
107
108    #[instrument(
109        name = "transport.offer.create",
110        skip_all,
111        fields(
112            room_id = %self.room_id(),
113            user_id = %self.user_id().path_segment(),
114            connection_id = self.connection_id().as_u64()
115        )
116    )]
117    pub(super) async fn run_initial_offer(&mut self) -> Result<UserOutput, UserError> {
118        let result = self.media.establish().await;
119        result.map_err(|e| self.negotiation_error(NegotiationKind::Offer, None, e))
120    }
121}
122
123pub(super) struct ServerMediaNegotiation {
124    session: MediaSession,
125    next: u64,
126    pending: Option<(RequestId, NegotiationKind)>,
127}
128
129enum AnswerError {
130    Protocol(RequestId, &'static str),
131    Session(NegotiationKind, RequestId, SessionError),
132}
133
134const EMPTY_SDP_ANSWER_LOG: &str = "received empty SDP answer for negotiation request";
135const UNKNOWN_ANSWER_LOG: &str = "received negotiation answer for an unknown or stale request";
136
137impl ServerMediaNegotiation {
138    pub(super) fn new(session: MediaSession) -> Self {
139        Self {
140            session,
141            next: 0,
142            pending: None,
143        }
144    }
145
146    pub(super) const fn session(&self) -> &MediaSession {
147        &self.session
148    }
149
150    async fn establish(&mut self) -> Result<UserOutput, SessionError> {
151        let offer = self.session.establish().await?;
152        Ok(self.issue(NegotiationKind::Offer, offer))
153    }
154
155    async fn renegotiate(&mut self) -> Result<UserOutput, SessionError> {
156        let offer = self.session.renegotiate().await?;
157        Ok(self.issue(NegotiationKind::Renegotiate, offer))
158    }
159
160    async fn publish(&mut self, intent: SourcePublishIntent) -> Result<UserOutput, SessionError> {
161        let offer = self.session.publish(intent).await?;
162        Ok(self.issue(NegotiationKind::Renegotiate, offer))
163    }
164
165    async fn deactivate_publication(&mut self, intent: SourceDeactivateIntent) -> UserOutput {
166        self.session.deactivate_publication(intent).await;
167        UserOutput::new()
168    }
169
170    async fn answer(
171        &mut self,
172        response_to: RequestId,
173        response: ClientResponse,
174    ) -> Result<UserOutput, AnswerError> {
175        let (kind, answer) = match response {
176            ClientResponse::Offer(answer) => (NegotiationKind::Offer, answer),
177            ClientResponse::Renegotiate(answer) => (NegotiationKind::Renegotiate, answer),
178        };
179        if answer.sdp.is_empty() {
180            return Err(AnswerError::Protocol(response_to, EMPTY_SDP_ANSWER_LOG));
181        }
182        if !self.expects(&response_to, kind) {
183            return Err(AnswerError::Protocol(response_to, UNKNOWN_ANSWER_LOG));
184        }
185        let offer = match self.session.answer(&answer.sdp).await {
186            Ok(offer) => offer,
187            Err(error) => return Err(AnswerError::Session(kind, response_to, error)),
188        };
189        self.pending = None;
190        Ok(self.issue(NegotiationKind::Renegotiate, offer))
191    }
192
193    pub(super) async fn close(&mut self) {
194        self.session.close().await;
195    }
196
197    fn issue(&mut self, kind: NegotiationKind, offer: Option<NegotiationOffer>) -> UserOutput {
198        let Some(offer) = offer else {
199            return UserOutput::new();
200        };
201        let req_id = RequestId::new(format!("server-{}", self.next));
202        self.next = self.next.saturating_add(1);
203        let payload = session_description_payload(offer);
204        let request = match kind {
205            NegotiationKind::Offer => ServerRequest::Offer(payload),
206            NegotiationKind::Renegotiate => ServerRequest::Renegotiate(payload),
207        };
208        self.pending = Some((req_id.clone(), kind));
209        vec![ServerEnvelope::Request {
210            request_id: req_id,
211            request,
212        }]
213    }
214
215    fn expects(&self, id: &RequestId, kind: NegotiationKind) -> bool {
216        matches!(self.pending.as_ref(), Some((req_id, pending_kind)) if req_id == id && *pending_kind == kind)
217    }
218}
219
220fn session_description_payload(offer: NegotiationOffer) -> SessionDescriptionPayload {
221    SessionDescriptionPayload {
222        sdp: offer.sdp,
223        upload_slots: offer
224            .upload_slots
225            .into_iter()
226            .map(|slot| NegotiationUploadSlot {
227                mid: slot.mid,
228                kind: slot.kind,
229                codecs: slot.codecs,
230                simulcast_encodings: slot
231                    .simulcast_encodings
232                    .into_iter()
233                    .map(|encoding| NegotiationUploadEncoding {
234                        rid: encoding.rid,
235                        max_bitrate: encoding.max_bitrate.map(Bitrate::as_bps),
236                        resolution_scale: encoding.resolution_scale,
237                        max_framerate: encoding.max_framerate,
238                    })
239                    .collect(),
240            })
241            .collect(),
242    }
243}
244
245impl User {
246    fn answer_error(&self, error: AnswerError) -> UserError {
247        match error {
248            AnswerError::Protocol(response_to, message) => {
249                warn!(
250                    user_id = ?self.user_id(),
251                    connection_id = ?self.connection_id(),
252                    remote_address = self.remote_address.as_ref(),
253                    ?response_to,
254                    "{message}"
255                );
256                UserError::ProtocolViolation
257            }
258            AnswerError::Session(kind, response_to, error) => {
259                self.negotiation_error(kind, Some(&response_to), error)
260            }
261        }
262    }
263
264    fn negotiation_error(
265        &self,
266        kind: NegotiationKind,
267        response_to: Option<&RequestId>,
268        error: SessionError,
269    ) -> UserError {
270        let operation = match kind {
271            NegotiationKind::Offer => "initial_offer_create",
272            NegotiationKind::Renegotiate => "renegotiation_offer_create",
273        };
274        let outcome = match error {
275            SessionError::NoPendingRequest => "no_pending_media_request",
276            SessionError::Core(error) if error.is_client_error() => "client_negotiation_error",
277            SessionError::Core(_) => "transport_error",
278        };
279        warn!(
280            event = telemetry_event::NEGOTIATION_FAILED,
281            operation,
282            outcome,
283            user_id = ?self.user_id(),
284            connection_id = ?self.connection_id(),
285            remote_address = self.remote_address.as_ref(),
286            response_to = ?response_to,
287            ?error,
288            "media session command failed"
289        );
290        user_error(error)
291    }
292
293    fn publish_error(&self, stream_type: StreamType, error: SessionError) -> UserError {
294        warn!(
295            event = telemetry_event::PUBLISH_ABORTED,
296            operation = "publish_intent",
297            outcome = "publish_rejected",
298            user_id = ?self.user_id(),
299            connection_id = ?self.connection_id(),
300            remote_address = self.remote_address.as_ref(),
301            ?stream_type,
302            ?error,
303            "media session command failed"
304        );
305        user_error(error)
306    }
307
308    fn subscribe_error(&self, target_user_id: &UserId, error: SessionError) -> UserError {
309        let outcome = match error {
310            SessionError::Core(SfuCoreError::SubscriptionUpdateRejected) => "stale_connection",
311            SessionError::NoPendingRequest | SessionError::Core(_) => "subscription_failed",
312        };
313        warn!(
314            event = telemetry_event::SUBSCRIBE_REJECTED,
315            operation = "consume_prepare",
316            outcome,
317            user_id = ?self.user_id(),
318            connection_id = ?self.connection_id(),
319            remote_address = self.remote_address.as_ref(),
320            ?target_user_id,
321            ?error,
322            "media session command failed"
323        );
324        user_error(error)
325    }
326}
327
328fn user_error(error: SessionError) -> UserError {
329    if error.is_client_error() {
330        UserError::ProtocolViolation
331    } else {
332        UserError::InternalError
333    }
334}