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}