Commit 8cf99f33 authored by Sebastian Dröge's avatar Sebastian Dröge
Browse files

Use weak references for "lower" actors storing "upper" actors addresses

Otherwise we create reference cycles.
parent b9df8bcc
Loading
Loading
Loading
Loading
+0 −63
Original line number Diff line number Diff line
@@ -1010,68 +1010,6 @@ dependencies = [
 "thiserror",
]

[[package]]
name = "gstreamer-app"
version = "0.16.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "330e3e40c0e27680e52e38c599353dddfb14c373a148dd9154243f7cf92ba3c2"
dependencies = [
 "bitflags",
 "futures-core",
 "futures-sink",
 "glib",
 "glib-sys",
 "gobject-sys",
 "gstreamer",
 "gstreamer-app-sys",
 "gstreamer-base",
 "gstreamer-sys",
 "libc",
 "once_cell",
]

[[package]]
name = "gstreamer-app-sys"
version = "0.9.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "813f64275c9e7b33b828b9efcf9dfa64b95996766d4de996e84363ac65b87e3d"
dependencies = [
 "glib-sys",
 "gstreamer-base-sys",
 "gstreamer-sys",
 "libc",
 "system-deps",
]

[[package]]
name = "gstreamer-base"
version = "0.16.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "082f83c11a9d568e94e3a29d9fb22bf86dda44becfd4ba1d18446766b678f07c"
dependencies = [
 "bitflags",
 "glib",
 "glib-sys",
 "gobject-sys",
 "gstreamer",
 "gstreamer-base-sys",
 "gstreamer-sys",
 "libc",
]

[[package]]
name = "gstreamer-base-sys"
version = "0.9.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a4b7b6dc2d6e160a1ae28612f602bd500b3fa474ce90bf6bb2f08072682beef5"
dependencies = [
 "glib-sys",
 "gobject-sys",
 "gstreamer-sys",
 "libc",
 "system-deps",
]

