Non blocking SignalClient
This commit is contained in:
+114
-21
@@ -3,36 +3,129 @@
|
||||
//
|
||||
|
||||
#include "signal_client.h"
|
||||
#include "proto/livekit_rtc.pb.h"
|
||||
#include <boost/regex.hpp>
|
||||
#include <iostream>
|
||||
#include <thread>
|
||||
#include <spdlog/spdlog.h>
|
||||
|
||||
namespace livekit {
|
||||
void SignalClient::Connect(const std::string& url, const std::string& token) {
|
||||
boost::regex reg("(ws|wss)://([^:]+):?([^/]*)(.*)");
|
||||
boost::match_results<std::string::const_iterator> regMatches;
|
||||
if (boost::regex_match(url, regMatches, reg))
|
||||
{
|
||||
const std::string protocol = regMatches[1]; // TODO Check if secure or not
|
||||
const std::string domain = regMatches[2];
|
||||
const std::string port = regMatches[3];
|
||||
|
||||
m_WebSocket = std::make_unique<websocket::stream<tcp::socket>>(m_IOContext);
|
||||
SignalClient::SignalClient() : m_Connected(false), m_Writing(false), m_Reading(false) {
|
||||
|
||||
tcp::resolver resolver{m_IOContext};
|
||||
auto const results = resolver.resolve(domain, port);
|
||||
net::connect(m_WebSocket->next_layer(), results.begin(), results.end());
|
||||
}
|
||||
|
||||
m_WebSocket->handshake(domain, "/rtc?access_token=" + token);
|
||||
SignalClient::~SignalClient() {
|
||||
Disconnect();
|
||||
}
|
||||
|
||||
while(true){
|
||||
beast::flat_buffer buffer;
|
||||
m_WebSocket->read(buffer);
|
||||
void SignalClient::Connect(const std::string &url, const std::string &token) {
|
||||
if (m_Connected)
|
||||
throw std::runtime_error{"already connected"};
|
||||
|
||||
m_URL = ParseURL(url);
|
||||
m_Token = token;
|
||||
|
||||
Start(); // We don't need a thread, everything is async ( + easier to maintain )
|
||||
}
|
||||
|
||||
void SignalClient::Update() {
|
||||
beast::error_code ec;
|
||||
m_IOContext.poll(ec);
|
||||
|
||||
if (ec)
|
||||
throw std::runtime_error{"SignalClient::Update - " + ec.message()};
|
||||
|
||||
if (m_Connected) {
|
||||
if (!m_WebSocket.is_open())
|
||||
throw std::runtime_error{"Websocket isn't open"}; // TODO Start reconnect
|
||||
|
||||
if (!m_Reading) {
|
||||
m_WebSocket.async_read(m_ReadBuffer, beast::bind_front_handler(&SignalClient::OnRead, this));
|
||||
m_Reading = true;
|
||||
}
|
||||
|
||||
}else{
|
||||
throw std::runtime_error{"Failed to parse url"};
|
||||
// Write pending messages
|
||||
if (!m_Writing && !m_WriteQueue.empty()) {
|
||||
auto req = m_WriteQueue.front();
|
||||
|
||||
unsigned long len = req.ByteSizeLong();
|
||||
uint8_t data[len];
|
||||
req.SerializeToArray(data, len);
|
||||
|
||||
m_WebSocket.async_write(net::buffer(&data, len),
|
||||
beast::bind_front_handler(&SignalClient::OnWrite, this));
|
||||
|
||||
m_WriteQueue.pop();
|
||||
m_Writing = true;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void SignalClient::Start() {
|
||||
m_Resolver.async_resolve(m_URL.host, m_URL.port, beast::bind_front_handler(&SignalClient::OnResolve, this));
|
||||
}
|
||||
|
||||
void SignalClient::Disconnect() {
|
||||
if (!m_Connected)
|
||||
return;
|
||||
|
||||
m_Connected = false;
|
||||
m_Work.reset();
|
||||
//m_IOContext.stop();
|
||||
m_WebSocket.close(websocket::close_code::normal); // TODO Close should be async
|
||||
}
|
||||
|
||||
void SignalClient::Send(SignalRequest req) {
|
||||
m_WriteQueue.emplace(req);
|
||||
}
|
||||
|
||||
void SignalClient::OnResolve(beast::error_code ec, tcp::resolver::results_type results) {
|
||||
if (ec)
|
||||
throw std::runtime_error{"SignalClient::OnResolve - " + ec.message()};
|
||||
|
||||
auto &layer = beast::get_lowest_layer(m_WebSocket);
|
||||
layer.expires_after(std::chrono::seconds(15));
|
||||
layer.async_connect(results, beast::bind_front_handler(&SignalClient::OnConnect, this));
|
||||
}
|
||||
|
||||
void SignalClient::OnConnect(beast::error_code ec, tcp::resolver::results_type::endpoint_type ep) {
|
||||
if (ec)
|
||||
throw std::runtime_error{"SignalClient::OnConnect - " + ec.message()};
|
||||
|
||||
beast::get_lowest_layer(m_WebSocket).expires_never();
|
||||
m_WebSocket.set_option(websocket::stream_base::timeout::suggested(beast::role_type::client));
|
||||
|
||||
m_WebSocket.async_handshake(m_URL.host, "/rtc?access_token=" + m_Token + "&protocol=7",
|
||||
beast::bind_front_handler(&SignalClient::OnHandshake, this));
|
||||
}
|
||||
|
||||
void SignalClient::OnHandshake(beast::error_code ec) {
|
||||
if (ec)
|
||||
throw std::runtime_error{
|
||||
"SignalClient::OnHandshake - " + ec.message()}; // TODO Callback for handling errors
|
||||
|
||||
m_Connected = true;
|
||||
spdlog::info("Connected to Websocket");
|
||||
}
|
||||
|
||||
void SignalClient::OnRead(beast::error_code ec, std::size_t bytesTransferred) {
|
||||
m_Reading = false;
|
||||
|
||||
if (ec)
|
||||
throw std::runtime_error{"SignalClient::OnRead - " + ec.message()};
|
||||
|
||||
SignalResponse res{};
|
||||
if (res.ParseFromArray(m_ReadBuffer.cdata().data(), bytesTransferred)) {
|
||||
m_ReadQueue.emplace(res);
|
||||
} else {
|
||||
spdlog::error("Failed to decode signal message");
|
||||
}
|
||||
|
||||
m_ReadBuffer.clear();
|
||||
}
|
||||
|
||||
void SignalClient::OnWrite(beast::error_code ec, std::size_t bytesTransferred) {
|
||||
m_Writing = false;
|
||||
|
||||
if (ec)
|
||||
throw std::runtime_error{"SignalClient::OnWrite - " + ec.message()};
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user