Skip to main content

o_sfu_core/engine/room/effects/
batch.rs

1use super::{output::RoomOutputPlan, transport::RoomTransportPlan};
2use crate::engine::{
3    media_transport::MediaTransport,
4    room::{
5        Room,
6        media_graph::{
7            ConsumerSetupOrigin, ProducerActivityCommit, PublishCommit, ReceiverRouteCommit,
8            ReceiverRouteWork,
9        },
10        source_policy::SourcePolicyTurn,
11        state::{
12            ConnectionCloseCommit, DisconnectCommit, JoinCommit, PresenceCommit, UserJoinedFanout,
13        },
14    },
15};
16
17#[derive(Debug, Clone, Copy)]
18pub struct RoomEffectContext<'a> {
19    media_transport: Option<&'a MediaTransport>,
20    route_effects: bool,
21    joined_fanout: UserJoinedFanout,
22}
23
24impl<'a> RoomEffectContext<'a> {
25    pub const fn runtime(media_transport: &'a MediaTransport) -> Self {
26        Self {
27            media_transport: Some(media_transport),
28            route_effects: true,
29            joined_fanout: UserJoinedFanout::Emit,
30        }
31    }
32
33    #[cfg(any(test, feature = "testing-transport"))]
34    pub const fn state_only(media_transport: Option<&'a MediaTransport>) -> Self {
35        Self {
36            media_transport,
37            route_effects: false,
38            joined_fanout: UserJoinedFanout::Suppress,
39        }
40    }
41
42    pub(in crate::engine::room) const fn user_joined_fanout(self) -> UserJoinedFanout {
43        self.joined_fanout
44    }
45
46    fn media_transport(self) -> Option<&'a MediaTransport> {
47        self.media_transport
48    }
49
50    fn route_transport(self) -> Option<&'a MediaTransport> {
51        self.route_effects.then_some(self.media_transport).flatten()
52    }
53}
54
55/// batches post-lock transport, signaling and policy side-effects for room state transitions
56///
57/// ```text
58/// room state mutation (holds RoomState write lock)
59///   - mutate in-memory graph
60///   - return commit data
61///             |
62///             v  drop RoomState write lock
63/// RoomEffects::from_*(commit) -> execute_inner
64///             |
65///             v  step 1: transport execution
66///   +-----------------------------------------------------------+
67///   | - create/remove local consumer routes on workers          |
68///   | - register/remove cross-worker relay route targets        |
69///   | - dispatch session teardowns                              |
70///   +-----------------------------------------------------------+
71///             |
72///             v  step 2: RoomOutputPlan pre-policy fanout
73///   +-----------------------------------------------------------+
74///   | - send track snapshots and presence user-info fanout      |
75///   +-----------------------------------------------------------+
76///             |
77///             v  step 3: source policy turn
78///   +-----------------------------------------------------------+
79///   | - re-evaluate audio admission and video bandwidth solver  |
80///   | - commit packet gate and BWE target updates               |
81///   +-----------------------------------------------------------+
82///             |
83///             v  step 4: RoomOutputPlan post-policy fanout
84///   +-----------------------------------------------------------+
85///   | - send user-info and lifecycle close/track/fanout output  |
86///   +-----------------------------------------------------------+
87/// ```
88///
89/// The diagram shows the normal batch order. Only an active `ProducerActivityCommit` from
90/// `from_publication_activity` runs its pre-policy user-info output and source policy before transport.
91/// Callers of `from_publish` and `from_publication_activity` use
92/// `execute_with_source_policy_guard` while holding `source_policy_turn`.
93#[derive(Debug, Default)]
94#[must_use = "room effect batches must be executed after the state transition commits"]
95pub struct RoomEffects {
96    policy_before_transport: bool,
97    transport: RoomTransportPlan,
98    output: RoomOutputPlan,
99    source_policy: SourcePolicyTurn,
100}
101
102impl RoomEffects {
103    pub(in crate::engine::room) fn from_join(commit: JoinCommit) -> Self {
104        let JoinCommit {
105            effects,
106            transport_plan,
107            ..
108        } = commit;
109        let mut batch = Self::default();
110        batch.transport.extend(transport_plan);
111        batch.source_policy.request();
112        batch.output.push_lifecycle(effects);
113        batch
114    }
115
116    pub(in crate::engine::room) fn from_connection_close(commit: ConnectionCloseCommit) -> Self {
117        let mut batch = Self::default();
118        match commit {
119            ConnectionCloseCommit::Current {
120                session_teardown,
121                effects,
122                transport_plan,
123                ..
124            } => {
125                batch.transport.extend(transport_plan);
126                batch.output.push_lifecycle(effects);
127                batch.source_policy.request();
128                batch.transport.extend_teardown(session_teardown);
129            }
130            ConnectionCloseCommit::StalePlacement { session_teardown } => {
131                batch.transport.extend_teardown([session_teardown]);
132            }
133        }
134        batch
135    }
136
137    pub(in crate::engine::room) fn from_disconnect(commit: DisconnectCommit) -> Self {
138        let mut batch = Self::default();
139        batch.transport.extend(commit.transport_plan);
140        batch.source_policy.request();
141        batch.output.push_lifecycle(commit.effects);
142        batch.transport.extend_teardown(commit.session_teardowns);
143        batch
144    }
145
146    pub(in crate::engine::room) fn from_presence(commit: PresenceCommit) -> Self {
147        let mut batch = Self::default();
148        batch.output.push_user_info(commit.fanout);
149        batch.source_policy.request();
150        batch
151    }
152
153    pub(in crate::engine::room) fn from_publish(commit: PublishCommit) -> Self {
154        let mut batch = Self::default();
155        batch
156            .transport
157            .push_receiver_work(commit.receiver_route_work, ConsumerSetupOrigin::Publish);
158        batch.push_presence_before_policy(commit.presence);
159        batch.source_policy.request();
160        batch
161    }
162
163    pub(in crate::engine::room) fn from_publication_activity(
164        commit: ProducerActivityCommit,
165    ) -> Self {
166        let ProducerActivityCommit {
167            source,
168            stream_id,
169            update,
170            remote_activity_effects,
171            track_snapshots,
172            presence,
173        } = commit;
174        let mut batch = Self {
175            policy_before_transport: update.activity().is_active(),
176            ..Self::default()
177        };
178        batch
179            .transport
180            .extend_remote_source_activity(remote_activity_effects);
181        batch.transport.push_producer(source, stream_id, update);
182        batch.output.push_track_snapshots(track_snapshots);
183        batch.push_presence_before_policy(presence);
184        batch.source_policy.request();
185        batch
186    }
187
188    pub(in crate::engine::room) fn from_receiver_intent(commit: ReceiverRouteCommit) -> Self {
189        let mut batch = Self::from_receiver_route(commit.work, ConsumerSetupOrigin::Subscribe);
190        batch.source_policy.request();
191        batch
192    }
193
194    pub(in crate::engine::room) fn from_consumer_readiness(commit: ReceiverRouteCommit) -> Self {
195        let ReceiverRouteCommit {
196            work,
197            track_snapshots,
198        } = commit;
199        let mut batch = Self::from_receiver_route(work, ConsumerSetupOrigin::Readiness);
200        batch.output.push_track_snapshots(track_snapshots);
201        batch.source_policy.request();
202        batch
203    }
204
205    fn push_presence_before_policy(&mut self, presence: Option<PresenceCommit>) {
206        if let Some(presence) = presence {
207            self.output.push_user_info_before_policy(presence.fanout);
208            self.source_policy.request();
209        }
210    }
211
212    fn from_receiver_route(work: ReceiverRouteWork, origin: ConsumerSetupOrigin) -> Self {
213        let mut batch = Self::default();
214        batch.transport.push_receiver_work(work, origin);
215        batch
216    }
217
218    /// preserves the room-wide side-effect order across transport and policy work
219    pub async fn execute(self, room: &Room, context: RoomEffectContext<'_>) {
220        if self.policy_before_transport {
221            let _guard = room.source_policy_turn.lock().await;
222            self.execute_inner(room, context, true).await;
223        } else {
224            self.execute_inner(room, context, false).await;
225        }
226    }
227
228    pub(in crate::engine::room) async fn execute_with_source_policy_guard(
229        self,
230        room: &Room,
231        context: RoomEffectContext<'_>,
232    ) {
233        self.execute_inner(room, context, true).await;
234    }
235
236    async fn execute_inner(
237        self,
238        room: &Room,
239        context: RoomEffectContext<'_>,
240        source_policy_guarded: bool,
241    ) {
242        let mut output = self.output;
243        let mut source_policy = self.source_policy;
244        if self.policy_before_transport {
245            output.emit_user_info_before_policy();
246            source_policy
247                .execute_guarded(room, context.media_transport(), None)
248                .await;
249            source_policy = SourcePolicyTurn::default();
250        }
251        self.transport
252            .execute(room, context.route_transport())
253            .await;
254        output.emit_before_policy();
255        if source_policy_guarded {
256            source_policy
257                .execute_guarded(room, context.media_transport(), None)
258                .await;
259        } else {
260            source_policy
261                .execute(room, context.media_transport(), None)
262                .await;
263        }
264        output.emit_after_policy();
265    }
266}