Skip to main content

o_sfu/application/user_session/
room_events.rs

1use o_sfu_protocol::wire::{
2    PeerInfoPayload, PeerLeftPayload, ServerBroadcastPayload, ServerEnvelope, ServerMessage,
3    TrackBinding,
4};
5
6use super::{User, UserError, UserOutput};
7use crate::{
8    application::stream_catalog::stream_type_for_stream_id,
9    runtime::room::{RoomEventMessage, UserOutbound},
10};
11
12impl User {
13    pub(crate) async fn apply_room_outbound(
14        &mut self,
15        outbound: UserOutbound,
16    ) -> Result<UserOutput, UserError> {
17        let mut output = UserOutput::new();
18        match outbound {
19            UserOutbound::Close(_) => {}
20            UserOutbound::Message(message) => match message {
21                RoomEventMessage::Broadcast { sender_id, message } => {
22                    output.push(ServerEnvelope::Message(ServerMessage::Broadcast(
23                        ServerBroadcastPayload {
24                            sender_id,
25                            message: message.to_json(),
26                        },
27                    )));
28                }
29                RoomEventMessage::UserJoined { user_id, info } => {
30                    output.push(ServerEnvelope::Message(ServerMessage::PeerJoined(
31                        PeerInfoPayload { user_id, info },
32                    )));
33                }
34                RoomEventMessage::UserDeparted { user_id } => {
35                    output.push(ServerEnvelope::Message(ServerMessage::PeerLeft(
36                        PeerLeftPayload { user_id },
37                    )));
38                }
39                RoomEventMessage::UserInfoChanged(snapshot) => {
40                    output.extend(snapshot.into_iter().map(|(user_id, info)| {
41                        ServerEnvelope::Message(ServerMessage::PeerInfo(PeerInfoPayload {
42                            user_id,
43                            info,
44                        }))
45                    }));
46                }
47                RoomEventMessage::RecordingStateChanged(state) => {
48                    output.push(ServerEnvelope::Message(ServerMessage::RecordingChange(
49                        state,
50                    )));
51                }
52            },
53            UserOutbound::RemoteTracks(snapshot) => {
54                output.push(ServerEnvelope::Message(ServerMessage::Tracks(
55                    snapshot
56                        .tracks
57                        .into_iter()
58                        .filter_map(|track| {
59                            Some(TrackBinding {
60                                mid: track.consumer_mid,
61                                user_id: track.user_id,
62                                stream_type: stream_type_for_stream_id(&track.stream_id)?,
63                                active: track.producer_active,
64                            })
65                        })
66                        .collect(),
67                )));
68                if snapshot.requires_negotiation {
69                    output.extend(self.renegotiate().await?);
70                }
71            }
72        }
73        Ok(output)
74    }
75}