[[package]]
name = "gstreamer-sdp"
version = "0.16.3"
@@ -2550,7 +2488,6 @@ dependencies = [
 "futures",
 "glib",
 "gstreamer",
 "gstreamer-app",
 "gstreamer-webrtc",
 "log",
 "openssl",
+23 −11
Original line number Diff line number Diff line
@@ -6,12 +6,12 @@ use crate::config::Config;
use crate::rooms::{Room, Rooms};
use crate::subscriber::Subscriber;

use std::sync::{Arc, Mutex};

use anyhow::Error;
use anyhow::{format_err, Error};

use actix::{Actor, Addr, AsyncContext, Handler, Message, StreamHandler};
use std::sync::{Arc, Mutex};

use actix::{Actor, Addr, Handler, Message, StreamHandler, WeakAddr};
use actix_web::dev::ConnectionInfo;
use actix_web_actors::ws;

use log::{debug, trace};
@@ -20,29 +20,41 @@ use log::{debug, trace};
#[derive(Debug)]
pub struct Publisher {
    cfg: Arc<Config>,
    rooms: Addr<Rooms>,
    rooms: WeakAddr<Rooms>,
    remote_addr: String,
    room: Mutex<Option<Addr<Room>>>,
}

impl Publisher {
    /// Create a new `Publisher` actor.
    pub fn new(cfg: Arc<Config>, rooms: Addr<Rooms>) -> Self {
        Publisher {
    pub fn new(
        cfg: Arc<Config>,
        rooms: Addr<Rooms>,
        connection_info: &ConnectionInfo,
    ) -> Result<Self, Error> {
        debug!("Creating new publisher {:?}", connection_info);

        let remote_addr = connection_info
            .realip_remote_addr()
            .ok_or_else(|| format_err!("WebSocket connection without remote address"))?;

        Ok(Publisher {
            cfg,
            rooms,
            rooms: rooms.downgrade(),
            remote_addr: String::from(remote_addr),
            room: Mutex::new(None),
        }
        })
    }
}

impl Actor for Publisher {
    type Context = ws::WebsocketContext<Self>;

    fn stopped(&mut self, ctx: &mut Self::Context) {
    fn stopped(&mut self, _ctx: &mut Self::Context) {
        // Drop reference to the joined room, if any
        self.room.lock().unwrap().take();

        trace!("Publisher {:?} stopped", ctx.address());
        trace!("Publisher {} stopped", self.remote_addr);
    }
}

+9 −4
Original line number Diff line number Diff line
@@ -6,7 +6,9 @@ use std::collections::{HashMap, HashSet};
use std::fmt;
use std::sync::Mutex;

use actix::{Actor, ActorContext, Addr, AsyncContext, Context, Handler, Message, MessageResult};
use actix::{
    Actor, ActorContext, Addr, AsyncContext, Context, Handler, Message, MessageResult, WeakAddr,
};
use anyhow::{bail, Error};
use log::{debug, error, info, trace};

@@ -188,7 +190,7 @@ pub struct Room {
    name: String,
    description: Option<String>,

    rooms: Addr<Rooms>,
    rooms: WeakAddr<Rooms>,

    publisher: Addr<Publisher>,
    subscribers: Mutex<HashSet<Addr<Subscriber>>>,
@@ -203,7 +205,7 @@ impl Room {
        publisher: Addr<Publisher>,
    ) -> Self {
        Room {
            rooms,
            rooms: rooms.downgrade(),
            id: RoomId(uuid::Uuid::new_v4()),
            name,
            description,
@@ -241,7 +243,10 @@ impl Handler<DeleteRoomMessage> for Room {
                subscriber.do_send(subscriber::RoomDeletedMessage);
            }
        }
        self.rooms.do_send(RoomDeletedMessage { room_id: self.id });

        if let Some(rooms) = self.rooms.upgrade() {
            rooms.do_send(RoomDeletedMessage { room_id: self.id });
        }

        ctx.stop();

+28 −10
Original line number Diff line number Diff line
@@ -12,6 +12,8 @@ use actix_files::NamedFile;
use actix_web::{web, App, HttpRequest, HttpResponse, HttpServer, Responder};
use actix_web_actors::ws;

use log::error;

/// Serve `index.html` for path `/`.
async fn index(cfg: web::Data<Config>) -> Result<NamedFile, actix_web::Error> {
    let full_path = cfg.static_files.join("index.html");
@@ -30,16 +32,32 @@ async fn ws(
    stream: web::Payload,
) -> impl Responder {
    match path.as_str() {
        "publish" => ws::start(
            Publisher::new(cfg.into_inner(), rooms.as_ref().clone()),
            &req,
            stream,
        ),
        "subscribe" => ws::start(
            Subscriber::new(cfg.into_inner(), rooms.as_ref().clone()),
            &req,
            stream,
        ),
        "publish" => {
            let publisher = Publisher::new(
                cfg.into_inner(),
                rooms.as_ref().clone(),
                &req.connection_info(),
            )
            .map_err(|err| {
                error!("Failed to create publisher: {}", err);
                HttpResponse::InternalServerError()
            })?;

            ws::start(publisher, &req, stream)
        }
        "subscribe" => {
            let subscriber = Subscriber::new(
                cfg.into_inner(),
                rooms.as_ref().clone(),
                &req.connection_info(),
            )
            .map_err(|err| {
                error!("Failed to create subscriber: {}", err);
                HttpResponse::InternalServerError()
            })?;

            ws::start(subscriber, &req, stream)
        }
        _ => Ok(HttpResponse::NotFound().finish()),
    }
}
+24 −10
Original line number Diff line number Diff line
@@ -5,10 +5,12 @@
use crate::config::Config;
use crate::rooms::{Room, Rooms};

use std::sync::{Arc, Mutex};
use anyhow::{format_err, Error};

use actix::{Actor, Addr, AsyncContext, Handler, Message, StreamHandler};
use std::sync::{Arc, Mutex};

use actix::{Actor, Addr, Handler, Message, StreamHandler, WeakAddr};
use actix_web::dev::ConnectionInfo;
use actix_web_actors::ws;

use log::{debug, trace};
@@ -17,29 +19,41 @@ use log::{debug, trace};
#[derive(Debug)]
pub struct Subscriber {
    cfg: Arc<Config>,
    rooms: Addr<Rooms>,
    room: Mutex<Option<Addr<Room>>>,
    rooms: WeakAddr<Rooms>,
    remote_addr: String,
    room: Mutex<Option<WeakAddr<Room>>>,
}

impl Subscriber {
    /// Create a new `Subscriber` actor.
    pub fn new(cfg: Arc<Config>, rooms: Addr<Rooms>) -> Self {
        Subscriber {
    pub fn new(
        cfg: Arc<Config>,
        rooms: Addr<Rooms>,
        connection_info: &ConnectionInfo,
    ) -> Result<Self, Error> {
        debug!("Creating new subscriber {:?}", connection_info);

        let remote_addr = connection_info
            .realip_remote_addr()
            .ok_or_else(|| format_err!("WebSocket connection without remote address"))?;

        Ok(Subscriber {
            cfg,
            rooms,
            rooms: rooms.downgrade(),
            remote_addr: String::from(remote_addr),
            room: Mutex::new(None),
        }
        })
    }
}

impl Actor for Subscriber {
    type Context = ws::WebsocketContext<Self>;

    fn stopped(&mut self, ctx: &mut Self::Context) {
    fn stopped(&mut self, _ctx: &mut Self::Context) {
        // Drop reference to the joined room, if any
        self.room.lock().unwrap().take();

        trace!("Subscriber {:?} stopped", ctx.address());
        trace!("Subscriber {} stopped", self.remote_addr);
    }
}