Merge pull request #257 from morphis/f/improve-adb-stability

Improve stability of the AdbMessageProcessor
This commit is contained in:
Simon Fels 2017-05-14 21:07:08 +02:00 • committed by GitHub
commit 52f4a7e8e1
2 changed files with 34 additions and 20 deletions

View file

@ -16,10 +16,10 @@
*/ */
#include "anbox/qemu/adb_message_processor.h" #include "anbox/qemu/adb_message_processor.h"
#include "anbox/logger.h"
#include "anbox/network/delegate_connection_creator.h" #include "anbox/network/delegate_connection_creator.h"
#include "anbox/network/delegate_message_processor.h" #include "anbox/network/delegate_message_processor.h"
#include "anbox/network/tcp_socket_messenger.h" #include "anbox/network/tcp_socket_messenger.h"
#include "anbox/utils.h"
#include <fstream> #include <fstream>
#include <functional> #include <functional>
@ -27,18 +27,23 @@
namespace { namespace {
const unsigned short default_adb_client_port{5037}; const unsigned short default_adb_client_port{5037};
const unsigned short default_host_listen_port{6664}; const unsigned short default_host_listen_port{6664};
constexpr const char *loopback_address{"127.0.0.1"};
const std::string accept_command{"accept"}; const std::string accept_command{"accept"};
const std::string ok_command{"ok"}; const std::string ok_command{"ok"};
const std::string ko_command{"ko"}; const std::string ko_command{"ko"};
const std::string start_command{"start"}; const std::string start_command{"start"};
// This timeount should be too high to not cause a too long wait time for the
// user until we connect to the adb host instance after it appeared and not
// too short to not put unnecessary burden on the CPU.
const boost::posix_time::seconds default_adb_wait_time{1}; const boost::posix_time::seconds default_adb_wait_time{1};
static std::mutex active_instance;
} }
using namespace std::placeholders; using namespace std::placeholders;
namespace anbox { namespace anbox {
namespace qemu { namespace qemu {
std::mutex AdbMessageProcessor::active_instance_{};
AdbMessageProcessor::AdbMessageProcessor( AdbMessageProcessor::AdbMessageProcessor(
const std::shared_ptr<Runtime> &rt, const std::shared_ptr<Runtime> &rt,
const std::shared_ptr<network::SocketMessenger> &messenger) const std::shared_ptr<network::SocketMessenger> &messenger)
@ -46,12 +51,14 @@ AdbMessageProcessor::AdbMessageProcessor(
state_(waiting_for_guest_accept_command), state_(waiting_for_guest_accept_command),
expected_command_(accept_command), expected_command_(accept_command),
messenger_(messenger), messenger_(messenger),
host_notify_timer_(rt->service()) {} host_notify_timer_(rt->service()),
lock_(active_instance_, std::defer_lock) {}
AdbMessageProcessor::~AdbMessageProcessor() { AdbMessageProcessor::~AdbMessageProcessor() {
state_ = closed_by_host; state_ = closed_by_host;
host_notify_timer_.cancel();
host_connector_.reset(); host_connector_.reset();
active_instance.unlock();
} }
void AdbMessageProcessor::advance_state() { void AdbMessageProcessor::advance_state() {
@ -61,9 +68,10 @@ void AdbMessageProcessor::advance_state() {
// running we don't have to do anything here until that one is done. // running we don't have to do anything here until that one is done.
// The container directly starts a second connection once the first // The container directly starts a second connection once the first
// one is established but will not use it until the active one is closed. // one is established but will not use it until the active one is closed.
active_instance.lock(); lock_.lock();
if (state_ == closed_by_host) { if (state_ == closed_by_host) {
host_notify_timer_.cancel();
host_connector_.reset(); host_connector_.reset();
return; return;
} }
@ -100,9 +108,12 @@ void AdbMessageProcessor::advance_state() {
} }
void AdbMessageProcessor::wait_for_host_connection() { void AdbMessageProcessor::wait_for_host_connection() {
if (state_ == closed_by_host || state_ == closed_by_container)
return;
if (!host_connector_) { if (!host_connector_) {
host_connector_ = std::make_shared<network::TcpSocketConnector>( host_connector_ = std::make_shared<network::TcpSocketConnector>(
boost::asio::ip::address_v4::from_string("127.0.0.1"), boost::asio::ip::address_v4::from_string(loopback_address),
default_host_listen_port, runtime_, default_host_listen_port, runtime_,
std::make_shared< std::make_shared<
network::DelegateConnectionCreator<boost::asio::ip::tcp>>( network::DelegateConnectionCreator<boost::asio::ip::tcp>>(
@ -113,18 +124,19 @@ void AdbMessageProcessor::wait_for_host_connection() {
// Notify the adb host instance so that it knows on which port our // Notify the adb host instance so that it knows on which port our
// proxy is waiting for incoming connections. // proxy is waiting for incoming connections.
auto messenger = std::make_shared<network::TcpSocketMessenger>( auto messenger = std::make_shared<network::TcpSocketMessenger>(
boost::asio::ip::address_v4::from_string("127.0.0.1"), boost::asio::ip::address_v4::from_string(loopback_address), default_adb_client_port, runtime_);
default_adb_client_port, runtime_); auto message = utils::string_format("host:emulator:%d", default_host_listen_port);
auto message = auto handshake = utils::string_format("%04x%s", message.size(), message.c_str());
utils::string_format("host:emulator:%d", default_host_listen_port);
auto handshake =
utils::string_format("%04x%s", message.size(), message.c_str());
messenger->send(handshake.data(), handshake.size()); messenger->send(handshake.data(), handshake.size());
} catch (std::exception &) { } catch (...) {
// Try again later when the host adb service is maybe available // Try again later when the host adb service is maybe available
host_notify_timer_.cancel();
host_notify_timer_.expires_from_now(default_adb_wait_time); host_notify_timer_.expires_from_now(default_adb_wait_time);
host_notify_timer_.async_wait( host_notify_timer_.async_wait([this](const boost::system::error_code &err) {
[&](const boost::system::error_code &) { wait_for_host_connection(); }); if (err)
return;
wait_for_host_connection();
});
} }
} }
@ -148,15 +160,12 @@ void AdbMessageProcessor::on_host_connection(std::shared_ptr<boost::asio::basic_
void AdbMessageProcessor::read_next_host_message() { void AdbMessageProcessor::read_next_host_message() {
auto callback = std::bind(&AdbMessageProcessor::on_host_read_size, this, _1, _2); auto callback = std::bind(&AdbMessageProcessor::on_host_read_size, this, _1, _2);
host_messenger_->async_receive_msg(callback, host_messenger_->async_receive_msg(callback, boost::asio::buffer(host_buffer_));
boost::asio::buffer(host_buffer_));
} }
void AdbMessageProcessor::on_host_read_size( void AdbMessageProcessor::on_host_read_size(const boost::system::error_code &error, std::size_t bytes_read) {
const boost::system::error_code &error, std::size_t bytes_read) {
if (error) { if (error) {
state_ = closed_by_host; state_ = closed_by_host;
messenger_->close();
return; return;
} }

View file

@ -27,6 +27,8 @@
#include <boost/asio.hpp> #include <boost/asio.hpp>
#include <mutex>
namespace anbox { namespace anbox {
namespace qemu { namespace qemu {
class AdbMessageProcessor : public network::MessageProcessor { class AdbMessageProcessor : public network::MessageProcessor {
@ -66,6 +68,9 @@ class AdbMessageProcessor : public network::MessageProcessor {
std::shared_ptr<network::TcpSocketMessenger> host_messenger_; std::shared_ptr<network::TcpSocketMessenger> host_messenger_;
std::array<std::uint8_t, 8192> host_buffer_; std::array<std::uint8_t, 8192> host_buffer_;
boost::asio::deadline_timer host_notify_timer_; boost::asio::deadline_timer host_notify_timer_;
std::unique_lock<std::mutex> lock_;
static std::mutex active_instance_;
}; };
} // namespace graphics } // namespace graphics
} // namespace anbox } // namespace anbox