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