Initial Room

Switching computer
This commit is contained in:
Théo Monnom
2022-10-29 00:11:10 +02:00
parent 2ecae706a9
commit 3630e93054
47 changed files with 1783 additions and 980 deletions
+87 -11
View File
@@ -1,18 +1,23 @@
use core::num::flt2dec::Sign;
use std::fmt::Debug;
use std::time::Duration;
use livekit_webrtc::peer_connection_factory::{
ContinualGatheringPolicy, ICEServer, IceTransportsType, RTCConfiguration,
};
use thiserror::Error;
use tokio::sync::mpsc;
use tokio_tungstenite::tungstenite::Error as WsError;
use crate::event::{Emitter, Events};
use crate::proto::{signal_request, signal_response};
use crate::proto::{signal_request, signal_response, JoinResponse};
use crate::signal_client::signal_stream::SignalStream;
mod signal_stream;
type SignalEmitter = Emitter<SignalEvent>;
type SignalEvents = Events<SignalEvent>;
type SignalResult<T> = Result<T, SignalError>;
pub(crate) type SignalEmitter = mpsc::Sender<SignalEvent>;
pub(crate) type SignalEvents = mpsc::Receiver<SignalEvent>;
pub(crate) type SignalResult<T> = Result<T, SignalError>;
pub const JOIN_RESPONSE_TIMEOUT: Duration = Duration::from_secs(5);
#[derive(Error, Debug)]
pub enum SignalError {
@@ -22,10 +27,12 @@ pub enum SignalError {
UrlParse(#[from] url::ParseError),
#[error("failed to decode messages from server")]
ProtoParse(#[from] prost::DecodeError),
#[error("{0}")]
Timeout(String),
}
/// Events used by the RTCEngine who will handle the reconnection logic
#[derive(Clone, Debug)]
#[derive(Debug)]
pub(crate) enum SignalEvent {
Open,
Signal(signal_response::Message),
@@ -40,6 +47,17 @@ pub(crate) struct SignalOptions {
adaptive_stream: bool,
}
impl Default for SignalOptions {
fn default() -> Self {
Self {
reconnect: false,
auto_subscribe: true,
sid: "".to_string(),
adaptive_stream: true,
}
}
}
#[derive(Debug)]
pub struct SignalClient {
stream: SignalStream,
@@ -47,15 +65,16 @@ pub struct SignalClient {
}
impl SignalClient {
pub async fn connect(
pub(crate) async fn connect(
url: &str,
token: &str,
options: SignalOptions,
) -> SignalResult<(Self, SignalEvents)> {
// TODO(theomonnom) Retry initial connection
let (emitter, receiver) = SignalEmitter::new();
let events = SignalEvents::new(receiver);
let (emitter, events) = mpsc::channel(8);
let stream = SignalStream::connect(url, token, options, emitter.clone()).await?;
// TODO(theomonnom) Retry initial connection
Ok((Self { stream, emitter }, events))
}
@@ -69,3 +88,60 @@ impl SignalClient {
// TODO(theomonnom) Close & recreate SignalStream, also send the queue if needed
}
}
impl From<JoinResponse> for RTCConfiguration {
fn from(join_response: JoinResponse) -> Self {
Self {
ice_servers: {
let mut servers = vec![];
for ice_server in join_response.ice_servers.clone() {
servers.push(ICEServer {
urls: ice_server.urls,
username: ice_server.username,
password: ice_server.credential,
})
}
servers
},
continual_gathering_policy: ContinualGatheringPolicy::GatherContinually,
ice_transport_type: IceTransportsType::All,
}
}
}
pub mod utils {
use crate::proto::{signal_response, JoinResponse};
use crate::signal_client::{SignalError, SignalEvent, SignalResult, JOIN_RESPONSE_TIMEOUT};
use tokio::sync::mpsc;
use tokio::time::timeout;
use tokio_tungstenite::tungstenite::Error as WsError;
use tracing::{event, Level};
pub(crate) async fn next_join_response(
receiver: &mut mpsc::Receiver<SignalEvent>,
) -> SignalResult<JoinResponse> {
let join = async {
while let Some(event) = receiver.recv().await {
match event {
SignalEvent::Signal(signal_response::Message::Join(join)) => return Ok(join),
SignalEvent::Close => break,
SignalEvent::Open => continue,
_ => {
event!(
Level::WARN,
"received unexpected message while waiting for JoinResponse: {:?}",
event
);
continue;
}
}
}
Err(WsError::ConnectionClosed)?
};
timeout(JOIN_RESPONSE_TIMEOUT, join)
.await
.map_err(|_| SignalError::Timeout("failed to receive JoinResponse".to_string()))?
}
}
@@ -1,13 +1,13 @@
use futures_util::{SinkExt, StreamExt};
use futures_util::stream::{SplitSink, SplitStream};
use futures_util::{SinkExt, StreamExt};
use prost::Message as ProstMessage;
use tokio::net::TcpStream;
use tokio::sync::{mpsc, oneshot};
use tokio::task::JoinHandle;
use tokio_tungstenite::{connect_async, MaybeTlsStream, WebSocketStream};
use tokio_tungstenite::tungstenite::Message;
use tokio_tungstenite::tungstenite::protocol::CloseFrame;
use tokio_tungstenite::tungstenite::protocol::frame::coding::CloseCode;
use tokio_tungstenite::tungstenite::protocol::CloseFrame;
use tokio_tungstenite::tungstenite::Message;
use tokio_tungstenite::{connect_async, MaybeTlsStream, WebSocketStream};
use tracing::{event, Level};
use crate::proto::{signal_request, SignalRequest, SignalResponse};
@@ -72,7 +72,7 @@ impl SignalStream {
event!(Level::DEBUG, "connecting to websocket: {}", lk_url);
let (ws_stream, _) = connect_async(lk_url).await?;
event!(Level::DEBUG, "connected to websocket");
emitter.event(SignalEvent::Open);
let _ = emitter.send(SignalEvent::Open).await;
let (ws_writer, ws_reader) = ws_stream.split();
let (internal_tx, internal_rx) = mpsc::channel::<InternalMessage>(8);
@@ -119,7 +119,7 @@ impl SignalStream {
/// This task is used to send messages to the websocket
/// It is also responsible for closing the connection
pub async fn handle_write(
async fn handle_write(
mut internal_rx: mpsc::Receiver<InternalMessage>,
mut ws_writer: SplitSink<WebSocket, Message>,
emitter: SignalEmitter,
@@ -136,7 +136,7 @@ impl SignalStream {
SignalRequest {
message: Some(signal),
}
.encode_to_vec(),
.encode_to_vec(),
);
if let Err(err) = ws_writer.send(data).await {
@@ -163,14 +163,14 @@ impl SignalStream {
}
let _ = ws_writer.close().await;
emitter.event(SignalEvent::Close);
let _ = emitter.send(SignalEvent::Close).await;
}
/// This task is used to read incoming messages from the websocket
/// and dispatch them through the EventEmitter.
///
/// It can also send messages to [handle_write] task ( Used e.g. answer to pings )
pub async fn handle_read(
async fn handle_read(
internal_tx: mpsc::Sender<InternalMessage>,
mut ws_reader: SplitStream<WebSocket>,
emitter: SignalEmitter,
@@ -181,8 +181,9 @@ impl SignalStream {
let res = SignalResponse::decode(data.as_slice())
.expect("failed to decode SignalResponse");
event!(Level::TRACE, "received SignalResponse: {:?}", res);
emitter.event(SignalEvent::Signal(res.message.unwrap()));
let msg = res.message.unwrap();
event!(Level::TRACE, "received SignalResponse: {:?}", msg);
let _ = emitter.send(SignalEvent::Signal(msg)).await;
}
Ok(Message::Ping(data)) => {
let _ = internal_tx