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