Commit 7050074a authored by Sebastian Dröge's avatar Sebastian Dröge
Browse files

Clean up messaging protocol a bit

Closing the websocket connection is sufficent for deleting/leaving
rooms.
parent 1c646e73
Loading
Loading
Loading
Loading
+0 −2
Original line number Diff line number Diff line
@@ -12,7 +12,6 @@ pub enum PublisherMessage {
        name: String,
        description: Option<String>,
    },
    DeleteRoom,
    Ice {
        candidate: String,
        #[serde(rename = "sdpMLineIndex")]
@@ -35,7 +34,6 @@ pub enum ServerMessage {
    RoomCreated {
        id: uuid::Uuid,
    },
    RoomDeleted,
    Ice {
        candidate: String,
        #[serde(rename = "sdpMLineIndex")]
+0 −2
Original line number Diff line number Diff line
@@ -11,7 +11,6 @@ pub enum SubscriberMessage {
    JoinRoom {
        id: uuid::Uuid,
    },
    LeaveRoom,
    Ice {
        candidate: String,
        #[serde(rename = "sdpMLineIndex")]
@@ -31,7 +30,6 @@ pub enum ServerMessage {
    Error {
        message: String,
    },
    LeftRoom,
    Ice {
        candidate: String,
        #[serde(rename = "sdpMLineIndex")]
+0 −7
Original line number Diff line number Diff line
@@ -272,13 +272,6 @@ impl Publisher {
                            error!("Got server error: {}", message);
                            bail!("Server error: {}", message);
                        }
                        ServerMessage::RoomDeleted => {
                            info!("Room got deleted");
                            websocket_sender.close_channel();
                            event_receiver.close();
                            websocket_send_task.await.context("Closing websocket")??;
                            break;
                        }
                        ServerMessage::RoomCreated { .. } => {
                            error!("Unexpected RoomCreated server message");
                            continue;
+0 −32
Original line number Diff line number Diff line
@@ -354,38 +354,6 @@ impl Publisher {
                // Asynchronously handle the SDP by starting the pipeline and creating the answer SDP.
                ctx.spawn(self.handle_sdp_future(sdp_message));
            }
            Ok(PublisherMessage::DeleteRoom) => {
                debug!("Publisher {} deleting room", self.remote_addr);

                let mut room_state = self.room.lock().unwrap();
                match &*room_state {
                    RoomState::Joined(room, _) => {
                        if let Some(room) = room.upgrade() {
                            room.do_send(crate::rooms::DeleteRoomMessage {
                                publisher: ctx.address(),
                            });
                        }
                        ctx.text(
                            serde_json::to_string(&ServerMessage::RoomDeleted)
                                .expect("Failed to serialize room deleted message"),
                        );
                        *room_state = RoomState::None;
                    }
                    RoomState::Joining => {
                        debug!("Aborting joining of room");
                        *room_state = RoomState::None;
                    }
                    RoomState::None => (),
                }
                self.pipeline.call_async(|pipeline| {
                    debug!("Stopping pipeline {}", pipeline.get_name());
                    let _ = pipeline.set_state(gst::State::Null);
                });

                let mut subscribers = self.subscribers.lock().unwrap();
                subscribers.subscribers.clear();
                subscribers.latency = gst::CLOCK_TIME_NONE;
            }
            Err(err) => {
                error!(
                    "Publisher {} has websocket error: {}",
+1 −49
Original line number Diff line number Diff line
@@ -304,35 +304,6 @@ impl Subscriber {
                    ctx.spawn(self.join_room_future(ctx.address(), rooms, RoomId(id)));
                }
            }
            Ok(SubscriberMessage::LeaveRoom) => {
                debug!("Subscriber {} leaving room", self.remote_addr);

                let mut room_state = self.room.lock().unwrap();
                match &*room_state {
                    RoomState::Joined(room, _) => {
                        if let Some(room) = room.upgrade() {
                            room.do_send(crate::rooms::LeaveRoomMessage {
                                subscriber: ctx.address(),
                            });
                        }
                        ctx.text(
                            serde_json::to_string(&ServerMessage::LeftRoom)
                                .expect("Failed to serialize left room message"),
                        );
                        *room_state = RoomState::None;
                    }
                    RoomState::Joining => {
                        debug!("Aborting joining of room");
                        *room_state = RoomState::None;
                    }
                    RoomState::None => (),
                }

                self.pipeline.call_async(|pipeline| {
                    debug!("Stopping pipeline");
                    let _ = pipeline.set_state(gst::State::Null);
                });
            }
            Ok(SubscriberMessage::Ice {
                candidate,
                sdp_mline_index,
@@ -522,26 +493,7 @@ impl Handler<RoomDeletedMessage> for Subscriber {
        ctx: &mut ws::WebsocketContext<Self>,
    ) -> Self::Result {
        debug!("Room deleted for subscriber {}", self.remote_addr);

        let mut room_state = self.room.lock().unwrap();
        match &*room_state {
            RoomState::Joined(_, _) => {
                ctx.text(
                    serde_json::to_string(&ServerMessage::LeftRoom)
                        .expect("Failed to serialize left room message"),
                );
                *room_state = RoomState::None;
            }
            RoomState::Joining => {
                *room_state = RoomState::None;
            }
            RoomState::None => (),
        }

        self.pipeline.call_async(|pipeline| {
            debug!("Stopping pipeline");
            let _ = pipeline.set_state(gst::State::Null);
        });
        self.shutdown(ctx);
    }
}