Skip to main content

o_sfu_core/engine/room/effects/
output.rs

1use 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}