Skip to main content

o_sfu_core/engine/room/transition/
publication.rs

1//! publication transitions keep unnegotiated media out of the room graph
2//!
3//! ```text
4//! publish intent
5//!   |
6//!   +-- existing producer --> activity commit --> effects after lock
7//!   |
8//!   +-- offer in flight ----> queued intent ---> answer ---> stage next offer
9//!   |
10//!   +-- new producer -------> StagedPublish ---> answer-proven RTP
11//!                              |                  |
12//!                              |                  v
13//!                              |                room graph commit
14//!                              |                  |
15//!                              v                  v
16//!                         rollback teardown   effects after lock
17//! ```
18//!
19//! only answer-proven RTP enters the room graph
20//! teardown, worker route updates and fanout run after state mutation releases
21//! the room lock
22
23use 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            // `publish_media` ran without the room lock. Revalidate connection
125            // identity and stream uniqueness before room state accepts it.
126            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                // Prevent another source-policy turn from interleaving with this
164                // activity commit and its ordered policy and transport effects.
165                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        // Keep publication activity and its policy effects in one serialized
210        // turn so policy cannot observe the state change without its transport work.
211        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    /// Resolves this connection's staged publishes against an accepted answer.
231    pub(crate) async fn commit_staged_publishes(self, applied_answer: &AppliedSessionAnswer) {
232        // Hold one source-policy turn across the answer batch. Policy must not
233        // observe a committed prefix while later answer-proven publishes remain
234        // outside the room graph.
235        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;