diff --git a/src/third_party/websocketpp/websocketpp/base64/base64.cpp b/src/third_party/websocketpp/websocketpp/base64/base64.cpp deleted file mode 100644 index 87388022..00000000 --- a/src/third_party/websocketpp/websocketpp/base64/base64.cpp +++ /dev/null @@ -1,123 +0,0 @@ -/* - base64.cpp and base64.h - - Copyright (C) 2004-2008 René Nyffenegger - - This source code is provided 'as-is', without any express or implied - warranty. In no event will the author be held liable for any damages - arising from the use of this software. - - Permission is granted to anyone to use this software for any purpose, - including commercial applications, and to alter it and redistribute it - freely, subject to the following restrictions: - - 1. The origin of this source code must not be misrepresented; you must not - claim that you wrote the original source code. If you use this source code - in a product, an acknowledgment in the product documentation would be - appreciated but is not required. - - 2. Altered source versions must be plainly marked as such, and must not be - misrepresented as being the original source code. - - 3. This notice may not be removed or altered from any source distribution. - - René Nyffenegger rene.nyffenegger@adp-gmbh.ch - -*/ - -#include "base64.h" -#include - -static const std::string base64_chars = - "ABCDEFGHIJKLMNOPQRSTUVWXYZ" - "abcdefghijklmnopqrstuvwxyz" - "0123456789+/"; - - -static inline bool is_base64(unsigned char c) { - return (isalnum(c) || (c == '+') || (c == '/')); -} - -std::string base64_encode(unsigned char const* bytes_to_encode, unsigned int in_len) { - std::string ret; - int i = 0; - int j = 0; - unsigned char char_array_3[3]; - unsigned char char_array_4[4]; - - while (in_len--) { - char_array_3[i++] = *(bytes_to_encode++); - if (i == 3) { - char_array_4[0] = (char_array_3[0] & 0xfc) >> 2; - char_array_4[1] = ((char_array_3[0] & 0x03) << 4) + ((char_array_3[1] & 0xf0) >> 4); - char_array_4[2] = ((char_array_3[1] & 0x0f) << 2) + ((char_array_3[2] & 0xc0) >> 6); - char_array_4[3] = char_array_3[2] & 0x3f; - - for(i = 0; (i <4) ; i++) - ret += base64_chars[char_array_4[i]]; - i = 0; - } - } - - if (i) - { - for(j = i; j < 3; j++) - char_array_3[j] = '\0'; - - char_array_4[0] = (char_array_3[0] & 0xfc) >> 2; - char_array_4[1] = ((char_array_3[0] & 0x03) << 4) + ((char_array_3[1] & 0xf0) >> 4); - char_array_4[2] = ((char_array_3[1] & 0x0f) << 2) + ((char_array_3[2] & 0xc0) >> 6); - char_array_4[3] = char_array_3[2] & 0x3f; - - for (j = 0; (j < i + 1); j++) - ret += base64_chars[char_array_4[j]]; - - while((i++ < 3)) - ret += '='; - - } - - return ret; - -} - -std::string base64_decode(std::string const& encoded_string) { - size_t in_len = encoded_string.size(); - int i = 0; - int j = 0; - int in_ = 0; - unsigned char char_array_4[4], char_array_3[3]; - std::string ret; - - while (in_len-- && ( encoded_string[in_] != '=') && is_base64(encoded_string[in_])) { - char_array_4[i++] = encoded_string[in_]; in_++; - if (i ==4) { - for (i = 0; i <4; i++) - char_array_4[i] = static_cast(base64_chars.find(char_array_4[i])); - - char_array_3[0] = (char_array_4[0] << 2) + ((char_array_4[1] & 0x30) >> 4); - char_array_3[1] = ((char_array_4[1] & 0xf) << 4) + ((char_array_4[2] & 0x3c) >> 2); - char_array_3[2] = ((char_array_4[2] & 0x3) << 6) + char_array_4[3]; - - for (i = 0; (i < 3); i++) - ret += char_array_3[i]; - i = 0; - } - } - - if (i) { - for (j = i; j <4; j++) - char_array_4[j] = 0; - - for (j = 0; j <4; j++) - char_array_4[j] = static_cast(base64_chars.find(char_array_4[j])); - - char_array_3[0] = (char_array_4[0] << 2) + ((char_array_4[1] & 0x30) >> 4); - char_array_3[1] = ((char_array_4[1] & 0xf) << 4) + ((char_array_4[2] & 0x3c) >> 2); - char_array_3[2] = ((char_array_4[2] & 0x3) << 6) + char_array_4[3]; - - for (j = 0; (j < i - 1); j++) ret += char_array_3[j]; - } - - return ret; -} diff --git a/src/third_party/websocketpp/websocketpp/base64/base64.h b/src/third_party/websocketpp/websocketpp/base64/base64.h deleted file mode 100644 index 65d5db8b..00000000 --- a/src/third_party/websocketpp/websocketpp/base64/base64.h +++ /dev/null @@ -1,4 +0,0 @@ -#include - -std::string base64_encode(unsigned char const* , unsigned int len); -std::string base64_decode(std::string const& s); diff --git a/src/third_party/websocketpp/websocketpp/base64/base64.hpp b/src/third_party/websocketpp/websocketpp/base64/base64.hpp index e0e9c4e7..6572c495 100644 --- a/src/third_party/websocketpp/websocketpp/base64/base64.hpp +++ b/src/third_party/websocketpp/websocketpp/base64/base64.hpp @@ -38,7 +38,9 @@ #include -static const std::string base64_chars = +namespace websocketpp { + +static std::string const base64_chars = "ABCDEFGHIJKLMNOPQRSTUVWXYZ" "abcdefghijklmnopqrstuvwxyz" "0123456789+/"; @@ -50,7 +52,7 @@ static inline bool is_base64(unsigned char c) { (c >= 97 && c <= 122)); // a-z } -inline std::string base64_encode(unsigned char const* bytes_to_encode, unsigned +inline std::string base64_encode(unsigned char const * bytes_to_encode, unsigned int in_len) { std::string ret; @@ -100,11 +102,11 @@ inline std::string base64_encode(unsigned char const* bytes_to_encode, unsigned return ret; } -inline std::string base64_encode(const std::string & data) { +inline std::string base64_encode(std::string const & data) { return base64_encode(reinterpret_cast(data.data()),data.size()); } -inline std::string base64_decode(std::string const& encoded_string) { +inline std::string base64_decode(std::string const & encoded_string) { size_t in_len = encoded_string.size(); int i = 0; int j = 0; @@ -149,4 +151,6 @@ inline std::string base64_decode(std::string const& encoded_string) { return ret; } +} // namespace websocketpp + #endif // _BASE64_HPP_ diff --git a/src/third_party/websocketpp/websocketpp/close.hpp b/src/third_party/websocketpp/websocketpp/close.hpp index d0e87663..eee44f61 100644 --- a/src/third_party/websocketpp/websocketpp/close.hpp +++ b/src/third_party/websocketpp/websocketpp/close.hpp @@ -51,6 +51,26 @@ namespace status { /// A blank value for internal use. static value const blank = 0; + /// Close the connection without a WebSocket close handshake. + /** + * This special value requests that the WebSocket connection be closed + * without performing the WebSocket closing handshake. This does not comply + * with RFC6455, but should be safe to do if necessary. This could be useful + * for clients that need to disconnect quickly and cannot afford the + * complete handshake. + */ + static value const omit_handshake = 1; + + /// Close the connection with a forced TCP drop. + /** + * This special value requests that the WebSocket connection be closed by + * forcibly dropping the TCP connection. This will leave the other side of + * the connection with a broken connection and some expensive timeouts. this + * should not be done except in extreme cases or in cases of malicious + * remote endpoints. + */ + static value const force_tcp_drop = 2; + /// Normal closure, meaning that the purpose for which the connection was /// established has been fulfilled. static value const normal = 1000; diff --git a/src/third_party/websocketpp/websocketpp/config/core.hpp b/src/third_party/websocketpp/websocketpp/config/core.hpp index 69744a19..e400da50 100644 --- a/src/third_party/websocketpp/websocketpp/config/core.hpp +++ b/src/third_party/websocketpp/websocketpp/config/core.hpp @@ -91,6 +91,11 @@ struct core { /// RNG policies typedef websocketpp::random::none::int_generator rng_type; + /// Controls compile time enabling/disabling of thread syncronization + /// code Disabling can provide a minor performance improvement to single + /// threaded applications + static bool const enable_multithreading = true; + struct transport_config { typedef type::concurrency_type concurrency_type; typedef type::elog_type elog_type; @@ -98,6 +103,11 @@ struct core { typedef type::request_type request_type; typedef type::response_type response_type; + /// Controls compile time enabling/disabling of thread syncronization + /// code Disabling can provide a minor performance improvement to single + /// threaded applications + static bool const enable_multithreading = true; + /// Default timer values (in ms) /// Length of time to wait for socket pre-initialization @@ -179,7 +189,7 @@ struct core { websocketpp::log::alevel::all ^ websocketpp::log::alevel::devel; /// - static const size_t connection_read_buffer_size = 512; + static const size_t connection_read_buffer_size = 16384; /// Drop connections immediately on protocol error. /** diff --git a/src/third_party/websocketpp/websocketpp/config/core_client.hpp b/src/third_party/websocketpp/websocketpp/config/core_client.hpp index b7329756..328bcc86 100644 --- a/src/third_party/websocketpp/websocketpp/config/core_client.hpp +++ b/src/third_party/websocketpp/websocketpp/config/core_client.hpp @@ -92,6 +92,11 @@ struct core_client { typedef websocketpp::random::random_device::int_generator rng_type; + /// Controls compile time enabling/disabling of thread syncronization code + /// Disabling can provide a minor performance improvement to single threaded + /// applications + static bool const enable_multithreading = true; + struct transport_config { typedef type::concurrency_type concurrency_type; typedef type::elog_type elog_type; @@ -99,6 +104,11 @@ struct core_client { typedef type::request_type request_type; typedef type::response_type response_type; + /// Controls compile time enabling/disabling of thread syncronization + /// code Disabling can provide a minor performance improvement to single + /// threaded applications + static bool const enable_multithreading = true; + /// Default timer values (in ms) /// Length of time to wait for socket pre-initialization @@ -180,7 +190,7 @@ struct core_client { websocketpp::log::alevel::all ^ websocketpp::log::alevel::devel; /// - static const size_t connection_read_buffer_size = 512; + static const size_t connection_read_buffer_size = 16384; /// Drop connections immediately on protocol error. /** diff --git a/src/third_party/websocketpp/websocketpp/config/debug.hpp b/src/third_party/websocketpp/websocketpp/config/debug.hpp index 0f2fc65c..691229a6 100644 --- a/src/third_party/websocketpp/websocketpp/config/debug.hpp +++ b/src/third_party/websocketpp/websocketpp/config/debug.hpp @@ -92,6 +92,11 @@ struct debug_core { /// RNG policies typedef websocketpp::random::none::int_generator rng_type; + /// Controls compile time enabling/disabling of thread syncronization + /// code Disabling can provide a minor performance improvement to single + /// threaded applications + static bool const enable_multithreading = true; + struct transport_config { typedef type::concurrency_type concurrency_type; typedef type::elog_type elog_type; @@ -99,6 +104,11 @@ struct debug_core { typedef type::request_type request_type; typedef type::response_type response_type; + /// Controls compile time enabling/disabling of thread syncronization + /// code Disabling can provide a minor performance improvement to single + /// threaded applications + static bool const enable_multithreading = true; + /// Default timer values (in ms) /// Length of time to wait for socket pre-initialization @@ -180,7 +190,7 @@ struct debug_core { websocketpp::log::alevel::all; /// - static const size_t connection_read_buffer_size = 512; + static const size_t connection_read_buffer_size = 16384; /// Drop connections immediately on protocol error. /** diff --git a/src/third_party/websocketpp/websocketpp/connection.hpp b/src/third_party/websocketpp/websocketpp/connection.hpp index c8f1bc60..c8be3633 100644 --- a/src/third_party/websocketpp/websocketpp/connection.hpp +++ b/src/third_party/websocketpp/websocketpp/connection.hpp @@ -149,6 +149,10 @@ typedef lib::function validate_handler; */ typedef lib::function http_handler; +// +typedef lib::function read_handler; +typedef lib::function write_frame_handler; + // constants related to the default WebSocket protocol versions available #ifdef _WEBSOCKETPP_INITIALIZER_LISTS_ // simplified C++11 version /// Container that stores the list of protocol versions supported @@ -278,7 +282,21 @@ public: explicit connection(bool is_server, std::string const & ua, alog_type& alog, elog_type& elog, rng_type & rng) : transport_con_type(is_server,alog,elog) + , m_handle_read_frame(lib::bind( + &type::handle_read_frame, + this, + lib::placeholders::_1, + lib::placeholders::_2 + )) + , m_write_frame_handler(lib::bind( + &type::handle_write_frame, + this, + lib::placeholders::_1 + )) , m_user_agent(ua) + , m_open_handshake_timeout_dur(config::timeout_open_handshake) + , m_close_handshake_timeout_dur(config::timeout_close_handshake) + , m_pong_timeout_dur(config::timeout_pong) , m_state(session::state::connecting) , m_internal_state(session::internal_state::USER_INIT) , m_msg_manager(new con_msg_manager_type()) @@ -437,6 +455,79 @@ public: m_message_handler = h; } + ///////////////////////// + // Connection timeouts // + ///////////////////////// + + /// Set open handshake timeout + /** + * Sets the length of time the library will wait after an opening handshake + * has been initiated before cancelling it. This can be used to prevent + * excessive wait times for outgoing clients or excessive resource usage + * from broken clients or DoS attacks on servers. + * + * Connections that time out will have their fail handlers called with the + * open_handshake_timeout error code. + * + * The default value is specified via the compile time config value + * 'timeout_open_handshake'. The default value in the core config + * is 5000ms. A value of 0 will disable the timer entirely. + * + * To be effective, the transport you are using must support timers. See + * the documentation for your transport policy for details about its + * timer support. + * + * @param dur The length of the open handshake timeout in ms + */ + void set_open_handshake_timeout(long dur) { + m_open_handshake_timeout_dur = dur; + } + + /// Set close handshake timeout + /** + * Sets the length of time the library will wait after a closing handshake + * has been initiated before cancelling it. This can be used to prevent + * excessive wait times for outgoing clients or excessive resource usage + * from broken clients or DoS attacks on servers. + * + * Connections that time out will have their close handlers called with the + * close_handshake_timeout error code. + * + * The default value is specified via the compile time config value + * 'timeout_close_handshake'. The default value in the core config + * is 5000ms. A value of 0 will disable the timer entirely. + * + * To be effective, the transport you are using must support timers. See + * the documentation for your transport policy for details about its + * timer support. + * + * @param dur The length of the close handshake timeout in ms + */ + void set_close_handshake_timeout(long dur) { + m_close_handshake_timeout_dur = dur; + } + + /// Set pong timeout + /** + * Sets the length of time the library will wait for a pong response to a + * ping. This can be used as a keepalive or to detect broken connections. + * + * Pong responses that time out will have the pong timeout handler called. + * + * The default value is specified via the compile time config value + * 'timeout_pong'. The default value in the core config + * is 5000ms. A value of 0 will disable the timer entirely. + * + * To be effective, the transport you are using must support timers. See + * the documentation for your transport policy for details about its + * timer support. + * + * @param dur The length of the pong timeout in ms + */ + void set_pong_timeout(long dur) { + m_pong_timeout_dur = dur; + } + ////////////////////////////////// // Uncategorized public methods // ////////////////////////////////// @@ -1034,7 +1125,7 @@ public: * @param ec A status code from the transport layer, zero on success, * non-zero otherwise. */ - void handle_write_frame(bool terminate, lib::error_code const & ec); + void handle_write_frame(lib::error_code const & ec); protected: void handle_transport_init(lib::error_code const & ec); @@ -1190,8 +1281,12 @@ private: */ void log_fail_result(); + // internal handler functions + read_handler m_handle_read_frame; + write_frame_handler m_write_frame_handler; + // static settings - const std::string m_user_agent; + std::string const m_user_agent; /// Pointer to the connection handle connection_hdl m_connection_hdl; @@ -1208,6 +1303,11 @@ private: validate_handler m_validate_handler; message_handler m_message_handler; + /// constant values + long m_open_handshake_timeout_dur; + long m_close_handshake_timeout_dur; + long m_pong_timeout_dur; + /// External connection state /** * Lock: m_connection_state_lock diff --git a/src/third_party/websocketpp/websocketpp/endpoint.hpp b/src/third_party/websocketpp/websocketpp/endpoint.hpp index 3c72b504..2df806a3 100644 --- a/src/third_party/websocketpp/websocketpp/endpoint.hpp +++ b/src/third_party/websocketpp/websocketpp/endpoint.hpp @@ -91,6 +91,9 @@ public: : m_alog(config::alog_level, &std::cout) , m_elog(config::elog_level, &std::cerr) , m_user_agent(::websocketpp::user_agent) + , m_open_handshake_timeout_dur(config::timeout_open_handshake) + , m_close_handshake_timeout_dur(config::timeout_close_handshake) + , m_pong_timeout_dur(config::timeout_pong) , m_is_server(is_server) { m_alog.set_channels(config::alog_level); @@ -269,6 +272,82 @@ public: m_message_handler = h; } + ///////////////////////// + // Connection timeouts // + ///////////////////////// + + /// Set open handshake timeout + /** + * Sets the length of time the library will wait after an opening handshake + * has been initiated before cancelling it. This can be used to prevent + * excessive wait times for outgoing clients or excessive resource usage + * from broken clients or DoS attacks on servers. + * + * Connections that time out will have their fail handlers called with the + * open_handshake_timeout error code. + * + * The default value is specified via the compile time config value + * 'timeout_open_handshake'. The default value in the core config + * is 5000ms. A value of 0 will disable the timer entirely. + * + * To be effective, the transport you are using must support timers. See + * the documentation for your transport policy for details about its + * timer support. + * + * @param dur The length of the open handshake timeout in ms + */ + void set_open_handshake_timeout(long dur) { + scoped_lock_type guard(m_mutex); + m_open_handshake_timeout_dur = dur; + } + + /// Set close handshake timeout + /** + * Sets the length of time the library will wait after a closing handshake + * has been initiated before cancelling it. This can be used to prevent + * excessive wait times for outgoing clients or excessive resource usage + * from broken clients or DoS attacks on servers. + * + * Connections that time out will have their close handlers called with the + * close_handshake_timeout error code. + * + * The default value is specified via the compile time config value + * 'timeout_close_handshake'. The default value in the core config + * is 5000ms. A value of 0 will disable the timer entirely. + * + * To be effective, the transport you are using must support timers. See + * the documentation for your transport policy for details about its + * timer support. + * + * @param dur The length of the close handshake timeout in ms + */ + void set_close_handshake_timeout(long dur) { + scoped_lock_type guard(m_mutex); + m_close_handshake_timeout_dur = dur; + } + + /// Set pong timeout + /** + * Sets the length of time the library will wait for a pong response to a + * ping. This can be used as a keepalive or to detect broken connections. + * + * Pong responses that time out will have the pong timeout handler called. + * + * The default value is specified via the compile time config value + * 'timeout_pong'. The default value in the core config + * is 5000ms. A value of 0 will disable the timer entirely. + * + * To be effective, the transport you are using must support timers. See + * the documentation for your transport policy for details about its + * timer support. + * + * @param dur The length of the pong timeout in ms + */ + void set_pong_timeout(long dur) { + scoped_lock_type guard(m_mutex); + m_pong_timeout_dur = dur; + } + /*************************************/ /* Connection pass through functions */ /*************************************/ @@ -415,13 +494,17 @@ private: validate_handler m_validate_handler; message_handler m_message_handler; + long m_open_handshake_timeout_dur; + long m_close_handshake_timeout_dur; + long m_pong_timeout_dur; + rng_type m_rng; // static settings bool const m_is_server; // endpoint state - mutex_type m_mutex; + mutable mutex_type m_mutex; }; } // namespace websocketpp diff --git a/src/third_party/websocketpp/websocketpp/http/constants.hpp b/src/third_party/websocketpp/websocketpp/http/constants.hpp index 3edc3eb0..be3006aa 100644 --- a/src/third_party/websocketpp/websocketpp/http/constants.hpp +++ b/src/third_party/websocketpp/websocketpp/http/constants.hpp @@ -53,7 +53,7 @@ namespace http { static char const header_delimiter[] = "\r\n"; /// Literal value of the HTTP header separator - static char const header_separator[] = ": "; + static char const header_separator[] = ":"; /// Literal value of an empty header static std::string const empty_header = ""; diff --git a/src/third_party/websocketpp/websocketpp/http/impl/parser.hpp b/src/third_party/websocketpp/websocketpp/http/impl/parser.hpp index 0ceac2ed..d1d90813 100644 --- a/src/third_party/websocketpp/websocketpp/http/impl/parser.hpp +++ b/src/third_party/websocketpp/websocketpp/http/impl/parser.hpp @@ -147,8 +147,8 @@ inline void parser::process_header(std::string::iterator begin, throw exception("Invalid header line",status_code::bad_request); } - append_header(std::string(begin,cursor), - std::string(cursor+sizeof(header_separator)-1,end)); + append_header(strip_lws(std::string(begin,cursor)), + strip_lws(std::string(cursor+sizeof(header_separator)-1,end))); } inline std::string parser::raw_headers() const { diff --git a/src/third_party/websocketpp/websocketpp/http/parser.hpp b/src/third_party/websocketpp/websocketpp/http/parser.hpp index 0bb52c2b..771c5397 100644 --- a/src/third_party/websocketpp/websocketpp/http/parser.hpp +++ b/src/third_party/websocketpp/websocketpp/http/parser.hpp @@ -367,6 +367,13 @@ InputIterator extract_parameters(InputIterator begin, InputIterator end, return cursor; } +inline std::string strip_lws(std::string const & input) { + std::string::const_iterator begin = extract_all_lws(input.begin(),input.end()); + std::string::const_reverse_iterator end = extract_all_lws(input.rbegin(),input.rend()); + + return std::string(begin,end.base()); +} + /// Base HTTP parser /** * Includes methods and data elements common to all types of HTTP messages such diff --git a/src/third_party/websocketpp/websocketpp/impl/connection_impl.hpp b/src/third_party/websocketpp/websocketpp/impl/connection_impl.hpp index 980e2ecc..0449e746 100644 --- a/src/third_party/websocketpp/websocketpp/impl/connection_impl.hpp +++ b/src/third_party/websocketpp/websocketpp/impl/connection_impl.hpp @@ -163,15 +163,17 @@ void connection::ping(const std::string& payload, lib::error_code& ec) { m_ping_timer->cancel(); } - m_ping_timer = transport_con_type::set_timer( - config::timeout_pong, - lib::bind( - &type::handle_pong_timeout, - type::get_shared(), - payload, - lib::placeholders::_1 - ) - ); + if (m_pong_timeout_dur > 0) { + m_ping_timer = transport_con_type::set_timer( + m_pong_timeout_dur, + lib::bind( + &type::handle_pong_timeout, + type::get_shared(), + payload, + lib::placeholders::_1 + ) + ); + } if (!m_ping_timer) { // Our transport doesn't support timers @@ -270,8 +272,8 @@ void connection::pong(const std::string& payload) { } template -void connection::close(const close::status::value code, - const std::string & reason, lib::error_code & ec) +void connection::close(close::status::value const code, + std::string const & reason, lib::error_code & ec) { m_alog.write(log::alevel::devel,"connection close"); @@ -288,8 +290,8 @@ void connection::close(const close::status::value code, } template -void connection::close(const close::status::value code, - const std::string & reason) +void connection::close(close::status::value const code, + std::string const & reason) { lib::error_code ec; close(code,reason,ec); @@ -850,11 +852,11 @@ void connection::handle_read_frame(const lib::error_code& ec, } // Boundaries checking. TODO: How much of this should be done? - if (bytes_transferred > config::connection_read_buffer_size) { + /*if (bytes_transferred > config::connection_read_buffer_size) { m_elog.write(log::elevel::fatal,"Fatal boundaries checking error"); this->terminate(make_error_code(error::general)); return; - } + }*/ size_t p = 0; @@ -942,12 +944,13 @@ void connection::handle_read_frame(const lib::error_code& ec, 1, m_buf, config::connection_read_buffer_size, - lib::bind( + /*lib::bind( &type::handle_read_frame, type::get_shared(), lib::placeholders::_1, lib::placeholders::_2 - ) + )*/ + m_handle_read_frame ); } @@ -1253,14 +1256,16 @@ void connection::send_http_request() { "Raw Handshake request:\n"+m_handshake_buffer); } - m_handshake_timer = transport_con_type::set_timer( - config::timeout_open_handshake, - lib::bind( - &type::handle_open_handshake_timeout, - type::get_shared(), - lib::placeholders::_1 - ) - ); + if (m_open_handshake_timeout_dur > 0) { + m_handshake_timer = transport_con_type::set_timer( + m_open_handshake_timeout_dur, + lib::bind( + &type::handle_open_handshake_timeout, + type::get_shared(), + lib::placeholders::_1 + ) + ); + } transport_con_type::async_write( m_handshake_buffer.data(), @@ -1429,7 +1434,7 @@ void connection::handle_close_handshake_timeout( } template -void connection::terminate(const lib::error_code & ec) { +void connection::terminate(lib::error_code const & ec) { if (m_alog.static_test(log::alevel::devel)) { m_alog.write(log::alevel::devel,"connection terminate"); } @@ -1459,6 +1464,8 @@ void connection::terminate(const lib::error_code & ec) { return; } + // TODO: choose between shutdown and close based on error code sent + transport_con_type::async_shutdown( lib::bind( &type::handle_terminate, @@ -1471,7 +1478,7 @@ void connection::terminate(const lib::error_code & ec) { template void connection::handle_terminate(terminate_status tstat, - const lib::error_code& ec) + lib::error_code const & ec) { if (m_alog.static_test(log::alevel::devel)) { m_alog.write(log::alevel::devel,"connection handle_terminate"); @@ -1540,8 +1547,8 @@ void connection::write_frame() { m_write_flag = true; } - const std::string& header = m_current_msg->get_header(); - const std::string& payload = m_current_msg->get_payload(); + std::string const & header = m_current_msg->get_header(); + std::string const & payload = m_current_msg->get_payload(); m_send_buffer.push_back(transport::buffer(header.c_str(),header.size())); m_send_buffer.push_back(transport::buffer(payload.c_str(),payload.size())); @@ -1563,25 +1570,28 @@ void connection::write_frame() { } } + transport_con_type::async_write( m_send_buffer, - lib::bind( + /*lib::bind( &type::handle_write_frame, type::get_shared(), m_current_msg->get_terminal(), lib::placeholders::_1 - ) + )*/ + m_write_frame_handler ); } template -void connection::handle_write_frame(bool terminate, - const lib::error_code& ec) +void connection::handle_write_frame(lib::error_code const & ec) { if (m_alog.static_test(log::alevel::devel)) { m_alog.write(log::alevel::devel,"connection handle_write_frame"); } + bool terminate = m_current_msg->get_terminal(); + m_send_buffer.clear(); m_current_msg.reset(); @@ -1791,16 +1801,19 @@ void connection::process_control_frame(typename template lib::error_code connection::send_close_ack(close::status::value code, - const std::string &reason) + std::string const & reason) { return send_close_frame(code,reason,true,m_is_server); } template lib::error_code connection::send_close_frame(close::status::value code, - const std::string &reason, bool ack, bool terminal) + std::string const & reason, bool ack, bool terminal) { m_alog.write(log::alevel::devel,"send_close_frame"); + + // check for special codes + // If silent close is set, respect it and blank out close information // Otherwise use whatever has been specified in the parameters. If // parameters specifies close::status::blank then determine what to do @@ -1861,14 +1874,16 @@ lib::error_code connection::send_close_frame(close::status::value code, // Start a timer so we don't wait forever for the acknowledgement close // frame - m_handshake_timer = transport_con_type::set_timer( - config::timeout_close_handshake, - lib::bind( - &type::handle_close_handshake_timeout, - type::get_shared(), - lib::placeholders::_1 - ) - ); + if (m_close_handshake_timeout_dur > 0) { + m_handshake_timer = transport_con_type::set_timer( + m_close_handshake_timeout_dur, + lib::bind( + &type::handle_close_handshake_timeout, + type::get_shared(), + lib::placeholders::_1 + ) + ); + } bool needs_writing = false; { diff --git a/src/third_party/websocketpp/websocketpp/impl/endpoint_impl.hpp b/src/third_party/websocketpp/websocketpp/impl/endpoint_impl.hpp index e30c0d9d..0dd14c5d 100644 --- a/src/third_party/websocketpp/websocketpp/impl/endpoint_impl.hpp +++ b/src/third_party/websocketpp/websocketpp/impl/endpoint_impl.hpp @@ -40,6 +40,7 @@ endpoint::create_connection() { return connection_ptr(); }*/ + //scoped_lock_type guard(m_mutex); // Create a connection on the heap and manage it using a shared pointer connection_ptr con(new connection_type(m_is_server,m_user_agent,m_alog, m_elog, m_rng)); @@ -63,6 +64,16 @@ endpoint::create_connection() { con->set_http_handler(m_http_handler); con->set_validate_handler(m_validate_handler); con->set_message_handler(m_message_handler); + + if (m_open_handshake_timeout_dur == config::timeout_open_handshake) { + con->set_open_handshake_timeout(m_open_handshake_timeout_dur); + } + if (m_close_handshake_timeout_dur == config::timeout_close_handshake) { + con->set_close_handshake_timeout(m_close_handshake_timeout_dur); + } + if (m_pong_timeout_dur == config::timeout_pong) { + con->set_pong_timeout(m_pong_timeout_dur); + } lib::error_code ec; diff --git a/src/third_party/websocketpp/websocketpp/sha1/sha1.cpp b/src/third_party/websocketpp/websocketpp/sha1/sha1.cpp deleted file mode 100644 index 661562a3..00000000 --- a/src/third_party/websocketpp/websocketpp/sha1/sha1.cpp +++ /dev/null @@ -1,185 +0,0 @@ -/* - Copyright (c) 2011, Micael Hildenborg - All rights reserved. - - Redistribution and use in source and binary forms, with or without - modification, are permitted provided that the following conditions are met: - * Redistributions of source code must retain the above copyright - notice, this list of conditions and the following disclaimer. - * Redistributions in binary form must reproduce the above copyright - notice, this list of conditions and the following disclaimer in the - documentation and/or other materials provided with the distribution. - * Neither the name of Micael Hildenborg nor the - names of its contributors may be used to endorse or promote products - derived from this software without specific prior written permission. - - THIS SOFTWARE IS PROVIDED BY Micael Hildenborg ''AS IS'' AND ANY - EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED - WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE ARE - DISCLAIMED. IN NO EVENT SHALL Micael Hildenborg BE LIABLE FOR ANY - DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES - (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; - LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND - ON ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT - (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS - SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE. - */ - -/* - Contributors: - Gustav - Several members in the gamedev.se forum. - Gregory Petrosyan - */ - -#include "sha1.h" - -namespace sha1 -{ - namespace // local - { - // Rotate an integer value to left. - inline const unsigned int rol(const unsigned int value, - const unsigned int steps) - { - return ((value << steps) | (value >> (32 - steps))); - } - - // Sets the first 16 integers in the buffert to zero. - // Used for clearing the W buffert. - inline void clearWBuffert(unsigned int* buffert) - { - for (int pos = 16; --pos >= 0;) - { - buffert[pos] = 0; - } - } - - void innerHash(unsigned int* result, unsigned int* w) - { - unsigned int a = result[0]; - unsigned int b = result[1]; - unsigned int c = result[2]; - unsigned int d = result[3]; - unsigned int e = result[4]; - - int round = 0; - - #define sha1macro(func,val) \ - { \ - const unsigned int t = rol(a, 5) + (func) + e + val + w[round]; \ - e = d; \ - d = c; \ - c = rol(b, 30); \ - b = a; \ - a = t; \ - } - - while (round < 16) - { - sha1macro((b & c) | (~b & d), 0x5a827999) - ++round; - } - while (round < 20) - { - w[round] = rol((w[round - 3] ^ w[round - 8] ^ w[round - 14] ^ w[round - 16]), 1); - sha1macro((b & c) | (~b & d), 0x5a827999) - ++round; - } - while (round < 40) - { - w[round] = rol((w[round - 3] ^ w[round - 8] ^ w[round - 14] ^ w[round - 16]), 1); - sha1macro(b ^ c ^ d, 0x6ed9eba1) - ++round; - } - while (round < 60) - { - w[round] = rol((w[round - 3] ^ w[round - 8] ^ w[round - 14] ^ w[round - 16]), 1); - sha1macro((b & c) | (b & d) | (c & d), 0x8f1bbcdc) - ++round; - } - while (round < 80) - { - w[round] = rol((w[round - 3] ^ w[round - 8] ^ w[round - 14] ^ w[round - 16]), 1); - sha1macro(b ^ c ^ d, 0xca62c1d6) - ++round; - } - - #undef sha1macro - - result[0] += a; - result[1] += b; - result[2] += c; - result[3] += d; - result[4] += e; - } - } // namespace - - void calc(const void* src, const int bytelength, unsigned char* hash) - { - // Init the result array. - unsigned int result[5] = { 0x67452301, 0xefcdab89, 0x98badcfe, 0x10325476, 0xc3d2e1f0 }; - - // Cast the void src pointer to be the byte array we can work with. - const unsigned char* sarray = (const unsigned char*) src; - - // The reusable round buffer - unsigned int w[80]; - - // Loop through all complete 64byte blocks. - const int endOfFullBlocks = bytelength - 64; - int endCurrentBlock; - int currentBlock = 0; - - while (currentBlock <= endOfFullBlocks) - { - endCurrentBlock = currentBlock + 64; - - // Init the round buffer with the 64 byte block data. - for (int roundPos = 0; currentBlock < endCurrentBlock; currentBlock += 4) - { - // This line will swap endian on big endian and keep endian on little endian. - w[roundPos++] = (unsigned int) sarray[currentBlock + 3] - | (((unsigned int) sarray[currentBlock + 2]) << 8) - | (((unsigned int) sarray[currentBlock + 1]) << 16) - | (((unsigned int) sarray[currentBlock]) << 24); - } - innerHash(result, w); - } - - // Handle the last and not full 64 byte block if existing. - endCurrentBlock = bytelength - currentBlock; - clearWBuffert(w); - int lastBlockBytes = 0; - for (;lastBlockBytes < endCurrentBlock; ++lastBlockBytes) - { - w[lastBlockBytes >> 2] |= (unsigned int) sarray[lastBlockBytes + currentBlock] << ((3 - (lastBlockBytes & 3)) << 3); - } - w[lastBlockBytes >> 2] |= 0x80 << ((3 - (lastBlockBytes & 3)) << 3); - if (endCurrentBlock >= 56) - { - innerHash(result, w); - clearWBuffert(w); - } - w[15] = bytelength << 3; - innerHash(result, w); - - // Store hash in result pointer, and make sure we get in in the correct order on both endian models. - for (int hashByte = 20; --hashByte >= 0;) - { - hash[hashByte] = (result[hashByte >> 2] >> (((3 - hashByte) & 0x3) << 3)) & 0xff; - } - } - - void toHexString(const unsigned char* hash, char* hexstring) - { - const char hexDigits[] = { "0123456789abcdef" }; - - for (int hashByte = 20; --hashByte >= 0;) - { - hexstring[hashByte << 1] = hexDigits[(hash[hashByte] >> 4) & 0xf]; - hexstring[(hashByte << 1) + 1] = hexDigits[hash[hashByte] & 0xf]; - } - hexstring[40] = 0; - } -} // namespace sha1 diff --git a/src/third_party/websocketpp/websocketpp/sha1/sha1.h b/src/third_party/websocketpp/websocketpp/sha1/sha1.h deleted file mode 100644 index 96d8ce4d..00000000 --- a/src/third_party/websocketpp/websocketpp/sha1/sha1.h +++ /dev/null @@ -1,49 +0,0 @@ -/* - Copyright (c) 2011, Micael Hildenborg - All rights reserved. - - Redistribution and use in source and binary forms, with or without - modification, are permitted provided that the following conditions are met: - * Redistributions of source code must retain the above copyright - notice, this list of conditions and the following disclaimer. - * Redistributions in binary form must reproduce the above copyright - notice, this list of conditions and the following disclaimer in the - documentation and/or other materials provided with the distribution. - * Neither the name of Micael Hildenborg nor the - names of its contributors may be used to endorse or promote products - derived from this software without specific prior written permission. - - THIS SOFTWARE IS PROVIDED BY Micael Hildenborg ''AS IS'' AND ANY - EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED - WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE ARE - DISCLAIMED. IN NO EVENT SHALL Micael Hildenborg BE LIABLE FOR ANY - DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES - (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; - LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND - ON ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT - (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS - SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE. - */ - -#ifndef SHA1_DEFINED -#define SHA1_DEFINED - -namespace sha1 -{ - - /** - @param src points to any kind of data to be hashed. - @param bytelength the number of bytes to hash from the src pointer. - @param hash should point to a buffer of at least 20 bytes of size for storing the sha1 result in. - */ - void calc(const void* src, const int bytelength, unsigned char* hash); - - /** - @param hash is 20 bytes of sha1 hash. This is the same data that is the result from the calc function. - @param hexstring should point to a buffer of at least 41 bytes of size for storing the hexadecimal representation of the hash. A zero will be written at position 40, so the buffer will be a valid zero ended string. - */ - void toHexString(const unsigned char* hash, char* hexstring); - -} // namespace sha1 - -#endif // SHA1_DEFINED diff --git a/src/third_party/websocketpp/websocketpp/transport/asio/base.hpp b/src/third_party/websocketpp/websocketpp/transport/asio/base.hpp index 35c80ef2..28ac933c 100644 --- a/src/third_party/websocketpp/websocketpp/transport/asio/base.hpp +++ b/src/third_party/websocketpp/websocketpp/transport/asio/base.hpp @@ -34,6 +34,10 @@ #include +#include +#include +#include + #include namespace websocketpp { @@ -45,6 +49,110 @@ namespace transport { */ namespace asio { +// + +// Class to manage the memory to be used for handler-based custom allocation. +// It contains a single block of memory which may be returned for allocation +// requests. If the memory is in use when an allocation request is made, the +// allocator delegates allocation to the global heap. +class handler_allocator + : private boost::noncopyable +{ +public: + handler_allocator() + : in_use_(false) + { + } + + void* allocate(std::size_t size) + { + if (!in_use_ && size < storage_.size) + { + in_use_ = true; + return storage_.address(); + } + else + { + return ::operator new(size); + } + } + + void deallocate(void* pointer) + { + if (pointer == storage_.address()) + { + in_use_ = false; + } + else + { + ::operator delete(pointer); + } + } + +private: + // Storage space used for handler-based custom memory allocation. + boost::aligned_storage<1024> storage_; + + // Whether the handler-based custom allocation storage has been used. + bool in_use_; +}; + +// Wrapper class template for handler objects to allow handler memory +// allocation to be customised. Calls to operator() are forwarded to the +// encapsulated handler. +template +class custom_alloc_handler +{ +public: + custom_alloc_handler(handler_allocator& a, Handler h) + : allocator_(a), + handler_(h) + { + } + + template + void operator()(Arg1 arg1) + { + handler_(arg1); + } + + template + void operator()(Arg1 arg1, Arg2 arg2) + { + handler_(arg1, arg2); + } + + friend void* asio_handler_allocate(std::size_t size, + custom_alloc_handler* this_handler) + { + return this_handler->allocator_.allocate(size); + } + + friend void asio_handler_deallocate(void* pointer, std::size_t /*size*/, + custom_alloc_handler* this_handler) + { + this_handler->allocator_.deallocate(pointer); + } + +private: + handler_allocator& allocator_; + Handler handler_; +}; + +// Helper function to wrap a handler object to add custom allocation. +template +inline custom_alloc_handler make_custom_alloc_handler( + handler_allocator& a, Handler h) +{ + return custom_alloc_handler(a, h); +} + + + + + + + // Forward declaration of class endpoint so that it can be friended/referenced // before being included. template @@ -53,6 +161,21 @@ class endpoint; typedef lib::function socket_shutdown_handler; +typedef lib::function async_read_handler; + +typedef lib::function async_write_handler; + +typedef lib::function pre_init_handler; + +// handle_timer: dynamic parameters, multiple copies +// handle_proxy_write +// handle_proxy_read +// handle_async_write +// handle_pre_init + + /// Asio transport errors namespace error { enum value { diff --git a/src/third_party/websocketpp/websocketpp/transport/asio/connection.hpp b/src/third_party/websocketpp/websocketpp/transport/asio/connection.hpp index c69ea604..f0aefa57 100644 --- a/src/third_party/websocketpp/websocketpp/transport/asio/connection.hpp +++ b/src/third_party/websocketpp/websocketpp/transport/asio/connection.hpp @@ -109,15 +109,47 @@ public: return socket_con_type::is_secure(); } - /// Sets the tcp init handler + /// Sets the tcp pre init handler /** - * The tcp init handler is called after the tcp connection has been - * established. + * The tcp pre init handler is called after the raw tcp connection has been + * established but before any additional wrappers (proxy connects, TLS + * handshakes, etc) have been performed. * - * @param h The handler to call on tcp init. + * @since 0.4.0-alpha1 + * + * @param h The handler to call on tcp pre init. + */ + void set_tcp_pre_init_handler(tcp_init_handler h) { + m_tcp_pre_init_handler = h; + } + + /// Sets the tcp pre init handler (deprecated) + /** + * The tcp pre init handler is called after the raw tcp connection has been + * established but before any additional wrappers (proxy connects, TLS + * handshakes, etc) have been performed. + * + * @deprecated Use set_tcp_pre_init_handler instead + * + * @param h The handler to call on tcp pre init. */ void set_tcp_init_handler(tcp_init_handler h) { - m_tcp_init_handler = h; + set_tcp_pre_init_handler(h); + } + + /// Sets the tcp post init handler + /** + * The tcp post init handler is called after the tcp connection has been + * established and all additional wrappers (proxy connects, TLS handshakes, + * etc have been performed. This is fired before any bytes are read or any + * WebSocket specific handshake logic has been performed. + * + * @since 0.4.0-alpha1 + * + * @param h The handler to call on tcp post init. + */ + void set_tcp_post_init_handler(tcp_init_handler h) { + m_tcp_post_init_handler = h; } /// Set the proxy to connect through (exception free) @@ -264,15 +296,21 @@ public: ) ); - new_timer->async_wait( - m_strand->wrap(lib::bind( - &type::handle_timer, - get_shared(), + if (config::enable_multithreading) { + new_timer->async_wait(m_strand->wrap(lib::bind( + &type::handle_timer, get_shared(), new_timer, callback, lib::placeholders::_1 - )) - ); + ))); + } else { + new_timer->async_wait(lib::bind( + &type::handle_timer, get_shared(), + new_timer, + callback, + lib::placeholders::_1 + )); + } return new_timer; } @@ -288,8 +326,8 @@ public: * @param callback The function to call back * @param ec The status code */ - void handle_timer(timer_ptr t, timer_handler callback, const - boost::system::error_code& ec) + void handle_timer(timer_ptr t, timer_handler callback, + boost::system::error_code const & ec) { if (ec) { if (ec == boost::asio::error::operation_aborted) { @@ -333,11 +371,12 @@ protected: // TODO: pre-init timeout. Right now no implemented socket policies // actually have an asyncronous pre-init + m_init_handler = callback; + socket_con_type::pre_init( lib::bind( &type::handle_pre_init, get_shared(), - callback, lib::placeholders::_1 ) ); @@ -379,34 +418,57 @@ protected: // do we need to store or use the io_service at this level? m_io_service = io_service; - m_strand.reset(new boost::asio::strand(*io_service)); + if (config::enable_multithreading) { + m_strand.reset(new boost::asio::strand(*io_service)); + + m_async_read_handler = m_strand->wrap(lib::bind( + &type::handle_async_read, get_shared(), + lib::placeholders::_1, lib::placeholders::_2 + )); + + m_async_write_handler = m_strand->wrap(lib::bind( + &type::handle_async_write, get_shared(), + lib::placeholders::_1, lib::placeholders::_2 + )); + } else { + // TODO: goal: not have this line here + //m_strand.reset(new boost::asio::strand(*io_service)); + + m_async_read_handler = lib::bind( + &type::handle_async_read, get_shared(), + lib::placeholders::_1, lib::placeholders::_2 + ); + + m_async_write_handler = lib::bind(&type::handle_async_write, + get_shared(), lib::placeholders::_1, lib::placeholders::_2); + } return socket_con_type::init_asio(io_service, m_strand, m_is_server); } - void handle_pre_init(init_handler callback, const lib::error_code& ec) { + void handle_pre_init(lib::error_code const & ec) { if (m_alog.static_test(log::alevel::devel)) { m_alog.write(log::alevel::devel,"asio connection handle pre_init"); } - if (m_tcp_init_handler) { - m_tcp_init_handler(m_connection_hdl); + if (m_tcp_pre_init_handler) { + m_tcp_pre_init_handler(m_connection_hdl); } if (ec) { - callback(ec); + m_init_handler(ec); } // If we have a proxy set issue a proxy connect, otherwise skip to // post_init if (!m_proxy.empty()) { - proxy_write(callback); + proxy_write(); } else { - post_init(callback); + post_init(); } } - void post_init(init_handler callback) { + void post_init() { if (m_alog.static_test(log::alevel::devel)) { m_alog.write(log::alevel::devel,"asio connection post_init"); } @@ -418,7 +480,7 @@ protected: &type::handle_post_init_timeout, get_shared(), post_timer, - callback, + m_init_handler, lib::placeholders::_1 ) ); @@ -428,7 +490,7 @@ protected: &type::handle_post_init, get_shared(), post_timer, - callback, + m_init_handler, lib::placeholders::_1 ) ); @@ -477,10 +539,14 @@ protected: m_alog.write(log::alevel::devel,"asio connection handle_post_init"); } + if (m_tcp_post_init_handler) { + m_tcp_post_init_handler(m_connection_hdl); + } + callback(ec); } - void proxy_write(init_handler callback) { + void proxy_write() { if (m_alog.static_test(log::alevel::devel)) { m_alog.write(log::alevel::devel,"asio connection proxy_write"); } @@ -488,7 +554,7 @@ protected: if (!m_proxy_data) { m_elog.write(log::elevel::library, "assertion failed: !m_proxy_data in asio::connection::proxy_write"); - callback(make_error_code(error::general)); + m_init_handler(make_error_code(error::general)); return; } @@ -505,25 +571,37 @@ protected: lib::bind( &type::handle_proxy_timeout, get_shared(), - callback, + m_init_handler, lib::placeholders::_1 ) ); // Send proxy request - boost::asio::async_write( - socket_con_type::get_next_layer(), - m_bufs, - m_strand->wrap(lib::bind( - &type::handle_proxy_write, - get_shared(), - callback, - lib::placeholders::_1 - )) - ); + if (config::enable_multithreading) { + boost::asio::async_write( + socket_con_type::get_next_layer(), + m_bufs, + m_strand->wrap(lib::bind( + &type::handle_proxy_write, get_shared(), + m_init_handler, + lib::placeholders::_1 + )) + ); + } else { + boost::asio::async_write( + socket_con_type::get_next_layer(), + m_bufs, + lib::bind( + &type::handle_proxy_write, get_shared(), + m_init_handler, + lib::placeholders::_1 + ) + ); + } } - void handle_proxy_timeout(init_handler callback, const lib::error_code & ec) { + void handle_proxy_timeout(init_handler callback, lib::error_code const & ec) + { if (ec == transport::error::operation_aborted) { m_alog.write(log::alevel::devel, "asio handle_proxy_write timer cancelled"); @@ -539,11 +617,12 @@ protected: } } - void handle_proxy_write(init_handler callback, const - boost::system::error_code& ec) + void handle_proxy_write(init_handler callback, + boost::system::error_code const & ec) { if (m_alog.static_test(log::alevel::devel)) { - m_alog.write(log::alevel::devel,"asio connection handle_proxy_write"); + m_alog.write(log::alevel::devel, + "asio connection handle_proxy_write"); } m_bufs.clear(); @@ -581,25 +660,37 @@ protected: return; } - boost::asio::async_read_until( - socket_con_type::get_next_layer(), - m_proxy_data->read_buf, - "\r\n\r\n", - m_strand->wrap(lib::bind( - &type::handle_proxy_read, - get_shared(), - callback, - lib::placeholders::_1, - lib::placeholders::_2 - )) - ); + if (config::enable_multithreading) { + boost::asio::async_read_until( + socket_con_type::get_next_layer(), + m_proxy_data->read_buf, + "\r\n\r\n", + m_strand->wrap(lib::bind( + &type::handle_proxy_read, get_shared(), + callback, + lib::placeholders::_1, lib::placeholders::_2 + )) + ); + } else { + boost::asio::async_read_until( + socket_con_type::get_next_layer(), + m_proxy_data->read_buf, + "\r\n\r\n", + lib::bind( + &type::handle_proxy_read, get_shared(), + callback, + lib::placeholders::_1, lib::placeholders::_2 + ) + ); + } } - void handle_proxy_read(init_handler callback, const - boost::system::error_code& ec, size_t bytes_transferred) + void handle_proxy_read(init_handler callback, + boost::system::error_code const & ec, size_t bytes_transferred) { if (m_alog.static_test(log::alevel::devel)) { - m_alog.write(log::alevel::devel,"asio connection handle_proxy_read"); + m_alog.write(log::alevel::devel, + "asio connection handle_proxy_read"); } // Timer expired or the operation was aborted for some reason. @@ -667,7 +758,7 @@ protected: m_proxy_data.reset(); // Continue with post proxy initialization - post_init(callback); + post_init(); } } @@ -685,93 +776,114 @@ protected: m_alog.write(log::alevel::devel,s.str()); } - if (num_bytes > len) { + if (!m_async_read_handler) { + m_alog.write(log::alevel::devel, + "async_read_at_least called after async_shutdown"); + handler(make_error_code(transport::error::action_after_shutdown), 0); + return; + } + + // TODO: safety vs speed ? + // maybe move into an if devel block + /*if (num_bytes > len) { m_elog.write(log::elevel::devel, "asio async_read_at_least error::invalid_num_bytes"); handler(make_error_code(transport::error::invalid_num_bytes), size_t(0)); return; - } + }*/ + + m_read_handler = handler; boost::asio::async_read( socket_con_type::get_socket(), boost::asio::buffer(buf,len), boost::asio::transfer_at_least(num_bytes), - m_strand->wrap(lib::bind( - &type::handle_async_read, - get_shared(), - handler, - lib::placeholders::_1, - lib::placeholders::_2 - )) + make_custom_alloc_handler( + m_read_handler_allocator, + m_async_read_handler + ) ); } - void handle_async_read(read_handler handler, const - boost::system::error_code& ec, size_t bytes_transferred) + void handle_async_read(const boost::system::error_code& ec, + size_t bytes_transferred) { if (!ec) { - handler(lib::error_code(), bytes_transferred); + m_read_handler(lib::error_code(), bytes_transferred); return; } // translate boost error codes into more lib::error_codes if (ec == boost::asio::error::eof) { - handler(make_error_code(transport::error::eof), + m_read_handler(make_error_code(transport::error::eof), bytes_transferred); } else if (ec.value() == 335544539) { - handler(make_error_code(transport::error::tls_short_read), + m_read_handler(make_error_code(transport::error::tls_short_read), bytes_transferred); } else { log_err(log::elevel::info,"asio async_read_at_least",ec); - handler(make_error_code(transport::error::pass_through), + m_read_handler(make_error_code(transport::error::pass_through), bytes_transferred); } } void async_write(const char* buf, size_t len, write_handler handler) { - m_bufs.push_back(boost::asio::buffer(buf,len)); + if (!m_async_write_handler) { + m_alog.write(log::alevel::devel, + "async_write (single) called after async_shutdown"); + handler(make_error_code(transport::error::action_after_shutdown)); + return; + } + + m_bufs.push_back(boost::asio::buffer(buf,len)); + + m_write_handler = handler; boost::asio::async_write( socket_con_type::get_socket(), m_bufs, - m_strand->wrap(lib::bind( - &type::handle_async_write, - get_shared(), - handler, - lib::placeholders::_1 - )) + make_custom_alloc_handler( + m_write_handler_allocator, + m_async_write_handler + ) ); } void async_write(const std::vector& bufs, write_handler handler) { + if (!m_async_write_handler) { + m_alog.write(log::alevel::devel, + "async_write (vector) called after async_shutdown"); + handler(make_error_code(transport::error::action_after_shutdown)); + return; + } std::vector::const_iterator it; for (it = bufs.begin(); it != bufs.end(); ++it) { m_bufs.push_back(boost::asio::buffer((*it).buf,(*it).len)); } + m_write_handler = handler; + boost::asio::async_write( socket_con_type::get_socket(), m_bufs, - m_strand->wrap(lib::bind( - &type::handle_async_write, - get_shared(), - handler, - lib::placeholders::_1 - )) + make_custom_alloc_handler( + m_write_handler_allocator, + m_async_write_handler + ) ); } - void handle_async_write(write_handler handler, const - boost::system::error_code& ec) + void handle_async_write(boost::system::error_code const & ec, + size_t bytes_transferred) { m_bufs.clear(); if (ec) { log_err(log::elevel::info,"asio async_write",ec); - handler(make_error_code(transport::error::pass_through)); + m_write_handler(make_error_code(transport::error::pass_through)); } else { - handler(lib::error_code()); + m_write_handler(lib::error_code()); } } @@ -792,12 +904,20 @@ protected: * This needs to be thread safe */ lib::error_code interrupt(interrupt_handler handler) { - m_io_service->post(m_strand->wrap(handler)); + if (config::enable_multithreading) { + m_io_service->post(m_strand->wrap(handler)); + } else { + m_io_service->post(handler); + } return lib::error_code(); } lib::error_code dispatch(dispatch_handler handler) { - m_io_service->post(m_strand->wrap(handler)); + if (config::enable_multithreading) { + m_io_service->post(m_strand->wrap(handler)); + } else { + m_io_service->post(handler); + } return lib::error_code(); } @@ -811,6 +931,13 @@ protected: m_alog.write(log::alevel::devel,"asio connection async_shutdown"); } + // Reset cached handlers now that we won't be reading or writing anymore + // These cached handlers store shared pointers to this connection and will leak + // the connection if not destroyed. + m_async_read_handler = 0; + m_async_write_handler = 0; + m_init_handler = 0; + timer_ptr shutdown_timer; shutdown_timer = set_timer( config::timeout_socket_shutdown, @@ -926,7 +1053,18 @@ private: std::vector m_bufs; // Handlers - tcp_init_handler m_tcp_init_handler; + tcp_init_handler m_tcp_pre_init_handler; + tcp_init_handler m_tcp_post_init_handler; + + handler_allocator m_read_handler_allocator; + handler_allocator m_write_handler_allocator; + + read_handler m_read_handler; + write_handler m_write_handler; + init_handler m_init_handler; + + async_read_handler m_async_read_handler; + async_write_handler m_async_write_handler; }; diff --git a/src/third_party/websocketpp/websocketpp/transport/asio/endpoint.hpp b/src/third_party/websocketpp/websocketpp/transport/asio/endpoint.hpp index e91dc58a..2728c66b 100644 --- a/src/third_party/websocketpp/websocketpp/transport/asio/endpoint.hpp +++ b/src/third_party/websocketpp/websocketpp/transport/asio/endpoint.hpp @@ -84,10 +84,13 @@ public: typedef lib::shared_ptr resolver_ptr; /// Type of timer handle typedef lib::shared_ptr timer_ptr; + /// Type of a shared pointer to an io_service work object + typedef lib::shared_ptr work_ptr; // generate and manage our own io_service explicit endpoint() : m_external_io_service(false) + , m_listen_backlog(0) , m_state(UNINITIALIZED) { //std::cout << "transport::asio::endpoint constructor" << std::endl; @@ -119,6 +122,7 @@ public: : m_io_service(src.m_io_service) , m_external_io_service(src.m_external_io_service) , m_acceptor(src.m_acceptor) + , m_listen_backlog(0) , m_state(src.m_state) { src.m_io_service = NULL; @@ -132,11 +136,13 @@ public: m_io_service = rhs.m_io_service; m_external_io_service = rhs.m_external_io_service; m_acceptor = rhs.m_acceptor; + m_listen_backlog = rhs.m_listen_backlog m_state = rhs.m_state; rhs.m_io_service = NULL; rhs.m_external_io_service = false; rhs.m_acceptor = NULL; + rhs.m_listen_backlog = 0; rhs.m_state = UNINITIALIZED; } return *this; @@ -216,6 +222,72 @@ public: m_external_io_service = false; } + /// Sets the tcp pre init handler + /** + * The tcp pre init handler is called after the raw tcp connection has been + * established but before any additional wrappers (proxy connects, TLS + * handshakes, etc) have been performed. + * + * @since 0.4.0-alpha1 + * + * @param h The handler to call on tcp pre init. + */ + void set_tcp_pre_init_handler(tcp_init_handler h) { + m_tcp_pre_init_handler = h; + } + + /// Sets the tcp pre init handler (deprecated) + /** + * The tcp pre init handler is called after the raw tcp connection has been + * established but before any additional wrappers (proxy connects, TLS + * handshakes, etc) have been performed. + * + * @deprecated Use set_tcp_pre_init_handler instead + * + * @param h The handler to call on tcp pre init. + */ + void set_tcp_init_handler(tcp_init_handler h) { + set_tcp_pre_init_handler(h); + } + + /// Sets the tcp post init handler + /** + * The tcp post init handler is called after the tcp connection has been + * established and all additional wrappers (proxy connects, TLS handshakes, + * etc have been performed. This is fired before any bytes are read or any + * WebSocket specific handshake logic has been performed. + * + * @since 0.4.0-alpha1 + * + * @param h The handler to call on tcp post init. + */ + void set_tcp_post_init_handler(tcp_init_handler h) { + m_tcp_post_init_handler = h; + } + + /// Sets the maximum length of the queue of pending connections. + /** + * Sets the maximum length of the queue of pending connections. Increasing + * this will allow WebSocket++ to queue additional incoming connections. + * Setting it higher may prevent failed connections at high connection rates + * but may cause additional latency. + * + * For this value to take effect you may need to adjust operating system + * settings. + * + * New values affect future calls to listen only. + * + * A value of zero will use the operating system default. This is the + * default value. + * + * @since 0.4.0-alpha1 + * + * @param backlog The maximum length of the queue of pending connections + */ + void set_listen_backlog(int backlog) { + m_listen_backlog = backlog; + } + /// Retrieve a reference to the endpoint's io_service /** * The io_service may be an internal or external one. This may be used to @@ -231,20 +303,6 @@ public: return *m_io_service; } - /// Sets the tcp init handler - /** - * The tcp init handler is called after the tcp connection has been - * established. - * - * @see WebSocket++ handler documentation for more information about - * handlers. - * - * @param h The handler to call on tcp init. - */ - void set_tcp_init_handler(tcp_init_handler h) { - m_tcp_init_handler = h; - } - /// Set up endpoint for listening manually (exception free) /** * Bind the internal acceptor using the specified settings. The endpoint @@ -268,7 +326,11 @@ public: m_acceptor->open(ep.protocol()); m_acceptor->set_option(boost::asio::socket_base::reuse_address(true)); m_acceptor->bind(ep); - m_acceptor->listen(); + if (m_listen_backlog == 0) { + m_acceptor->listen(); + } else { + m_acceptor->listen(m_listen_backlog); + } m_state = LISTENING; ec = lib::error_code(); } @@ -494,6 +556,34 @@ public: return m_io_service->stopped(); } + /// Marks the endpoint as perpetual, stopping it from exiting when empty + /** + * Marks the endpoint as perpetual. Perpetual endpoints will not + * automatically exit when they run out of connections to process. To stop + * a perpetual endpoint call `end_perpetual`. + * + * An endpoint may be marked perpetual at any time by any thread. It must be + * called either before the endpoint has run out of work or before it was + * started + * + * @since 0.4.0-alpha1 + */ + void start_perpetual() { + m_work.reset(new boost::asio::io_service::work(*m_io_service)); + } + + /// Clears the endpoint's perpetual flag, allowing it to exit when empty + /** + * Clears the endpoint's perpetual flag. This will cause the endpoint's run + * method to exit normally when it runs out of connections. If there are + * currently active connections it will not end until they are complete. + * + * @since 0.4.0-alpha1 + */ + void stop_perpetual() { + m_work.reset(); + } + /// Call back a function after a period of time. /** * Sets a timer that calls back a function after the specified period of @@ -572,15 +662,27 @@ public: m_alog->write(log::alevel::devel, "asio::async_accept"); - m_acceptor->async_accept( - tcon->get_raw_socket(), - tcon->get_strand()->wrap(lib::bind( - &type::handle_accept, - this, - callback, - lib::placeholders::_1 - )) - ); + if (config::enable_multithreading) { + m_acceptor->async_accept( + tcon->get_raw_socket(), + tcon->get_strand()->wrap(lib::bind( + &type::handle_accept, + this, + callback, + lib::placeholders::_1 + )) + ); + } else { + m_acceptor->async_accept( + tcon->get_raw_socket(), + lib::bind( + &type::handle_accept, + this, + callback, + lib::placeholders::_1 + ) + ); + } } /// Accept the next connection attempt and assign it to con. @@ -683,18 +785,33 @@ protected: ) ); - m_resolver->async_resolve( - query, - tcon->get_strand()->wrap(lib::bind( - &type::handle_resolve, - this, - tcon, - dns_timer, - cb, - lib::placeholders::_1, - lib::placeholders::_2 - )) - ); + if (config::enable_multithreading) { + m_resolver->async_resolve( + query, + tcon->get_strand()->wrap(lib::bind( + &type::handle_resolve, + this, + tcon, + dns_timer, + cb, + lib::placeholders::_1, + lib::placeholders::_2 + )) + ); + } else { + m_resolver->async_resolve( + query, + lib::bind( + &type::handle_resolve, + this, + tcon, + dns_timer, + cb, + lib::placeholders::_1, + lib::placeholders::_2 + ) + ); + } } void handle_resolve_timeout(timer_ptr dns_timer, connect_handler callback, @@ -767,18 +884,33 @@ protected: ) ); - boost::asio::async_connect( - tcon->get_raw_socket(), - iterator, - tcon->get_strand()->wrap(lib::bind( - &type::handle_connect, - this, - tcon, - con_timer, - callback, - lib::placeholders::_1 - )) - ); + if (config::enable_multithreading) { + boost::asio::async_connect( + tcon->get_raw_socket(), + iterator, + tcon->get_strand()->wrap(lib::bind( + &type::handle_connect, + this, + tcon, + con_timer, + callback, + lib::placeholders::_1 + )) + ); + } else { + boost::asio::async_connect( + tcon->get_raw_socket(), + iterator, + lib::bind( + &type::handle_connect, + this, + tcon, + con_timer, + callback, + lib::placeholders::_1 + ) + ); + } } void handle_connect_timeout(transport_con_ptr tcon, timer_ptr con_timer, @@ -857,7 +989,8 @@ protected: ec = tcon->init_asio(m_io_service); if (ec) {return ec;} - tcon->set_tcp_init_handler(m_tcp_init_handler); + tcon->set_tcp_pre_init_handler(m_tcp_pre_init_handler); + tcon->set_tcp_post_init_handler(m_tcp_post_init_handler); return lib::error_code(); } @@ -877,13 +1010,18 @@ private: }; // Handlers - tcp_init_handler m_tcp_init_handler; + tcp_init_handler m_tcp_pre_init_handler; + tcp_init_handler m_tcp_post_init_handler; // Network Resources io_service_ptr m_io_service; bool m_external_io_service; acceptor_ptr m_acceptor; resolver_ptr m_resolver; + work_ptr m_work; + + // Network constants + int m_listen_backlog; elog_type* m_elog; alog_type* m_alog; diff --git a/src/third_party/websocketpp/websocketpp/transport/asio/security/tls.hpp b/src/third_party/websocketpp/websocketpp/transport/asio/security/tls.hpp index d7da561e..2a743112 100644 --- a/src/third_party/websocketpp/websocketpp/transport/asio/security/tls.hpp +++ b/src/third_party/websocketpp/websocketpp/transport/asio/security/tls.hpp @@ -228,15 +228,25 @@ protected: m_ec = socket::make_error_code(socket::error::tls_handshake_timeout); // TLS handshake - m_socket->async_handshake( - get_handshake_type(), - m_strand->wrap(lib::bind( - &type::handle_init, - get_shared(), - callback, - lib::placeholders::_1 - )) - ); + if (m_strand) { + m_socket->async_handshake( + get_handshake_type(), + m_strand->wrap(lib::bind( + &type::handle_init, get_shared(), + callback, + lib::placeholders::_1 + )) + ); + } else { + m_socket->async_handshake( + get_handshake_type(), + lib::bind( + &type::handle_init, get_shared(), + callback, + lib::placeholders::_1 + ) + ); + } } /// Sets the connection handle diff --git a/src/third_party/websocketpp/websocketpp/transport/base/connection.hpp b/src/third_party/websocketpp/websocketpp/transport/base/connection.hpp index 1db6310d..70e1ec50 100644 --- a/src/third_party/websocketpp/websocketpp/transport/base/connection.hpp +++ b/src/third_party/websocketpp/websocketpp/transport/base/connection.hpp @@ -169,7 +169,10 @@ enum value { tls_short_read, /// Timer expired - timeout + timeout, + + /// read or write after shutdown + action_after_shutdown }; class category : public lib::error_category { @@ -198,6 +201,8 @@ class category : public lib::error_category { return "TLS Short Read"; case timeout: return "Timer Expired"; + case action_after_shutdown: + return "A transport action was requested after shutdown"; default: return "Unknown"; } diff --git a/src/third_party/websocketpp/websocketpp/uri.hpp b/src/third_party/websocketpp/websocketpp/uri.hpp index b0c7f38f..5471031a 100644 --- a/src/third_party/websocketpp/websocketpp/uri.hpp +++ b/src/third_party/websocketpp/websocketpp/uri.hpp @@ -282,7 +282,22 @@ public: return s.str(); } - // get query? + /// Return the query portion + /** + * Returns the query portion (after the ?) of the URI or an empty string if + * there is none. + * + * @return query portion of the URI. + */ + std::string get_query() const { + std::size_t found = m_resource.find('?'); + if (found != std::string::npos) { + return m_resource.substr(found + 1); + } else { + return ""; + } + } + // get fragment // hi <3