17 : server_id_(server_id)
18 , stop_future_(stop_promise_.get_future())
36 "Server already running",
42 cleanup_timer_.reset();
45 io_context_ = std::make_unique<asio::io_context>();
46 work_guard_ = std::make_unique<asio::executor_work_guard<asio::io_context::executor_type>>(
47 io_context_->get_executor());
49 acceptor_ = std::make_unique<asio::ip::tcp::acceptor>(
51 asio::ip::tcp::endpoint(asio::ip::tcp::v4(), port));
57 io_future_ = std::async(std::launch::async, [
this]() { run_io(); });
61 start_cleanup_timer();
64 }
catch (
const std::exception& e) {
67 std::string(
"Failed to start server: ") + e.what(),
77 "Server already running",
83 cleanup_timer_.reset();
86 io_context_ = std::make_unique<asio::io_context>();
87 work_guard_ = std::make_unique<asio::executor_work_guard<asio::io_context::executor_type>>(
88 io_context_->get_executor());
91 ssl_context_ = std::make_unique<asio::ssl::context>(asio::ssl::context::tls_server);
92 ssl_context_->set_options(
93 asio::ssl::context::default_workarounds |
94 asio::ssl::context::no_sslv2 |
95 asio::ssl::context::no_sslv3 |
96 asio::ssl::context::no_tlsv1 |
97 asio::ssl::context::no_tlsv1_1);
99 ssl_context_->use_certificate_file(
config.cert_file, asio::ssl::context::pem);
100 ssl_context_->use_private_key_file(
config.key_file, asio::ssl::context::pem);
102 if (!
config.ca_file.empty()) {
103 ssl_context_->load_verify_file(
config.ca_file);
106 if (
config.verify_client) {
107 ssl_context_->set_verify_mode(asio::ssl::verify_peer | asio::ssl::verify_fail_if_no_peer_cert);
111 SSL_CTX_set_alpn_select_cb(
112 ssl_context_->native_handle(),
113 [](
SSL* ,
const unsigned char** out,
unsigned char* outlen,
114 const unsigned char* in,
unsigned int inlen,
void* ) ->
int {
116 const unsigned char* client = in;
117 const unsigned char* end = in + inlen;
118 while (client < end) {
119 unsigned char len = *client++;
120 if (len == 2 && client + 2 <= end &&
121 client[0] ==
'h' && client[1] ==
'2') {
124 return SSL_TLSEXT_ERR_OK;
128 return SSL_TLSEXT_ERR_NOACK;
132 acceptor_ = std::make_unique<asio::ip::tcp::acceptor>(
134 asio::ip::tcp::endpoint(asio::ip::tcp::v4(), port));
140 io_future_ = std::async(std::launch::async, [
this]() { run_io(); });
144 start_cleanup_timer();
147 }
catch (
const std::exception& e) {
150 std::string(
"Failed to start TLS server: ") + e.what(),
169 std::lock_guard<std::mutex> lock(connections_mutex_);
170 for (
auto& [
id, conn] : connections_) {
171 conn->shutdown_transport();
180 if (acceptor_ && acceptor_->is_open()) {
182 acceptor_->close(ec);
187 std::lock_guard<std::mutex> lock(connections_mutex_);
188 for (
auto& [
id, conn] : connections_) {
191 connections_.clear();
195 if (cleanup_timer_) {
196 cleanup_timer_->cancel();
201 stop_promise_.set_value();
209 auto http2_server::is_running() const ->
bool
214 auto http2_server::wait() ->
void
221 request_handler_ = std::move(handler);
226 error_handler_ = std::move(handler);
231 std::lock_guard<std::mutex> lock(settings_mutex_);
233 encoder_->set_max_table_size(
settings.header_table_size);
234 decoder_->set_max_table_size(
settings.header_table_size);
239 std::lock_guard<std::mutex> lock(settings_mutex_);
243 auto http2_server::active_connections() const ->
size_t
245 std::lock_guard<std::mutex> lock(
const_cast<std::mutex&
>(connections_mutex_));
246 return connections_.size();
249 auto http2_server::active_streams() const ->
size_t
251 std::lock_guard<std::mutex> lock(
const_cast<std::mutex&
>(connections_mutex_));
253 for (
const auto& [
id, conn] : connections_) {
254 total += conn->stream_count();
259 auto http2_server::server_id() const -> std::string_view
264 auto http2_server::do_accept() ->
void
266 if (!is_running_ || !acceptor_) {
270 acceptor_->async_accept(
271 [
this](std::error_code ec, asio::ip::tcp::socket socket) {
272 handle_accept(ec, std::move(socket));
276 auto http2_server::do_accept_tls() ->
void
278 if (!is_running_ || !acceptor_) {
282 acceptor_->async_accept(
283 [
this](std::error_code ec, asio::ip::tcp::socket socket) {
284 handle_accept_tls(ec, std::move(socket));
288 auto http2_server::handle_accept(std::error_code ec, asio::ip::tcp::socket socket) ->
void
295 if (error_handler_) {
296 error_handler_(std::string(
"Accept error: ") + ec.message());
299 uint64_t conn_id = next_connection_id_++;
300 auto conn = std::make_shared<http2_server_connection>(
307 add_connection(conn);
315 auto http2_server::handle_accept_tls(std::error_code ec, asio::ip::tcp::socket socket) ->
void
322 if (error_handler_) {
323 error_handler_(std::string(
"Accept error: ") + ec.message());
329 auto tls_socket = std::make_unique<asio::ssl::stream<asio::ip::tcp::socket>>(
330 std::move(socket), *ssl_context_);
332 auto* raw_socket = tls_socket.get();
333 raw_socket->async_handshake(
334 asio::ssl::stream_base::server,
335 [
this, tls_socket = std::move(tls_socket)](std::error_code ec)
mutable {
337 if (error_handler_) {
338 error_handler_(std::string(
"TLS handshake error: ") + ec.message());
341 uint64_t conn_id = next_connection_id_++;
342 auto conn = std::make_shared<http2_server_connection>(
344 std::move(tls_socket),
349 add_connection(conn);
358 auto http2_server::add_connection(std::shared_ptr<http2_server_connection> conn) ->
void
360 std::lock_guard<std::mutex> lock(connections_mutex_);
361 connections_[conn->connection_id()] = std::move(conn);
364 auto http2_server::remove_connection(uint64_t connection_id) ->
void
366 std::lock_guard<std::mutex> lock(connections_mutex_);
367 connections_.erase(connection_id);
370 auto http2_server::cleanup_dead_connections() ->
void
372 std::lock_guard<std::mutex> lock(connections_mutex_);
373 for (
auto it = connections_.begin(); it != connections_.end();) {
374 if (!it->second->is_alive()) {
375 it = connections_.erase(it);
382 auto http2_server::start_cleanup_timer() ->
void
388 cleanup_timer_ = std::make_unique<asio::steady_timer>(*io_context_);
389 cleanup_timer_->expires_after(std::chrono::seconds(30));
390 cleanup_timer_->async_wait([
this](std::error_code ec) {
391 if (!ec && is_running_) {
392 cleanup_dead_connections();
393 start_cleanup_timer();
398 auto http2_server::run_io() ->
void
405 auto http2_server::stop_io() ->
void
415 if (io_future_.valid()) {
424 http2_server_connection::http2_server_connection(
425 uint64_t connection_id,
426 asio::ip::tcp::socket socket,
430 : connection_id_(connection_id)
432 , plain_socket_(std::make_unique<asio::ip::tcp::socket>(std::move(socket)))
436 , request_handler_(std::move(request_handler))
437 , error_handler_(std::move(error_handler))
438 , frame_header_buffer_{}
443 uint64_t connection_id,
444 std::unique_ptr<asio::ssl::stream<asio::ip::tcp::socket>> socket,
448 : connection_id_(connection_id)
450 , tls_socket_(std::move(socket))
454 , request_handler_(std::move(request_handler))
455 , error_handler_(std::move(error_handler))
456 , frame_header_buffer_{}
468 read_connection_preface();
474 std::lock_guard<std::mutex> lock(transport_shutdown_mutex_);
475 if (!is_alive_.exchange(
false)) {
481 if (use_tls_ && tls_socket_) {
482 tls_socket_->lowest_layer().close(ec);
483 }
else if (plain_socket_) {
484 plain_socket_->close(ec);
495 std::lock_guard<std::mutex> lock(transport_shutdown_mutex_);
497 if (use_tls_ && tls_socket_) {
498 tls_socket_->lowest_layer().shutdown(asio::ip::tcp::socket::shutdown_both, ec);
499 }
else if (plain_socket_) {
500 plain_socket_->shutdown(asio::ip::tcp::socket::shutdown_both, ec);
516 std::lock_guard<std::mutex> lock(
const_cast<std::mutex&
>(
streams_mutex_));
522 constexpr size_t PREFACE_SIZE = 24;
523 auto buffer = std::make_shared<std::vector<uint8_t>>(PREFACE_SIZE);
525 auto read_handler = [
this, self = shared_from_this(), buffer](std::error_code ec, std::size_t bytes_read) {
526 if (ec || bytes_read != 24) {
527 if (error_handler_) {
528 error_handler_(
"Failed to read connection preface");
535 const std::string expected =
"PRI * HTTP/2.0\r\n\r\nSM\r\n\r\n";
536 if (std::string(buffer->begin(), buffer->end()) != expected) {
537 if (error_handler_) {
538 error_handler_(
"Invalid connection preface");
544 preface_received_ =
true;
547 auto result = send_settings();
548 if (result.is_err()) {
549 if (error_handler_) {
550 error_handler_(
"Failed to send settings");
561 asio::async_read(*tls_socket_, asio::buffer(*buffer), read_handler);
563 asio::async_read(*plain_socket_, asio::buffer(*buffer), read_handler);
569 std::vector<setting_parameter> params = {
571 local_settings_.max_concurrent_streams},
573 local_settings_.initial_window_size},
575 local_settings_.max_frame_size},
577 local_settings_.header_table_size},
581 return send_frame(
frame);
586 if (
frame.is_ack()) {
592 for (
const auto& param :
frame.settings()) {
595 remote_settings_.header_table_size = param.value;
596 encoder_.set_max_table_size(param.value);
599 remote_settings_.enable_push = (param.value != 0);
602 remote_settings_.max_concurrent_streams = param.value;
605 remote_settings_.initial_window_size = param.value;
608 remote_settings_.max_frame_size = param.value;
611 remote_settings_.max_header_list_size = param.value;
617 return send_settings_ack();
623 return send_frame(
frame);
637 auto read_handler = [
this, self = shared_from_this()](std::error_code ec, std::size_t bytes_read) {
638 if (ec || bytes_read != 9) {
639 if (ec != asio::error::eof && ec != asio::error::operation_aborted) {
640 if (error_handler_) {
641 error_handler_(std::string(
"Read error: ") + ec.message());
649 if (header_result.is_err()) {
650 if (error_handler_) {
651 error_handler_(
"Failed to parse frame header");
657 auto header = header_result.value();
658 if (header.length > 0) {
659 read_frame_payload(header.length);
663 if (frame_result.is_ok()) {
664 process_frame(std::move(frame_result.value()));
671 asio::async_read(*tls_socket_, asio::buffer(frame_header_buffer_), read_handler);
673 asio::async_read(*plain_socket_, asio::buffer(frame_header_buffer_), read_handler);
683 read_buffer_.resize(9 + length);
684 std::copy(frame_header_buffer_.begin(), frame_header_buffer_.end(), read_buffer_.begin());
686 auto payload_buffer = asio::buffer(read_buffer_.data() + 9, length);
688 auto read_handler = [
this, self = shared_from_this()](std::error_code ec, std::size_t ) {
690 if (ec != asio::error::eof && ec != asio::error::operation_aborted) {
691 if (error_handler_) {
692 error_handler_(std::string(
"Read payload error: ") + ec.message());
700 if (frame_result.is_ok()) {
701 process_frame(std::move(frame_result.value()));
703 if (error_handler_) {
704 error_handler_(
"Failed to parse frame");
712 asio::async_read(*tls_socket_, payload_buffer, read_handler);
714 asio::async_read(*plain_socket_, payload_buffer, read_handler);
720 auto data = f.serialize();
724 asio::write(*tls_socket_, asio::buffer(
data), ec);
726 asio::write(*plain_socket_, asio::buffer(
data), ec);
732 std::string(
"Failed to send frame: ") + ec.message(),
733 "http2_server_connection");
741 auto& header = f->header();
743 switch (header.type) {
747 return handle_settings_frame(*settings_f);
754 return handle_headers_frame(*headers_f);
759 auto* data_f =
dynamic_cast<data_frame*
>(f.get());
761 return handle_data_frame(*data_f);
768 return handle_rst_stream_frame(*rst_f);
773 auto* ping_f =
dynamic_cast<ping_frame*
>(f.get());
775 return handle_ping_frame(*ping_f);
782 return handle_goaway_frame(*goaway_f);
789 return handle_window_update_frame(*wu_f);
803 std::lock_guard<std::mutex> lock(streams_mutex_);
805 auto it = streams_.find(stream_id);
806 if (it != streams_.end()) {
814 stream.
window_size =
static_cast<int32_t
>(local_settings_.initial_window_size);
816 auto [iter, inserted] = streams_.emplace(stream_id, std::move(stream));
817 if (stream_id > last_stream_id_) {
818 last_stream_id_ = stream_id;
821 return &iter->second;
826 std::lock_guard<std::mutex> lock(streams_mutex_);
827 streams_.erase(stream_id);
832 uint32_t stream_id = f.header().stream_id;
833 auto* stream = get_or_create_stream(stream_id);
838 "Failed to create stream",
839 "http2_server_connection");
843 auto header_block = f.header_block();
844 std::vector<uint8_t> block_data(header_block.begin(), header_block.end());
845 auto headers_result = decoder_.decode(block_data);
847 if (headers_result.is_err()) {
854 headers_result.error().code,
855 headers_result.error().message,
856 "http2_server_connection");
859 stream->request_headers = headers_result.value();
861 if (f.is_end_stream()) {
863 dispatch_request(stream_id);
864 }
else if (f.is_end_headers()) {
865 stream->headers_complete =
true;
873 uint32_t stream_id = f.header().stream_id;
875 std::unique_lock<std::mutex> lock(streams_mutex_);
876 auto it = streams_.find(stream_id);
877 if (it == streams_.end()) {
880 return send_frame(rst);
883 auto& stream = it->second;
884 auto data = f.data();
885 stream.request_body.insert(stream.request_body.end(),
data.begin(),
data.end());
888 size_t data_size =
data.size();
896 send_frame(stream_wu);
899 if (f.is_end_stream()) {
901 stream.body_complete =
true;
905 dispatch_request(stream_id);
913 close_stream(f.header().stream_id);
925 return send_frame(ack);
937 uint32_t stream_id = f.header().stream_id;
938 int32_t increment =
static_cast<int32_t
>(f.window_size_increment());
940 if (stream_id == 0) {
942 connection_window_size_ += increment;
945 std::lock_guard<std::mutex> lock(streams_mutex_);
946 auto it = streams_.find(stream_id);
947 if (it != streams_.end()) {
948 it->second.window_size += increment;
957 if (!request_handler_) {
963 std::lock_guard<std::mutex> lock(streams_mutex_);
964 auto it = streams_.find(stream_id);
965 if (it == streams_.end()) {
968 stream = &it->second;
976 auto encoder_ptr = std::make_shared<hpack_encoder>(encoder_);
977 auto weak_this = weak_from_this();
984 auto self = weak_this.lock();
989 "http2_server_stream");
991 return self->send_frame(f);
993 remote_settings_.max_frame_size);
997 request_handler_(server_stream, server_stream.request());
998 }
catch (
const std::exception& e) {
999 if (error_handler_) {
1000 error_handler_(std::string(
"Request handler exception: ") + e.what());
1005 close_stream(stream_id);
DATA frame (RFC 7540 Section 6.1)
Base class for HTTP/2 frames.
static auto parse(std::span< const uint8_t > data) -> Result< std::unique_ptr< frame > >
Parse frame from raw bytes.
GOAWAY frame (RFC 7540 Section 6.8)
HPACK header decoder (RFC 7541)
HPACK header encoder (RFC 7541)
http2_server_connection(uint64_t connection_id, asio::ip::tcp::socket socket, const http2_settings &settings, http2_server::request_handler_t request_handler, http2_server::error_handler_t error_handler)
Construct server connection with plain socket.
auto handle_data_frame(const data_frame &f) -> VoidResult
auto close_stream(uint32_t stream_id) -> void
auto start_reading() -> void
auto get_or_create_stream(uint32_t stream_id) -> http2_stream *
auto read_frame_payload(uint32_t length) -> void
~http2_server_connection()
Destructor.
auto read_connection_preface() -> void
auto start() -> VoidResult
Start connection handling.
auto stream_count() const -> size_t
Get number of active streams.
auto handle_goaway_frame(const goaway_frame &f) -> VoidResult
auto stop() -> VoidResult
Stop connection.
auto connection_id() const -> uint64_t
Get connection identifier.
auto send_settings_ack() -> VoidResult
auto handle_window_update_frame(const window_update_frame &f) -> VoidResult
auto shutdown_transport() -> void
auto process_frame(std::unique_ptr< frame > f) -> VoidResult
auto handle_settings_frame(const settings_frame &frame) -> VoidResult
auto handle_rst_stream_frame(const rst_stream_frame &f) -> VoidResult
auto send_frame(const frame &f) -> VoidResult
auto send_settings() -> VoidResult
auto handle_headers_frame(const headers_frame &f) -> VoidResult
auto read_frame_header() -> void
auto is_alive() const -> bool
Check if connection is alive.
std::atomic< bool > is_alive_
auto dispatch_request(uint32_t stream_id) -> void
uint64_t connection_id_
Connection identifier.
std::map< uint32_t, http2_stream > streams_
auto handle_ping_frame(const ping_frame &f) -> VoidResult
std::mutex streams_mutex_
Server-side HTTP/2 stream for sending responses.
http2_server(std::string_view server_id)
Construct HTTP/2 server.
std::function< void( const std::string &error_message)> error_handler_t
Error handler function type.
~http2_server()
Destructor - stops server gracefully.
auto stop() -> VoidResult
Stop the server.
std::atomic< bool > is_running_
auto start(unsigned short port) -> VoidResult
Start HTTP/2 server without TLS (h2c - HTTP/2 cleartext)
auto start_tls(unsigned short port, const tls_config &config) -> VoidResult
Start HTTP/2 server with TLS.
std::function< void( http2_server_stream &stream, const http2_request &request)> request_handler_t
Request handler function type.
PING frame (RFC 7540 Section 6.7)
RST_STREAM frame (RFC 7540 Section 6.4)
SETTINGS frame (RFC 7540 Section 6.5)
WINDOW_UPDATE frame (RFC 7540 Section 6.9)
constexpr int internal_error
constexpr int already_exists
constexpr int bind_failed
constexpr int send_failed
@ compression_error
Compression state not updated.
@ stream_closed
Frame received for closed stream.
setting_identifier
SETTINGS frame parameter identifiers (RFC 7540 Section 6.5.2)
@ max_frame_size
SETTINGS_MAX_FRAME_SIZE.
@ max_header_list_size
SETTINGS_MAX_HEADER_LIST_SIZE.
@ max_concurrent_streams
SETTINGS_MAX_CONCURRENT_STREAMS.
@ enable_push
SETTINGS_ENABLE_PUSH.
@ header_table_size
SETTINGS_HEADER_TABLE_SIZE.
@ initial_window_size
SETTINGS_INITIAL_WINDOW_SIZE.
@ open
Stream open and active.
@ half_closed_remote
Remote end closed, local can send.
@ settings
SETTINGS frame.
@ rst_stream
RST_STREAM frame.
@ window_update
WINDOW_UPDATE frame.
VoidResult error_void(int code, const std::string &message, const std::string &source="network_system", const std::string &details="")
static auto from_headers(const std::vector< http_header > &parsed_headers) -> http2_request
Create http2_request from parsed headers.
HTTP/2 connection settings.
HTTP/2 stream state and data.
std::vector< uint8_t > request_body
Request body.
std::atomic< stream_state > state
Current state.
std::vector< http_header > request_headers
Request headers.
uint32_t stream_id
Stream identifier.
int32_t window_size
Flow control window.
TLS configuration for HTTP/2 server.