o_sfu_core/engine/room/effects/
batch.rs1use 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#[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 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}