Network System 0.1.1
High-performance modular networking library for scalable client-server applications
Loading...
Searching...
No Matches
http2_server.cpp
Go to the documentation of this file.
1// BSD 3-Clause License
2// Copyright (c) 2024, 🍀☀🌕🌥 🌊
3// See the LICENSE file in the project root for full license information.
4
6
7#include <algorithm>
8#include <stdexcept>
9
11{
12 // =========================================================================
13 // http2_server implementation
14 // =========================================================================
15
16 http2_server::http2_server(std::string_view server_id)
17 : server_id_(server_id)
18 , stop_future_(stop_promise_.get_future())
19 , encoder_(std::make_shared<hpack_encoder>(settings_.header_table_size))
20 , decoder_(std::make_shared<hpack_decoder>(settings_.header_table_size))
21 {
22 }
23
25 {
26 if (is_running_) {
27 stop();
28 }
29 }
30
31 auto http2_server::start(unsigned short port) -> VoidResult
32 {
33 if (is_running_) {
34 return error_void(
36 "Server already running",
37 "http2_server");
38 }
39
40 try {
41 stop_io();
42 cleanup_timer_.reset();
43 acceptor_.reset();
44 ssl_context_.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());
48
49 acceptor_ = std::make_unique<asio::ip::tcp::acceptor>(
50 *io_context_,
51 asio::ip::tcp::endpoint(asio::ip::tcp::v4(), port));
52
53 use_tls_ = false;
54 is_running_ = true;
55
56 // Start I/O thread
57 io_future_ = std::async(std::launch::async, [this]() { run_io(); });
58
59 // Start accepting connections
60 do_accept();
61 start_cleanup_timer();
62
63 return ok();
64 } catch (const std::exception& e) {
65 return error_void(
67 std::string("Failed to start server: ") + e.what(),
68 "http2_server");
69 }
70 }
71
72 auto http2_server::start_tls(unsigned short port, const tls_config& config) -> VoidResult
73 {
74 if (is_running_) {
75 return error_void(
77 "Server already running",
78 "http2_server");
79 }
80
81 try {
82 stop_io();
83 cleanup_timer_.reset();
84 acceptor_.reset();
85 ssl_context_.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());
89
90 // Setup TLS context
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);
98
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);
101
102 if (!config.ca_file.empty()) {
103 ssl_context_->load_verify_file(config.ca_file);
104 }
105
106 if (config.verify_client) {
107 ssl_context_->set_verify_mode(asio::ssl::verify_peer | asio::ssl::verify_fail_if_no_peer_cert);
108 }
109
110 // Set ALPN callback for HTTP/2
111 SSL_CTX_set_alpn_select_cb(
112 ssl_context_->native_handle(),
113 [](SSL* /*ssl*/, const unsigned char** out, unsigned char* outlen,
114 const unsigned char* in, unsigned int inlen, void* /*arg*/) -> int {
115 // Look for "h2" in client's ALPN list
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') {
122 *out = client;
123 *outlen = 2;
124 return SSL_TLSEXT_ERR_OK;
125 }
126 client += len;
127 }
128 return SSL_TLSEXT_ERR_NOACK;
129 },
130 nullptr);
131
132 acceptor_ = std::make_unique<asio::ip::tcp::acceptor>(
133 *io_context_,
134 asio::ip::tcp::endpoint(asio::ip::tcp::v4(), port));
135
136 use_tls_ = true;
137 is_running_ = true;
138
139 // Start I/O thread
140 io_future_ = std::async(std::launch::async, [this]() { run_io(); });
141
142 // Start accepting connections
143 do_accept_tls();
144 start_cleanup_timer();
145
146 return ok();
147 } catch (const std::exception& e) {
148 return error_void(
150 std::string("Failed to start TLS server: ") + e.what(),
151 "http2_server");
152 }
153 }
154
155 auto http2_server::stop() -> VoidResult
156 {
157 if (!is_running_) {
158 return ok();
159 }
160
161 is_running_ = false;
162
163 // Prevent another callback from starting, and interrupt any response
164 // writer before waiting for the current callback to finish.
165 if (io_context_) {
166 io_context_->stop();
167 }
168 {
169 std::lock_guard<std::mutex> lock(connections_mutex_);
170 for (auto& [id, conn] : connections_) {
171 conn->shutdown_transport();
172 }
173 }
174
175 // A read or request callback can still be using a connection. Join
176 // the I/O worker before closing its sockets or releasing the map.
177 stop_io();
178
179 // Close acceptor
180 if (acceptor_ && acceptor_->is_open()) {
181 std::error_code ec;
182 acceptor_->close(ec);
183 }
184
185 // Stop all connections
186 {
187 std::lock_guard<std::mutex> lock(connections_mutex_);
188 for (auto& [id, conn] : connections_) {
189 conn->stop();
190 }
191 connections_.clear();
192 }
193
194 // Stop cleanup timer
195 if (cleanup_timer_) {
196 cleanup_timer_->cancel();
197 }
198
199 // Signal stop
200 try {
201 stop_promise_.set_value();
202 } catch (...) {
203 // Already set
204 }
205
206 return ok();
207 }
208
209 auto http2_server::is_running() const -> bool
210 {
211 return is_running_;
212 }
213
214 auto http2_server::wait() -> void
215 {
216 stop_future_.wait();
217 }
218
219 auto http2_server::set_request_handler(request_handler_t handler) -> void
220 {
221 request_handler_ = std::move(handler);
222 }
223
224 auto http2_server::set_error_handler(error_handler_t handler) -> void
225 {
226 error_handler_ = std::move(handler);
227 }
228
229 auto http2_server::set_settings(const http2_settings& settings) -> void
230 {
231 std::lock_guard<std::mutex> lock(settings_mutex_);
232 settings_ = settings;
233 encoder_->set_max_table_size(settings.header_table_size);
234 decoder_->set_max_table_size(settings.header_table_size);
235 }
236
237 auto http2_server::get_settings() const -> http2_settings
238 {
239 std::lock_guard<std::mutex> lock(settings_mutex_);
240 return settings_;
241 }
242
243 auto http2_server::active_connections() const -> size_t
244 {
245 std::lock_guard<std::mutex> lock(const_cast<std::mutex&>(connections_mutex_));
246 return connections_.size();
247 }
248
249 auto http2_server::active_streams() const -> size_t
250 {
251 std::lock_guard<std::mutex> lock(const_cast<std::mutex&>(connections_mutex_));
252 size_t total = 0;
253 for (const auto& [id, conn] : connections_) {
254 total += conn->stream_count();
255 }
256 return total;
257 }
258
259 auto http2_server::server_id() const -> std::string_view
260 {
261 return server_id_;
262 }
263
264 auto http2_server::do_accept() -> void
265 {
266 if (!is_running_ || !acceptor_) {
267 return;
268 }
269
270 acceptor_->async_accept(
271 [this](std::error_code ec, asio::ip::tcp::socket socket) {
272 handle_accept(ec, std::move(socket));
273 });
274 }
275
276 auto http2_server::do_accept_tls() -> void
277 {
278 if (!is_running_ || !acceptor_) {
279 return;
280 }
281
282 acceptor_->async_accept(
283 [this](std::error_code ec, asio::ip::tcp::socket socket) {
284 handle_accept_tls(ec, std::move(socket));
285 });
286 }
287
288 auto http2_server::handle_accept(std::error_code ec, asio::ip::tcp::socket socket) -> void
289 {
290 if (!is_running_) {
291 return;
292 }
293
294 if (ec) {
295 if (error_handler_) {
296 error_handler_(std::string("Accept error: ") + ec.message());
297 }
298 } else {
299 uint64_t conn_id = next_connection_id_++;
300 auto conn = std::make_shared<http2_server_connection>(
301 conn_id,
302 std::move(socket),
303 get_settings(),
304 request_handler_,
305 error_handler_);
306
307 add_connection(conn);
308 conn->start();
309 }
310
311 // Continue accepting
312 do_accept();
313 }
314
315 auto http2_server::handle_accept_tls(std::error_code ec, asio::ip::tcp::socket socket) -> void
316 {
317 if (!is_running_) {
318 return;
319 }
320
321 if (ec) {
322 if (error_handler_) {
323 error_handler_(std::string("Accept error: ") + ec.message());
324 }
325 do_accept_tls();
326 return;
327 }
328
329 auto tls_socket = std::make_unique<asio::ssl::stream<asio::ip::tcp::socket>>(
330 std::move(socket), *ssl_context_);
331
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 {
336 if (ec) {
337 if (error_handler_) {
338 error_handler_(std::string("TLS handshake error: ") + ec.message());
339 }
340 } else {
341 uint64_t conn_id = next_connection_id_++;
342 auto conn = std::make_shared<http2_server_connection>(
343 conn_id,
344 std::move(tls_socket),
345 get_settings(),
346 request_handler_,
347 error_handler_);
348
349 add_connection(conn);
350 conn->start();
351 }
352
353 // Continue accepting
354 do_accept_tls();
355 });
356 }
357
358 auto http2_server::add_connection(std::shared_ptr<http2_server_connection> conn) -> void
359 {
360 std::lock_guard<std::mutex> lock(connections_mutex_);
361 connections_[conn->connection_id()] = std::move(conn);
362 }
363
364 auto http2_server::remove_connection(uint64_t connection_id) -> void
365 {
366 std::lock_guard<std::mutex> lock(connections_mutex_);
367 connections_.erase(connection_id);
368 }
369
370 auto http2_server::cleanup_dead_connections() -> void
371 {
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);
376 } else {
377 ++it;
378 }
379 }
380 }
381
382 auto http2_server::start_cleanup_timer() -> void
383 {
384 if (!io_context_) {
385 return;
386 }
387
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();
394 }
395 });
396 }
397
398 auto http2_server::run_io() -> void
399 {
400 if (io_context_) {
401 io_context_->run();
402 }
403 }
404
405 auto http2_server::stop_io() -> void
406 {
407 if (work_guard_) {
408 work_guard_.reset();
409 }
410
411 if (io_context_) {
412 io_context_->stop();
413 }
414
415 if (io_future_.valid()) {
416 io_future_.wait();
417 }
418 }
419
420 // =========================================================================
421 // http2_server_connection implementation
422 // =========================================================================
423
424 http2_server_connection::http2_server_connection(
425 uint64_t connection_id,
426 asio::ip::tcp::socket socket,
427 const http2_settings& settings,
428 http2_server::request_handler_t request_handler,
429 http2_server::error_handler_t error_handler)
430 : connection_id_(connection_id)
431 , use_tls_(false)
432 , plain_socket_(std::make_unique<asio::ip::tcp::socket>(std::move(socket)))
433 , local_settings_(settings)
434 , encoder_(settings.header_table_size)
435 , decoder_(settings.header_table_size)
436 , request_handler_(std::move(request_handler))
437 , error_handler_(std::move(error_handler))
438 , frame_header_buffer_{}
439 {
440 }
441
443 uint64_t connection_id,
444 std::unique_ptr<asio::ssl::stream<asio::ip::tcp::socket>> socket,
446 http2_server::request_handler_t request_handler,
447 http2_server::error_handler_t error_handler)
448 : connection_id_(connection_id)
449 , use_tls_(true)
450 , tls_socket_(std::move(socket))
451 , local_settings_(settings)
452 , encoder_(settings.header_table_size)
453 , decoder_(settings.header_table_size)
454 , request_handler_(std::move(request_handler))
455 , error_handler_(std::move(error_handler))
456 , frame_header_buffer_{}
457 {
458 }
459
464
466 {
467 // Read connection preface from client
468 read_connection_preface();
469 return ok();
470 }
471
473 {
474 std::lock_guard<std::mutex> lock(transport_shutdown_mutex_);
475 if (!is_alive_.exchange(false)) {
476 return ok();
477 }
478
479 // Close socket
480 std::error_code ec;
481 if (use_tls_ && tls_socket_) {
482 tls_socket_->lowest_layer().close(ec);
483 } else if (plain_socket_) {
484 plain_socket_->close(ec);
485 }
486
487 return ok();
488 }
489
491 {
492 // Shutdown interrupts synchronous writes without destroying the
493 // descriptor that the I/O worker may still be using. A read callback
494 // may close it concurrently, so serialize the descriptor operations.
495 std::lock_guard<std::mutex> lock(transport_shutdown_mutex_);
496 std::error_code ec;
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);
501 }
502 }
503
505 {
506 return is_alive_;
507 }
508
510 {
511 return connection_id_;
512 }
513
515 {
516 std::lock_guard<std::mutex> lock(const_cast<std::mutex&>(streams_mutex_));
517 return streams_.size();
518 }
519
521 {
522 constexpr size_t PREFACE_SIZE = 24; // "PRI * HTTP/2.0\r\n\r\nSM\r\n\r\n"
523 auto buffer = std::make_shared<std::vector<uint8_t>>(PREFACE_SIZE);
524
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");
529 }
530 stop();
531 return;
532 }
533
534 // Verify 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");
539 }
540 stop();
541 return;
542 }
543
544 preface_received_ = true;
545
546 // Send our SETTINGS frame
547 auto result = send_settings();
548 if (result.is_err()) {
549 if (error_handler_) {
550 error_handler_("Failed to send settings");
551 }
552 stop();
553 return;
554 }
555
556 // Start reading frames
557 start_reading();
558 };
559
560 if (use_tls_) {
561 asio::async_read(*tls_socket_, asio::buffer(*buffer), read_handler);
562 } else {
563 asio::async_read(*plain_socket_, asio::buffer(*buffer), read_handler);
564 }
565 }
566
568 {
569 std::vector<setting_parameter> params = {
570 {static_cast<uint16_t>(setting_identifier::max_concurrent_streams),
571 local_settings_.max_concurrent_streams},
572 {static_cast<uint16_t>(setting_identifier::initial_window_size),
573 local_settings_.initial_window_size},
574 {static_cast<uint16_t>(setting_identifier::max_frame_size),
575 local_settings_.max_frame_size},
576 {static_cast<uint16_t>(setting_identifier::header_table_size),
577 local_settings_.header_table_size},
578 };
579
580 settings_frame frame(params, false);
581 return send_frame(frame);
582 }
583
585 {
586 if (frame.is_ack()) {
587 // Settings acknowledgment received
588 return ok();
589 }
590
591 // Apply peer settings
592 for (const auto& param : frame.settings()) {
593 switch (static_cast<setting_identifier>(param.identifier)) {
595 remote_settings_.header_table_size = param.value;
596 encoder_.set_max_table_size(param.value);
597 break;
599 remote_settings_.enable_push = (param.value != 0);
600 break;
602 remote_settings_.max_concurrent_streams = param.value;
603 break;
605 remote_settings_.initial_window_size = param.value;
606 break;
608 remote_settings_.max_frame_size = param.value;
609 break;
611 remote_settings_.max_header_list_size = param.value;
612 break;
613 }
614 }
615
616 // Send ACK
617 return send_settings_ack();
618 }
619
621 {
622 settings_frame frame({}, true);
623 return send_frame(frame);
624 }
625
627 {
628 read_frame_header();
629 }
630
632 {
633 if (!is_alive_) {
634 return;
635 }
636
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());
642 }
643 }
644 stop();
645 return;
646 }
647
648 auto header_result = frame_header::parse(frame_header_buffer_);
649 if (header_result.is_err()) {
650 if (error_handler_) {
651 error_handler_("Failed to parse frame header");
652 }
653 stop();
654 return;
655 }
656
657 auto header = header_result.value();
658 if (header.length > 0) {
659 read_frame_payload(header.length);
660 } else {
661 // Process frame with empty payload
662 auto frame_result = frame::parse(frame_header_buffer_);
663 if (frame_result.is_ok()) {
664 process_frame(std::move(frame_result.value()));
665 }
666 read_frame_header();
667 }
668 };
669
670 if (use_tls_) {
671 asio::async_read(*tls_socket_, asio::buffer(frame_header_buffer_), read_handler);
672 } else {
673 asio::async_read(*plain_socket_, asio::buffer(frame_header_buffer_), read_handler);
674 }
675 }
676
678 {
679 if (!is_alive_) {
680 return;
681 }
682
683 read_buffer_.resize(9 + length);
684 std::copy(frame_header_buffer_.begin(), frame_header_buffer_.end(), read_buffer_.begin());
685
686 auto payload_buffer = asio::buffer(read_buffer_.data() + 9, length);
687
688 auto read_handler = [this, self = shared_from_this()](std::error_code ec, std::size_t /*bytes_read*/) {
689 if (ec) {
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());
693 }
694 }
695 stop();
696 return;
697 }
698
699 auto frame_result = frame::parse(read_buffer_);
700 if (frame_result.is_ok()) {
701 process_frame(std::move(frame_result.value()));
702 } else {
703 if (error_handler_) {
704 error_handler_("Failed to parse frame");
705 }
706 }
707
708 read_frame_header();
709 };
710
711 if (use_tls_) {
712 asio::async_read(*tls_socket_, payload_buffer, read_handler);
713 } else {
714 asio::async_read(*plain_socket_, payload_buffer, read_handler);
715 }
716 }
717
719 {
720 auto data = f.serialize();
721
722 std::error_code ec;
723 if (use_tls_) {
724 asio::write(*tls_socket_, asio::buffer(data), ec);
725 } else {
726 asio::write(*plain_socket_, asio::buffer(data), ec);
727 }
728
729 if (ec) {
730 return error_void(
732 std::string("Failed to send frame: ") + ec.message(),
733 "http2_server_connection");
734 }
735
736 return ok();
737 }
738
739 auto http2_server_connection::process_frame(std::unique_ptr<frame> f) -> VoidResult
740 {
741 auto& header = f->header();
742
743 switch (header.type) {
745 auto* settings_f = dynamic_cast<settings_frame*>(f.get());
746 if (settings_f) {
747 return handle_settings_frame(*settings_f);
748 }
749 break;
750 }
751 case frame_type::headers: {
752 auto* headers_f = dynamic_cast<headers_frame*>(f.get());
753 if (headers_f) {
754 return handle_headers_frame(*headers_f);
755 }
756 break;
757 }
758 case frame_type::data: {
759 auto* data_f = dynamic_cast<data_frame*>(f.get());
760 if (data_f) {
761 return handle_data_frame(*data_f);
762 }
763 break;
764 }
766 auto* rst_f = dynamic_cast<rst_stream_frame*>(f.get());
767 if (rst_f) {
768 return handle_rst_stream_frame(*rst_f);
769 }
770 break;
771 }
772 case frame_type::ping: {
773 auto* ping_f = dynamic_cast<ping_frame*>(f.get());
774 if (ping_f) {
775 return handle_ping_frame(*ping_f);
776 }
777 break;
778 }
779 case frame_type::goaway: {
780 auto* goaway_f = dynamic_cast<goaway_frame*>(f.get());
781 if (goaway_f) {
782 return handle_goaway_frame(*goaway_f);
783 }
784 break;
785 }
787 auto* wu_f = dynamic_cast<window_update_frame*>(f.get());
788 if (wu_f) {
789 return handle_window_update_frame(*wu_f);
790 }
791 break;
792 }
793 default:
794 // Ignore unknown frame types
795 break;
796 }
797
798 return ok();
799 }
800
802 {
803 std::lock_guard<std::mutex> lock(streams_mutex_);
804
805 auto it = streams_.find(stream_id);
806 if (it != streams_.end()) {
807 return &it->second;
808 }
809
810 // Create new stream
811 http2_stream stream;
812 stream.stream_id = stream_id;
813 stream.state = stream_state::open;
814 stream.window_size = static_cast<int32_t>(local_settings_.initial_window_size);
815
816 auto [iter, inserted] = streams_.emplace(stream_id, std::move(stream));
817 if (stream_id > last_stream_id_) {
818 last_stream_id_ = stream_id;
819 }
820
821 return &iter->second;
822 }
823
824 auto http2_server_connection::close_stream(uint32_t stream_id) -> void
825 {
826 std::lock_guard<std::mutex> lock(streams_mutex_);
827 streams_.erase(stream_id);
828 }
829
831 {
832 uint32_t stream_id = f.header().stream_id;
833 auto* stream = get_or_create_stream(stream_id);
834
835 if (!stream) {
836 return error_void(
838 "Failed to create stream",
839 "http2_server_connection");
840 }
841
842 // Decode headers
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);
846
847 if (headers_result.is_err()) {
848 // Send GOAWAY with COMPRESSION_ERROR
849 goaway_frame goaway(last_stream_id_,
850 static_cast<uint32_t>(error_code::compression_error));
851 send_frame(goaway);
852 stop();
853 return error_void(
854 headers_result.error().code,
855 headers_result.error().message,
856 "http2_server_connection");
857 }
858
859 stream->request_headers = headers_result.value();
860
861 if (f.is_end_stream()) {
862 stream->state = stream_state::half_closed_remote;
863 dispatch_request(stream_id);
864 } else if (f.is_end_headers()) {
865 stream->headers_complete = true;
866 }
867
868 return ok();
869 }
870
872 {
873 uint32_t stream_id = f.header().stream_id;
874
875 std::unique_lock<std::mutex> lock(streams_mutex_);
876 auto it = streams_.find(stream_id);
877 if (it == streams_.end()) {
878 // Stream doesn't exist
879 rst_stream_frame rst(stream_id, static_cast<uint32_t>(error_code::stream_closed));
880 return send_frame(rst);
881 }
882
883 auto& stream = it->second;
884 auto data = f.data();
885 stream.request_body.insert(stream.request_body.end(), data.begin(), data.end());
886
887 // Send WINDOW_UPDATE if needed
888 size_t data_size = data.size();
889 if (data_size > 0) {
890 // Update connection window
891 window_update_frame conn_wu(0, static_cast<uint32_t>(data_size));
892 send_frame(conn_wu);
893
894 // Update stream window
895 window_update_frame stream_wu(stream_id, static_cast<uint32_t>(data_size));
896 send_frame(stream_wu);
897 }
898
899 if (f.is_end_stream()) {
901 stream.body_complete = true;
902
903 // Need to unlock before dispatching
904 lock.unlock();
905 dispatch_request(stream_id);
906 }
907
908 return ok();
909 }
910
912 {
913 close_stream(f.header().stream_id);
914 return ok();
915 }
916
918 {
919 if (f.is_ack()) {
920 return ok();
921 }
922
923 // Send PING ACK
924 ping_frame ack(f.opaque_data(), true);
925 return send_frame(ack);
926 }
927
929 {
930 // Client is closing connection
931 stop();
932 return ok();
933 }
934
936 {
937 uint32_t stream_id = f.header().stream_id;
938 int32_t increment = static_cast<int32_t>(f.window_size_increment());
939
940 if (stream_id == 0) {
941 // Connection-level flow control
942 connection_window_size_ += increment;
943 } else {
944 // Stream-level flow control
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;
949 }
950 }
951
952 return ok();
953 }
954
955 auto http2_server_connection::dispatch_request(uint32_t stream_id) -> void
956 {
957 if (!request_handler_) {
958 return;
959 }
960
961 http2_stream* stream = nullptr;
962 {
963 std::lock_guard<std::mutex> lock(streams_mutex_);
964 auto it = streams_.find(stream_id);
965 if (it == streams_.end()) {
966 return;
967 }
968 stream = &it->second;
969 }
970
971 // Create request from headers
972 auto request = http2_request::from_headers(stream->request_headers);
973 request.body = stream->request_body;
974
975 // Create server stream for response
976 auto encoder_ptr = std::make_shared<hpack_encoder>(encoder_);
977 auto weak_this = weak_from_this();
978
979 http2_server_stream server_stream(
980 stream_id,
981 std::move(request),
982 encoder_ptr,
983 [weak_this](const frame& f) -> VoidResult {
984 auto self = weak_this.lock();
985 if (!self) {
986 return error_void(
988 "Connection closed",
989 "http2_server_stream");
990 }
991 return self->send_frame(f);
992 },
993 remote_settings_.max_frame_size);
994
995 // Call request handler
996 try {
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());
1001 }
1002 }
1003
1004 // Close stream after handler completes
1005 close_stream(stream_id);
1006 }
1007
1008} // namespace kcenon::network::protocols::http2
DATA frame (RFC 7540 Section 6.1)
Definition frame.h:139
Base class for HTTP/2 frames.
Definition frame.h:82
static auto parse(std::span< const uint8_t > data) -> Result< std::unique_ptr< frame > >
Parse frame from raw bytes.
Definition frame.cpp:94
GOAWAY frame (RFC 7540 Section 6.8)
Definition frame.h:387
HEADERS frame (RFC 7540 Section 6.2)
Definition frame.h:189
HPACK header decoder (RFC 7541)
Definition hpack.h:209
HPACK header encoder (RFC 7541)
Definition hpack.h:159
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 get_or_create_stream(uint32_t stream_id) -> http2_stream *
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 connection_id() const -> uint64_t
Get connection identifier.
auto handle_window_update_frame(const window_update_frame &f) -> VoidResult
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 handle_headers_frame(const headers_frame &f) -> VoidResult
auto is_alive() const -> bool
Check if connection is alive.
auto handle_ping_frame(const ping_frame &f) -> VoidResult
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.
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)
Definition frame.h:344
RST_STREAM frame (RFC 7540 Section 6.4)
Definition frame.h:307
SETTINGS frame (RFC 7540 Section 6.5)
Definition frame.h:265
WINDOW_UPDATE frame (RFC 7540 Section 6.9)
Definition frame.h:438
struct ssl_st SSL
Definition crypto.h:21
tracing_config config
Definition exporters.cpp:29
@ 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)
Definition frame.h:248
@ max_header_list_size
SETTINGS_MAX_HEADER_LIST_SIZE.
@ max_concurrent_streams
SETTINGS_MAX_CONCURRENT_STREAMS.
@ initial_window_size
SETTINGS_INITIAL_WINDOW_SIZE.
@ half_closed_remote
Remote end closed, local can send.
VoidResult error_void(int code, const std::string &message, const std::string &source="network_system", const std::string &details="")
VoidResult ok()
static auto parse(std::span< const uint8_t > data) -> Result< frame_header >
Parse frame header from raw bytes.
Definition frame.cpp:52
static auto from_headers(const std::vector< http_header > &parsed_headers) -> http2_request
Create http2_request from parsed headers.
std::vector< uint8_t > request_body
Request body.
std::atomic< stream_state > state
Current state.
std::vector< http_header > request_headers
Request headers.
TLS configuration for HTTP/2 server.