o_sfu_core/engine/room/effects/
output.rs1use crate::engine::room::{
2 UserOutbound,
3 outbound::{MessageFanout, OutboundSender, VersionedRemoteTrackSnapshot},
4 state::LifecycleEffects,
5};
6
7#[derive(Debug, Default)]
8pub(super) struct RoomOutputPlan {
9 track_snapshots: Vec<(OutboundSender, VersionedRemoteTrackSnapshot)>,
10 user_info_before_policy: Vec<MessageFanout>,
11 user_info: Vec<MessageFanout>,
12 lifecycle: Vec<LifecycleEffects>,
13}
14
15impl RoomOutputPlan {
16 pub(super) fn push_track_snapshots(
17 &mut self,
18 snapshots: Vec<(OutboundSender, VersionedRemoteTrackSnapshot)>,
19 ) {
20 self.track_snapshots.extend(snapshots);
21 }
22
23 pub(super) fn push_user_info(&mut self, fanout: MessageFanout) {
24 self.user_info.push(fanout);
25 }
26
27 pub(super) fn push_user_info_before_policy(&mut self, fanout: MessageFanout) {
28 self.user_info_before_policy.push(fanout);
29 }
30
31 pub(super) fn push_lifecycle(&mut self, effects: LifecycleEffects) {
32 self.lifecycle.push(effects);
33 }
34
35 pub(super) fn emit_before_policy(&mut self) {
36 for (recipient, snapshot) in self.track_snapshots.drain(..) {
37 let _ = recipient.send_remote_tracks(snapshot);
38 }
39 self.emit_user_info_before_policy();
40 }
41
42 pub(super) fn emit_user_info_before_policy(&mut self) {
43 for fanout in self.user_info_before_policy.drain(..) {
44 fanout.emit();
45 }
46 }
47
48 pub(super) fn emit_after_policy(self) {
49 for fanout in self.user_info {
50 fanout.emit();
51 }
52 for effects in self.lifecycle {
53 for close_request in effects.close_requests {
54 let _ = close_request
55 .sender
56 .send(UserOutbound::Close(close_request.reason));
57 }
58 for (recipient, snapshot) in effects.track_snapshots {
59 let _ = recipient.send_remote_tracks(snapshot);
60 }
61 for fanout in effects.fanouts {
62 fanout.emit();
63 }
64 }
65 }
66}