o_sfu_core/engine/room/transition/
publication.rs1use o_sfu_router::rtp::MediaStream as RouterRtpParameters;
24use tracing::warn;
25
26use super::super::{
27 Room, RoomUserOperation, SourcePolicyGuard,
28 effects::batch::{RoomEffectContext, RoomEffects},
29 media_graph::{ProducerActivityCommit, PublishIntentPlan, ValidatedPublish},
30};
31#[cfg(any(test, feature = "testing-transport"))]
32use crate::engine::{ConnectionId, UserId};
33use crate::engine::{
34 media_transport::{AppliedSessionAnswer, MediaTransport, TransportAdapterError},
35 source_model::{SourceDeactivateIntent, SourcePublishIntent, UserStreamId},
36};
37
38mod staging;
39#[cfg(any(test, feature = "testing-transport"))]
40#[path = "TESTS/publication_support.rs"]
41mod test_support;
42
43pub use staging::{StagedPublish, StagedPublishes};
44
45#[derive(Debug, Clone, Copy, PartialEq, Eq)]
46pub enum PublishIntentOutcome {
47 Noop,
48 Queue,
49 Activated,
50 Staged,
51}
52
53#[derive(Debug, Clone, Copy, PartialEq, Eq)]
54pub enum DeactivateIntentOutcome {
55 Noop,
56 RolledBack,
57 Deactivated,
58}
59
60#[derive(Debug, Clone, Copy, PartialEq, Eq)]
61pub enum PublishStageOutcome {
62 Staged,
63 Duplicate,
64 DuplicateAfterReservation,
65 #[cfg(test)]
66 Rejected,
67}
68
69impl RoomUserOperation<'_> {
70 #[cfg(test)]
71 pub async fn stage_negotiated_publish(
72 self,
73 intent: &SourcePublishIntent,
74 ) -> Result<PublishStageOutcome, TransportAdapterError> {
75 let Some(validated_descriptor) = ({
76 let state = self.room.state.read().await;
77 state.validate_publish(self.user_id, self.connection_id, intent)
78 }) else {
79 return Ok(PublishStageOutcome::Rejected);
80 };
81 self.stage_validated_publish(validated_descriptor).await
82 }
83
84 async fn stage_validated_publish(
85 self,
86 validated_descriptor: ValidatedPublish,
87 ) -> Result<PublishStageOutcome, TransportAdapterError> {
88 let is_duplicate = {
89 let state = self.room.state.read().await;
90 state.staged_publishes.contains(
91 self.user_id,
92 self.connection_id,
93 validated_descriptor.intent.stream_id(),
94 )
95 };
96 if is_duplicate {
97 return Ok(PublishStageOutcome::Duplicate);
98 }
99 let rtp_parameters = RouterRtpParameters::default();
100 let media = match self
101 .media_transport
102 .publish_media(
103 &validated_descriptor.session_key,
104 validated_descriptor.intent.media_kind(),
105 &rtp_parameters,
106 )
107 .await
108 {
109 Ok(media) => media,
110 Err(error) => {
111 warn!(
112 user_id = ?self.user_id,
113 connection_id = ?self.connection_id,
114 stream_id = %validated_descriptor.intent.stream_id(),
115 media_kind = ?validated_descriptor.intent.media_kind(),
116 "failed to stage negotiated publish stream"
117 );
118 return Err(error);
119 }
120 };
121 let reserved_publish = StagedPublish::new(validated_descriptor, media);
122 let duplicate = {
123 let mut state = self.room.state.write().await;
124 if state
127 .validate_publish_commit(&reserved_publish.descriptor, reserved_publish.media)
128 .is_some()
129 {
130 state.staged_publishes.stage(reserved_publish)
131 } else {
132 Some(reserved_publish)
133 }
134 };
135 if let Some(duplicate) = duplicate {
136 duplicate.release_reserved_media(self.media_transport).await;
137 return Ok(PublishStageOutcome::DuplicateAfterReservation);
138 }
139 Ok(PublishStageOutcome::Staged)
140 }
141
142 pub(crate) async fn start_publish(
143 self,
144 intent: &SourcePublishIntent,
145 can_stage: bool,
146 ) -> Result<PublishIntentOutcome, TransportAdapterError> {
147 let has_staged_publish = {
148 let state = self.room.state.read().await;
149 state
150 .staged_publishes
151 .contains(self.user_id, self.connection_id, intent.stream_id())
152 };
153 if has_staged_publish {
154 return Ok(PublishIntentOutcome::Noop);
155 }
156 let source_policy_guard = self.room.lock_source_policy().await;
157 let plan = {
158 let mut state = source_policy_guard.room().state.write().await;
159 state.apply_publish_intent(self.user_id, self.connection_id, intent, can_stage)
160 };
161 match plan {
162 PublishIntentPlan::Activate(commit) => {
163 execute_publication_activity(&source_policy_guard, self.media_transport, commit)
166 .await;
167 drop(source_policy_guard);
168 Ok(PublishIntentOutcome::Activated)
169 }
170 PublishIntentPlan::Noop => {
171 drop(source_policy_guard);
172 Ok(PublishIntentOutcome::Noop)
173 }
174 PublishIntentPlan::Queue => {
175 drop(source_policy_guard);
176 Ok(PublishIntentOutcome::Queue)
177 }
178 PublishIntentPlan::Stage(validated) => {
179 drop(source_policy_guard);
180 if self.stage_validated_publish(validated).await? == PublishStageOutcome::Staged {
181 Ok(PublishIntentOutcome::Staged)
182 } else {
183 Ok(PublishIntentOutcome::Noop)
184 }
185 }
186 }
187 }
188
189 pub async fn rollback_staged_publish(self, stream_id: &UserStreamId) -> bool {
190 let Some(staged) = ({
191 let mut state = self.room.state.write().await;
192 state
193 .staged_publishes
194 .take(self.user_id, self.connection_id, stream_id)
195 }) else {
196 return false;
197 };
198 staged.release_reserved_media(self.media_transport).await;
199 true
200 }
201
202 pub(crate) async fn deactivate_publication(
203 self,
204 intent: &SourceDeactivateIntent,
205 ) -> DeactivateIntentOutcome {
206 if self.rollback_staged_publish(intent.stream_id()).await {
207 return DeactivateIntentOutcome::RolledBack;
208 }
209 let source_policy_guard = self.room.lock_source_policy().await;
212 let commit = {
213 let mut state = source_policy_guard.room().state.write().await;
214 state.apply_publication_activity(
215 self.user_id,
216 self.connection_id,
217 intent.stream_id(),
218 false,
219 intent.presence(),
220 )
221 };
222 let Ok(commit) = commit else {
223 return DeactivateIntentOutcome::Noop;
224 };
225 execute_publication_activity(&source_policy_guard, self.media_transport, commit).await;
226 drop(source_policy_guard);
227 DeactivateIntentOutcome::Deactivated
228 }
229
230 pub(crate) async fn commit_staged_publishes(self, applied_answer: &AppliedSessionAnswer) {
232 let source_policy_guard = self.room.lock_source_policy().await;
236 let staged = {
237 let mut state = source_policy_guard.room().state.write().await;
238 state
239 .staged_publishes
240 .take_for_connection(self.user_id, self.connection_id)
241 };
242 for publish in staged {
243 publish
244 .commit_answer_guarded(&source_policy_guard, self.media_transport, applied_answer)
245 .await;
246 }
247 drop(source_policy_guard);
248 }
249}
250
251async fn execute_publication_activity(
252 guard: &SourcePolicyGuard<'_>,
253 media_transport: &MediaTransport,
254 commit: ProducerActivityCommit,
255) {
256 RoomEffects::from_publication_activity(commit)
257 .execute_with_source_policy_guard(guard, RoomEffectContext::runtime(media_transport))
258 .await;
259}
260
261impl Room {
262 #[cfg(any(test, feature = "testing-transport"))]
263 #[must_use]
264 pub async fn has_staged_publish(
265 &self,
266 user_id: &UserId,
267 connection_id: ConnectionId,
268 stream_id: &UserStreamId,
269 ) -> bool {
270 self.state
271 .read()
272 .await
273 .staged_publishes
274 .contains(user_id, connection_id, stream_id)
275 }
276}
277
278#[cfg(test)]
279#[path = "TESTS/publication.rs"]
280mod tests;