1use std::sync::Arc;
11
12use o_sfu_rfc::webrtc::sdp;
13use o_sfu_router::rtp::MediaStream as RouterRtpParameters;
14use str0m::{
15 change::{SdpAnswer, SdpOffer},
16 media::{KeyframeRequestKind, MediaKind, Rid},
17};
18use tokio::sync::{mpsc, oneshot};
19
20use super::{
21 codec::{ParsedAnswerRids, RepairSummary, validate_answer_sdp},
22 state::{
23 relay_registry::{RelayPacketMailbox, RelayTargetId},
24 route_control::PacketLayerGate,
25 },
26};
27use crate::engine::{
28 media_transport::{
29 ActiveSpeakerSource, AppliedSessionAnswer, ConsumerRouteControl,
30 ConsumerRouteControlOutcome, ProducerRouteControl, ReceiverBweTargetUpdate, SessionOffer,
31 SessionUploadSlot, SourceActivityUpdate, TransportAdapterError, TransportConsumerRoute,
32 TransportMediaId, TransportResult, TransportSessionKey, TransportSourceDiagnosticsSnapshot,
33 TransportSourceKey,
34 },
35 metrics::{
36 RtcMetricsRecorder, RtcRemoteControlDropKind, RtcRemotePacketGateConvergence,
37 RtcWorkerObservationKind,
38 },
39};
40
41#[derive(Debug, Clone)]
52pub struct RemoteSourceControl {
53 tx: mpsc::Sender<RtcWorkerCommand>,
54 target_id: RelayTargetId,
55 rtc_metrics: Arc<RtcMetricsRecorder>,
56}
57
58#[derive(Debug, Clone, Copy, PartialEq, Eq)]
59pub(super) enum RemoteControlSendOutcome {
60 Forwarded,
61 Full,
62 Closed,
63}
64
65impl RemoteSourceControl {
66 pub(super) fn new(
68 tx: mpsc::Sender<RtcWorkerCommand>,
69 target_id: RelayTargetId,
70 rtc_metrics: Arc<RtcMetricsRecorder>,
71 ) -> Self {
72 Self {
73 tx,
74 target_id,
75 rtc_metrics,
76 }
77 }
78
79 pub(super) fn request_kf(
84 &self,
85 source: &TransportSourceKey,
86 rid: Option<Rid>,
87 kind: KeyframeRequestKind,
88 ) -> RemoteControlSendOutcome {
89 self.send_command(
90 RouteControlRequest::RequestRemoteKeyframe {
91 source: source.clone(),
92 target_id: self.target_id,
93 rid,
94 kind,
95 },
96 RtcRemoteControlDropKind::Keyframe,
97 )
98 }
99
100 pub(super) fn set_pkt_gate(
104 &self,
105 source: &TransportSourceKey,
106 packet_gate: PacketLayerGate,
107 ) -> RemoteControlSendOutcome {
108 self.send_command(
109 RouteControlRequest::SetRemoteSourcePacketGate {
110 source: source.clone(),
111 target_id: self.target_id,
112 packet_gate,
113 },
114 RtcRemoteControlDropKind::PacketGate,
115 )
116 }
117
118 pub(super) fn record_pkt_gate_retry(&self) {
119 self.rtc_metrics
120 .record_rtc_remote_packet_gate_convergence(RtcRemotePacketGateConvergence::Retry);
121 }
122
123 pub(super) fn record_pkt_gate_flushed(&self) {
124 self.rtc_metrics
125 .record_rtc_remote_packet_gate_convergence(RtcRemotePacketGateConvergence::Flushed);
126 }
127
128 #[inline]
129 fn send_command(
130 &self,
131 request: RouteControlRequest,
132 drop_kind: RtcRemoteControlDropKind,
133 ) -> RemoteControlSendOutcome {
134 match self.tx.try_send(RtcWorkerCommand::RouteControl {
135 request,
136 response: None,
137 }) {
138 Ok(()) => RemoteControlSendOutcome::Forwarded,
139 Err(mpsc::error::TrySendError::Full(_)) => {
140 self.rtc_metrics.record_rtc_remote_control_drop(drop_kind);
141 RemoteControlSendOutcome::Full
142 }
143 Err(mpsc::error::TrySendError::Closed(_)) => {
144 self.rtc_metrics.record_rtc_remote_control_drop(drop_kind);
145 RemoteControlSendOutcome::Closed
146 }
147 }
148 }
149}
150
151pub type RtcWorkerResponse<T> = oneshot::Sender<TransportResult<T>>;
156
157pub struct ParsedSessionAnswer {
159 pub(super) answer: SdpAnswer,
160 pub(super) rids: ParsedAnswerRids,
161 pub(super) repair: RepairSummary,
162}
163
164impl ParsedSessionAnswer {
165 pub(in crate::engine::media_transport) fn parse(answer_sdp: &str) -> TransportResult<Self> {
172 let repair = validate_answer_sdp(answer_sdp)?;
173 let answer = SdpAnswer::from_sdp_string(answer_sdp)
174 .map_err(|_error| TransportAdapterError::InvalidInput)?;
175 let rids = ParsedAnswerRids::parse(answer_sdp, &answer);
176 Ok(Self {
177 answer,
178 rids,
179 repair,
180 })
181 }
182}
183
184pub struct RtcSessionOffer {
185 offer: SdpOffer,
186 upload_slots: Vec<SessionUploadSlot>,
187}
188
189impl RtcSessionOffer {
190 pub(super) fn new(offer: SdpOffer, upload_slots: Vec<SessionUploadSlot>) -> Self {
191 Self {
192 offer,
193 upload_slots,
194 }
195 }
196
197 pub(in crate::engine::media_transport) fn into_session_offer(self) -> SessionOffer {
198 let mut sdp = self.offer.to_sdp_string();
199 let media_line_start = sdp.find("\r\nm=").map_or(sdp.len(), |index| index + 2);
200 sdp.insert_str(media_line_start, sdp::EOC_LINE);
201 SessionOffer::new(sdp).with_upload_slots(self.upload_slots)
202 }
203}
204
205#[derive(Debug)]
206pub enum WorkerMediaControlBatch {
207 ReceiverBwe(Vec<(usize, ReceiverBweTargetUpdate)>),
208 ProducerActivity(Vec<(usize, ProducerRouteControl)>),
209 ConsumerGates {
210 source: TransportSourceKey,
211 updates: Vec<(usize, TransportConsumerRoute, PacketLayerGate)>,
212 },
213 ConsumerFollowUp(Vec<(usize, ConsumerRouteControl)>),
214}
215
216#[derive(Debug)]
217pub enum WorkerMediaControlBatchOutcome {
218 Applied(Vec<TransportResult<()>>),
219 Consumers(Vec<ConsumerRouteControlOutcome>),
220}
221
222pub enum RouteControlRequest {
223 AddRelayTarget {
224 source: TransportSourceKey,
225 target_id: RelayTargetId,
226 target: RelayPacketMailbox,
227 },
228 RemoveRelayTarget {
229 source: TransportSourceKey,
230 target_id: RelayTargetId,
231 },
232 SetRelayTargetActive {
233 source: TransportSourceKey,
234 target_id: RelayTargetId,
235 active: bool,
236 },
237 SetRemoteSourceActivity {
238 source: TransportSourceKey,
239 update: SourceActivityUpdate,
240 },
241 RequestRemoteKeyframe {
242 source: TransportSourceKey,
243 target_id: RelayTargetId,
244 rid: Option<Rid>,
245 kind: KeyframeRequestKind,
246 },
247 SetRemoteSourcePacketGate {
248 source: TransportSourceKey,
249 target_id: RelayTargetId,
250 packet_gate: PacketLayerGate,
251 },
252}
253
254pub enum RtcWorkerCommand {
261 CreateInitialSessionOffer {
267 room_id: Arc<str>,
268 session_key: TransportSessionKey,
269 response: RtcWorkerResponse<RtcSessionOffer>,
270 },
271 CreateSessionRenegotiationOffer {
278 session_key: TransportSessionKey,
279 response: RtcWorkerResponse<RtcSessionOffer>,
280 },
281 ActiveSpeakerSourceSnapshot {
286 response: RtcWorkerResponse<Vec<ActiveSpeakerSource>>,
287 },
288 SourceDiagnosticsSnapshot {
290 transport_media_ids: Vec<TransportMediaId>,
291 response: RtcWorkerResponse<TransportSourceDiagnosticsSnapshot>,
292 },
293 ApplySessionAnswer {
299 session_key: TransportSessionKey,
300 answer: ParsedSessionAnswer,
301 response: RtcWorkerResponse<AppliedSessionAnswer>,
302 },
303 CloseSession {
308 session_key: TransportSessionKey,
309 response: RtcWorkerResponse<()>,
310 },
311 RemoveMedia {
318 session_key: TransportSessionKey,
319 transport_media_id: TransportMediaId,
320 response: RtcWorkerResponse<()>,
321 },
322 #[cfg(test)]
328 ResolveNegotiatedProducerParameters {
329 session_key: TransportSessionKey,
330 transport_media_id: TransportMediaId,
331 response: RtcWorkerResponse<RouterRtpParameters>,
332 },
333 #[cfg(test)]
338 ResolveMediaMid {
339 transport_media_id: TransportMediaId,
340 response: RtcWorkerResponse<Option<String>>,
341 },
342 AddRecvMedia {
350 session_key: TransportSessionKey,
351 media_kind: MediaKind,
352 rtp_parameters: RouterRtpParameters,
353 response: RtcWorkerResponse<TransportMediaId>,
354 },
355 AddSendMedia {
363 consumer_key: TransportSessionKey,
364 media_kind: MediaKind,
365 source: TransportSourceKey,
366 remote_source_control: Option<RemoteSourceControl>,
367 consumer_rtp_parameters: RouterRtpParameters,
368 active: bool,
369 response: RtcWorkerResponse<(TransportMediaId, String)>,
370 },
371 ApplyMediaControlBatch {
372 batch: WorkerMediaControlBatch,
373 response: RtcWorkerResponse<WorkerMediaControlBatchOutcome>,
374 },
375 RouteControl {
376 request: RouteControlRequest,
377 response: Option<RtcWorkerResponse<()>>,
378 },
379}
380
381impl RtcWorkerCommand {
382 pub(super) fn observation_kind(&self) -> Option<RtcWorkerObservationKind> {
384 match self {
385 Self::ActiveSpeakerSourceSnapshot { .. } => {
386 Some(RtcWorkerObservationKind::ActiveSpeakerSources)
387 }
388 Self::SourceDiagnosticsSnapshot { .. } => {
389 Some(RtcWorkerObservationKind::SourceDiagnostics)
390 }
391 #[cfg(test)]
392 Self::ResolveMediaMid { .. } => Some(RtcWorkerObservationKind::ResolveMediaMid),
393 #[cfg(test)]
394 Self::ResolveNegotiatedProducerParameters { .. } => {
395 Some(RtcWorkerObservationKind::NegotiatedProducerParameters)
396 }
397 Self::CreateInitialSessionOffer { .. }
398 | Self::CreateSessionRenegotiationOffer { .. }
399 | Self::ApplySessionAnswer { .. }
400 | Self::CloseSession { .. }
401 | Self::RemoveMedia { .. }
402 | Self::AddRecvMedia { .. }
403 | Self::AddSendMedia { .. }
404 | Self::ApplyMediaControlBatch { .. }
405 | Self::RouteControl { .. } => None,
406 }
407 }
408
409 pub(super) fn may_change_demux_topology(&self) -> bool {
410 match self {
411 Self::CreateInitialSessionOffer { .. }
412 | Self::ApplySessionAnswer { .. }
413 | Self::CloseSession { .. } => true,
414 Self::CreateSessionRenegotiationOffer { .. }
415 | Self::ActiveSpeakerSourceSnapshot { .. }
416 | Self::SourceDiagnosticsSnapshot { .. }
417 | Self::RemoveMedia { .. }
418 | Self::AddRecvMedia { .. }
419 | Self::AddSendMedia { .. }
420 | Self::ApplyMediaControlBatch { .. }
421 | Self::RouteControl { .. } => false,
422 #[cfg(test)]
423 Self::ResolveMediaMid { .. } => false,
424 #[cfg(test)]
425 Self::ResolveNegotiatedProducerParameters { .. } => false,
426 }
427 }
428}