Network System 0.1.1
High-performance modular networking library for scalable client-server applications
Loading...
Searching...
No Matches
http2_client.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
10
11#include <algorithm>
12#include <limits>
13#include <thread>
14
16{
17 // http2_response implementation
18 auto http2_response::get_header(const std::string& name) const -> std::optional<std::string>
19 {
20 std::string lower_name = name;
21 std::transform(lower_name.begin(), lower_name.end(), lower_name.begin(),
22 [](unsigned char c) { return std::tolower(c); });
23
24 for (const auto& header : headers)
25 {
26 std::string lower_header_name = header.name;
27 std::transform(lower_header_name.begin(), lower_header_name.end(),
28 lower_header_name.begin(),
29 [](unsigned char c) { return std::tolower(c); });
30
31 if (lower_header_name == lower_name)
32 {
33 return header.value;
34 }
35 }
36 return std::nullopt;
37 }
38
39 auto http2_response::get_body_string() const -> std::string
40 {
41 return std::string(body.begin(), body.end());
42 }
43
44 // http2_client implementation
45 http2_client::http2_client(std::string_view client_id)
46 : client_id_(client_id)
47 , encoder_(local_settings_.header_table_size)
48 , decoder_(local_settings_.header_table_size)
49 {
50 }
51
53 {
54 try
55 {
56 disconnect();
57 }
58 catch (...)
59 {
60 // Suppress exceptions in destructor
61 }
62 }
63
64 auto http2_client::connect(const std::string& host, unsigned short port) -> VoidResult
65 {
66 std::lock_guard<std::mutex> lifecycle_lock(lifecycle_mutex_);
67 // Create tracing span for connect operation
68 auto span = tracing::is_tracing_enabled()
69 ? std::make_optional(tracing::trace_context::create_span("http2.client.connect"))
70 : std::nullopt;
71 if (span)
72 {
73 span->set_attribute("net.peer.name", host)
74 .set_attribute("net.peer.port", static_cast<int64_t>(port))
75 .set_attribute("net.transport", "tcp")
76 .set_attribute("http.flavor", "2.0")
77 .set_attribute("client.id", client_id_);
78 }
79
80 if (is_connected_)
81 {
82 if (span)
83 {
84 span->set_error("Already connected");
85 }
87 "Already connected", "http2_client::connect");
88 }
89
90 if (host.empty())
91 {
92 if (span)
93 {
94 span->set_error("Host cannot be empty");
95 }
97 "Host cannot be empty", "http2_client::connect");
98 }
99
100 host_ = host;
101 port_ = port;
102
103 try
104 {
105 // A failed connection can still own a socket and work guard.
106 // Release them while their original execution context is alive.
107 stop_io();
108 socket_.reset();
109 work_guard_.reset();
110 ssl_context_.reset();
111 io_context_.reset();
112 goaway_received_ = false;
113
114 // Create I/O context
115 io_context_ = std::make_unique<asio::io_context>();
116
117 // Create SSL context with TLS 1.3
118 ssl_context_ = std::make_unique<asio::ssl::context>(asio::ssl::context::tlsv13_client);
119
120 // Set ALPN for HTTP/2
121 SSL_CTX* native_ctx = ssl_context_->native_handle();
122 static const unsigned char alpn_protos[] = { 2, 'h', '2' };
123 SSL_CTX_set_alpn_protos(native_ctx, alpn_protos, sizeof(alpn_protos));
124
125 // Verify peer certificate
126 ssl_context_->set_default_verify_paths();
127 ssl_context_->set_verify_mode(asio::ssl::verify_peer);
128
129 // Create socket
130 asio::ip::tcp::resolver resolver(*io_context_);
131 auto endpoints = resolver.resolve(host, std::to_string(port));
132
133 socket_ = std::make_unique<asio::ssl::stream<asio::ip::tcp::socket>>(
134 *io_context_, *ssl_context_);
135
136 // Set SNI hostname
137 SSL_set_tlsext_host_name(socket_->native_handle(), host.c_str());
138
139 // Connect TCP
140 asio::connect(socket_->lowest_layer(), endpoints);
141
142 // Perform SSL handshake
143 socket_->handshake(asio::ssl::stream_base::client);
144
145 // Verify ALPN negotiation
146 const unsigned char* alpn_result = nullptr;
147 unsigned int alpn_len = 0;
148 SSL_get0_alpn_selected(socket_->native_handle(), &alpn_result, &alpn_len);
149
150 if (alpn_len != 2 || alpn_result == nullptr ||
151 std::string(reinterpret_cast<const char*>(alpn_result), alpn_len) != "h2")
152 {
153 socket_->lowest_layer().close();
154 if (span)
155 {
156 span->set_error("Server does not support HTTP/2 via ALPN");
157 }
159 "Server does not support HTTP/2 via ALPN",
160 "http2_client::connect");
161 }
162
163 is_running_ = true;
164
165 // Send connection preface
166 auto preface_result = send_connection_preface();
167 if (preface_result.is_err())
168 {
169 is_connected_ = false;
170 is_running_ = false;
171 if (span)
172 {
173 span->set_error(preface_result.error().message);
174 }
175 return preface_result;
176 }
177
178 // Send initial SETTINGS
179 auto settings_result = send_settings();
180 if (settings_result.is_err())
181 {
182 is_connected_ = false;
183 is_running_ = false;
184 if (span)
185 {
186 span->set_error(settings_result.error().message);
187 }
188 return settings_result;
189 }
190
191 // The server connection preface must start with a non-ACK
192 // SETTINGS frame on stream zero (RFC 9113 section 3.4).
193 // Bound both reads by one deadline so a partial frame cannot
194 // leave connect() blocked indefinitely.
195 const auto deadline = std::chrono::steady_clock::now() + timeout_.load();
196 auto read_preface_bytes = [&](asio::mutable_buffer buffer) {
197 asio::steady_timer timer(*io_context_, deadline);
198 std::error_code read_error;
199 bool expired = false;
200 timer.async_wait([&](std::error_code ec) {
201 if (!ec) {
202 expired = true;
203 std::error_code ignored;
204 socket_->lowest_layer().cancel(ignored);
205 }
206 });
207 asio::async_read(*socket_, buffer,
208 [&](std::error_code ec, std::size_t) {
209 read_error = ec;
210 timer.cancel();
211 });
212 io_context_->restart();
213 io_context_->run();
214 if (expired) throw std::runtime_error("Server SETTINGS timed out");
215 if (read_error) throw std::system_error(read_error);
216 };
217 std::vector<uint8_t> initial_header(FRAME_HEADER_SIZE);
218 read_preface_bytes(asio::buffer(initial_header));
219 auto header = frame_header::parse(initial_header);
220 if (header.is_err() || header.value().type != frame_type::settings ||
221 header.value().stream_id != 0 ||
222 (header.value().flags & frame_flags::ack) != 0 ||
223 header.value().length > local_settings_.max_frame_size) {
224 throw std::runtime_error("Invalid server SETTINGS preface");
225 }
226 std::vector<uint8_t> initial_payload(header.value().length);
227 if (!initial_payload.empty()) read_preface_bytes(asio::buffer(initial_payload));
228 initial_header.insert(initial_header.end(), initial_payload.begin(), initial_payload.end());
229 auto initial_frame = frame::parse(initial_header);
230 if (initial_frame.is_err()) throw std::runtime_error("Invalid server SETTINGS payload");
231 auto initial_result = process_frame(std::move(initial_frame.value()));
232 if (initial_result.is_err()) throw std::runtime_error(initial_result.error().message);
233
234 // All subsequent SSL operations run on this one executor thread.
235 // Publish the connected state while submissions are locked so no
236 // caller writes synchronously between startup and the first read.
237 work_guard_ = std::make_unique<asio::executor_work_guard<asio::io_context::executor_type>>(
238 asio::make_work_guard(*io_context_));
239 io_context_->restart();
240 {
241 std::lock_guard<std::mutex> lock(io_submission_mutex_);
242 io_active_ = true;
243 is_connected_ = true;
244 io_future_ = std::async(std::launch::async, [this]() { run_io(); });
245 }
246
247 if (span)
248 {
249 span->set_status(tracing::span_status::ok);
250 }
251 return ok();
252 }
253 catch (const std::exception& e)
254 {
255 is_connected_ = false;
256 is_running_ = false;
257 stop_io();
258 if (span)
259 {
260 span->set_error(std::string("Connection failed: ") + e.what());
261 }
263 std::string("Connection failed: ") + e.what(),
264 "http2_client::connect");
265 }
266 }
267
269 {
270 std::lock_guard<std::mutex> lifecycle_lock(lifecycle_mutex_);
271 if (!is_connected_)
272 {
273 stop_io();
274 return ok();
275 }
276
277 try
278 {
279 // Send GOAWAY frame
280 if (socket_ && socket_->lowest_layer().is_open())
281 {
282 goaway_frame goaway(next_stream_id_ - 2,
283 static_cast<uint32_t>(error_code::no_error));
284 send_frame(goaway);
285 }
286 }
287 catch (...)
288 {
289 // Ignore errors during GOAWAY
290 }
291
292 stop_io();
293
294 is_connected_ = false;
295
296 return ok();
297 }
298
299 auto http2_client::is_connected() const -> bool
300 {
302 }
303
304 auto http2_client::get(const std::string& path,
305 const std::vector<http_header>& headers)
307 {
308 return send_request("GET", path, headers, {});
309 }
310
311 auto http2_client::post(const std::string& path,
312 const std::string& body,
313 const std::vector<http_header>& headers)
315 {
316 std::vector<uint8_t> body_bytes(body.begin(), body.end());
317 return send_request("POST", path, headers, body_bytes);
318 }
319
320 auto http2_client::post(const std::string& path,
321 const std::vector<uint8_t>& body,
322 const std::vector<http_header>& headers)
324 {
325 return send_request("POST", path, headers, body);
326 }
327
328 auto http2_client::put(const std::string& path,
329 const std::string& body,
330 const std::vector<http_header>& headers)
332 {
333 std::vector<uint8_t> body_bytes(body.begin(), body.end());
334 return send_request("PUT", path, headers, body_bytes);
335 }
336
337 auto http2_client::del(const std::string& path,
338 const std::vector<http_header>& headers)
340 {
341 return send_request("DELETE", path, headers, {});
342 }
343
344 auto http2_client::set_timeout(std::chrono::milliseconds timeout) -> void
345 {
346 timeout_ = timeout;
347 }
348
349 auto http2_client::get_timeout() const -> std::chrono::milliseconds
350 {
351 return timeout_.load();
352 }
353
355 {
356 return local_settings_;
357 }
358
360 {
361 local_settings_ = settings;
362 encoder_.set_max_table_size(settings.header_table_size);
363 decoder_.set_max_table_size(settings.header_table_size);
364 }
365
366 auto http2_client::start_stream(const std::string& path,
367 const std::vector<http_header>& headers,
368 std::function<void(std::vector<uint8_t>)> on_data,
369 std::function<void(std::vector<http_header>)> on_headers,
370 std::function<void(int)> on_complete)
372 {
373 if (!is_connected())
374 {
376 "Not connected",
377 "http2_client::start_stream");
378 }
379
380 // Create stream with callbacks
381 http2_stream& stream = create_stream();
382 stream.state = stream_state::open;
383 stream.is_streaming = true;
384 stream.on_data = std::move(on_data);
385 stream.on_headers = std::move(on_headers);
386 stream.on_complete = std::move(on_complete);
387
388 // Build headers for POST (streaming requests use POST)
389 auto request_headers = build_headers("POST", path, headers);
390 stream.request_headers = request_headers;
391
392 // Encode headers with HPACK
393 auto encoded_headers = encoder_.encode(request_headers);
394
395 // Send HEADERS frame (not end_stream since we might send more data)
396 headers_frame hf(stream.stream_id, std::move(encoded_headers), false, true);
397
398 auto send_result = send_frame(hf);
399 if (send_result.is_err())
400 {
401 close_stream(stream.stream_id);
402 const auto& err = send_result.error();
403 return error<uint32_t>(err.code, err.message,
404 "http2_client::start_stream",
405 get_error_details(err));
406 }
407
408 uint32_t result_id = stream.stream_id;
409 return ok(std::move(result_id));
410 }
411
412 auto http2_client::write_stream(uint32_t stream_id,
413 const std::vector<uint8_t>& data,
414 bool end_stream) -> VoidResult
415 {
416 if (!is_connected())
417 {
419 "Not connected",
420 "http2_client::write_stream");
421 }
422
423 auto* stream = get_stream(stream_id);
424 if (!stream)
425 {
427 "Stream not found",
428 "http2_client::write_stream");
429 }
430
431 if (stream->state == stream_state::closed ||
432 stream->state == stream_state::half_closed_local)
433 {
435 "Stream not writable",
436 "http2_client::write_stream");
437 }
438
439 // Send DATA frame
440 if (end_stream) stream->state = stream_state::half_closed_local;
441 data_frame df(stream_id, std::vector<uint8_t>(data), end_stream);
442 auto send_result = send_frame(df);
443 if (send_result.is_err())
444 {
445 return send_result;
446 }
447
448 return ok();
449 }
450
452 {
453 if (!is_connected())
454 {
456 "Not connected",
457 "http2_client::close_stream_writer");
458 }
459
460 auto* stream = get_stream(stream_id);
461 if (!stream)
462 {
464 "Stream not found",
465 "http2_client::close_stream_writer");
466 }
467
468 if (stream->state == stream_state::closed ||
469 stream->state == stream_state::half_closed_local)
470 {
471 return ok(); // Already closed
472 }
473
474 // Publish the state before the peer can respond to END_STREAM.
475 stream->state = stream_state::half_closed_local;
476 // Send empty DATA frame with END_STREAM
477 data_frame df(stream_id, {}, true);
478 auto send_result = send_frame(df);
479 if (send_result.is_err())
480 {
481 return send_result;
482 }
483
484 return ok();
485 }
486
487 auto http2_client::cancel_stream(uint32_t stream_id) -> VoidResult
488 {
489 if (!is_connected())
490 {
492 "Not connected",
493 "http2_client::cancel_stream");
494 }
495
496 auto* stream = get_stream(stream_id);
497 if (!stream)
498 {
500 "Stream not found",
501 "http2_client::cancel_stream");
502 }
503
504 if (stream->state == stream_state::closed)
505 {
506 return ok(); // Already closed
507 }
508
509 // Send RST_STREAM frame
510 rst_stream_frame rsf(stream_id, static_cast<uint32_t>(error_code::cancel));
511 auto send_result = send_frame(rsf);
512 if (send_result.is_err())
513 {
514 return send_result;
515 }
516
517 stream->state = stream_state::closed;
518 return ok();
519 }
520
521 // Private methods
522
524 {
525 try
526 {
527 asio::write(*socket_,
528 asio::buffer(CONNECTION_PREFACE.data(), CONNECTION_PREFACE.size()));
529 return ok();
530 }
531 catch (const std::exception& e)
532 {
534 std::string("Failed to send connection preface: ") + e.what(),
535 "http2_client::send_connection_preface");
536 }
537 }
538
540 {
541 std::vector<setting_parameter> params = {
542 {static_cast<uint16_t>(setting_identifier::header_table_size),
543 local_settings_.header_table_size},
544 {static_cast<uint16_t>(setting_identifier::enable_push),
545 local_settings_.enable_push ? 1u : 0u},
546 {static_cast<uint16_t>(setting_identifier::max_concurrent_streams),
547 local_settings_.max_concurrent_streams},
548 {static_cast<uint16_t>(setting_identifier::initial_window_size),
549 local_settings_.initial_window_size},
550 {static_cast<uint16_t>(setting_identifier::max_frame_size),
551 local_settings_.max_frame_size},
552 {static_cast<uint16_t>(setting_identifier::max_header_list_size),
553 local_settings_.max_header_list_size}
554 };
555
556 settings_frame settings(params, false);
557 return send_frame(settings);
558 }
559
561 {
562 if (frame.is_ack())
563 {
564 // Our settings were acknowledged
565 return ok();
566 }
567
568 // Apply remote settings
569 for (const auto& param : frame.settings())
570 {
571 switch (static_cast<setting_identifier>(param.identifier))
572 {
574 remote_settings_.header_table_size = param.value;
575 encoder_.set_max_table_size(param.value);
576 break;
578 remote_settings_.enable_push = (param.value != 0);
579 break;
581 remote_settings_.max_concurrent_streams = param.value;
582 break;
584 remote_settings_.initial_window_size = param.value;
585 break;
587 remote_settings_.max_frame_size = param.value;
588 break;
590 remote_settings_.max_header_list_size = param.value;
591 break;
592 }
593 }
594
595 return send_settings_ack();
596 }
597
599 {
600 settings_frame ack({}, true);
601 return send_frame(ack);
602 }
603
605 {
606 std::unique_lock<std::mutex> lock(io_submission_mutex_);
607 if (!socket_ || !is_running_)
608 {
610 "Connection closed", "http2_client::send_frame");
611 }
612
613 try
614 {
615 auto data = f.serialize();
616 if (!io_active_)
617 {
618 // Connection preface only: the reader has not started yet.
619 asio::write(*socket_, asio::buffer(data));
620 return ok();
621 }
622
623 auto completion = std::make_shared<std::promise<VoidResult>>();
624 auto result = completion->get_future();
625 const bool on_io_thread = io_context_->get_executor().running_in_this_thread();
626 asio::post(*io_context_, [this, data = std::move(data), completion]() mutable {
627 if (!is_running_)
628 {
629 completion->set_value(error_void(
631 "Connection closed", "http2_client::send_frame"));
632 return;
633 }
634 const bool idle = pending_writes_.empty();
635 pending_writes_.push_back({std::move(data), completion});
636 if (idle) write_next_frame();
637 });
638 lock.unlock();
639 // Frame handlers enqueue protocol replies without waiting on their
640 // own executor. Public callers retain synchronous send results.
641 return on_io_thread ? ok() : result.get();
642 }
643 catch (const std::exception& e)
644 {
646 std::string("Failed to send frame: ") + e.what(),
647 "http2_client::send_frame");
648 }
649 }
650
652 {
653 if (pending_writes_.empty()) return;
654 asio::async_write(*socket_, asio::buffer(pending_writes_.front().data),
655 [this](std::error_code ec, std::size_t) {
656 auto completion = std::move(pending_writes_.front().completion);
657 pending_writes_.pop_front();
658 if (ec)
659 {
660 is_connected_ = false;
661 completion->set_value(error_void(
662 error_codes::network_system::send_failed, ec.message(),
663 "http2_client::send_frame"));
664 while (!pending_writes_.empty())
665 {
666 pending_writes_.front().completion->set_value(error_void(
667 error_codes::network_system::send_failed, ec.message(),
668 "http2_client::send_frame"));
669 pending_writes_.pop_front();
670 }
671 return;
672 }
673 completion->set_value(ok());
674 write_next_frame();
675 });
676 }
677
678 auto http2_client::read_next_frame() -> void
679 {
680 if (!is_running_) return;
681 auto bytes = std::make_shared<std::vector<uint8_t>>(FRAME_HEADER_SIZE);
682 asio::async_read(*socket_, asio::buffer(*bytes),
683 [this, bytes](std::error_code ec, std::size_t) {
684 if (ec || !is_running_)
685 {
686 if (ec) is_connected_ = false;
687 return;
688 }
689 auto header = frame_header::parse(*bytes);
690 if (header.is_err())
691 {
692 is_connected_ = false;
693 return;
694 }
695 const auto length = header.value().length;
696 bytes->resize(FRAME_HEADER_SIZE + length);
697 auto complete = [this, bytes](std::error_code payload_error, std::size_t) {
698 if (payload_error || !is_running_)
699 {
700 if (payload_error) is_connected_ = false;
701 return;
702 }
703 try
704 {
705 auto parsed = frame::parse(*bytes);
706 if (parsed.is_err())
707 {
708 is_connected_ = false;
709 return;
710 }
711 process_frame(std::move(parsed.value()));
712 read_next_frame();
713 }
714 catch (...)
715 {
716 is_connected_ = false;
717 }
718 };
719 if (length == 0)
720 complete({}, 0);
721 else
722 asio::async_read(*socket_,
723 asio::buffer(bytes->data() + FRAME_HEADER_SIZE, length),
724 std::move(complete));
725 });
726 }
727
728 auto http2_client::process_frame(std::unique_ptr<frame> f) -> VoidResult
729 {
730 if (!f)
731 {
733 "Null frame", "http2_client::process_frame");
734 }
735
736 switch (f->header().type)
737 {
738 case frame_type::settings:
739 if (auto* sf = dynamic_cast<settings_frame*>(f.get()))
740 {
741 return handle_settings_frame(*sf);
742 }
743 break;
744
745 case frame_type::headers:
746 if (auto* hf = dynamic_cast<headers_frame*>(f.get()))
747 {
748 return handle_headers_frame(*hf);
749 }
750 break;
751
752 case frame_type::data:
753 if (auto* df = dynamic_cast<data_frame*>(f.get()))
754 {
755 return handle_data_frame(*df);
756 }
757 break;
758
759 case frame_type::rst_stream:
760 if (auto* rf = dynamic_cast<rst_stream_frame*>(f.get()))
761 {
762 return handle_rst_stream_frame(*rf);
763 }
764 break;
765
766 case frame_type::goaway:
767 if (auto* gf = dynamic_cast<goaway_frame*>(f.get()))
768 {
769 return handle_goaway_frame(*gf);
770 }
771 break;
772
773 case frame_type::window_update:
774 if (auto* wf = dynamic_cast<window_update_frame*>(f.get()))
775 {
776 return handle_window_update_frame(*wf);
777 }
778 break;
779
780 case frame_type::ping:
781 if (auto* pf = dynamic_cast<ping_frame*>(f.get()))
782 {
783 return handle_ping_frame(*pf);
784 }
785 break;
786
787 default:
788 // Ignore unknown frame types
789 break;
790 }
791
792 return ok();
793 }
794
795 auto http2_client::allocate_stream_id() -> uint32_t
796 {
797 uint32_t id = next_stream_id_.fetch_add(2);
798 return id;
799 }
800
801 auto http2_client::get_stream(uint32_t stream_id) -> http2_stream*
802 {
803 std::lock_guard<std::mutex> lock(streams_mutex_);
804 auto it = streams_.find(stream_id);
805 if (it != streams_.end())
806 {
807 return &it->second;
808 }
809 return nullptr;
810 }
811
812 auto http2_client::create_stream() -> http2_stream&
813 {
814 std::lock_guard<std::mutex> lock(streams_mutex_);
815 uint32_t stream_id = allocate_stream_id();
816 auto& stream = streams_[stream_id];
817 stream.stream_id = stream_id;
818 stream.state = stream_state::idle;
819 stream.window_size = remote_settings_.initial_window_size;
820 return stream;
821 }
822
823 auto http2_client::close_stream(uint32_t stream_id) -> void
824 {
825 std::lock_guard<std::mutex> lock(streams_mutex_);
826 auto it = streams_.find(stream_id);
827 if (it != streams_.end())
828 {
829 it->second.state = stream_state::closed;
830 }
831 }
832
833 auto http2_client::send_request(const std::string& method,
834 const std::string& path,
835 const std::vector<http_header>& headers,
836 const std::vector<uint8_t>& body)
838 {
839 // Create tracing span for HTTP request
840 auto span = tracing::is_tracing_enabled()
841 ? std::make_optional(tracing::trace_context::create_span("http2.request"))
842 : std::nullopt;
843 if (span)
844 {
845 span->set_attribute("http.method", method)
846 .set_attribute("http.url", "https://" + host_ + ":" + std::to_string(port_) + path)
847 .set_attribute("http.flavor", "2.0")
848 .set_attribute("http.request.body_size", static_cast<int64_t>(body.size()))
849 .set_attribute("net.peer.name", host_)
850 .set_attribute("net.peer.port", static_cast<int64_t>(port_));
851 }
852
853 if (!is_connected())
854 {
855 if (span)
856 {
857 span->set_error("Not connected");
858 }
860 "Not connected",
861 "http2_client::send_request");
862 }
863
864 // Create stream
865 http2_stream& stream = create_stream();
866 stream.state = stream_state::open;
867
868 if (span)
869 {
870 span->set_attribute("http2.stream_id", static_cast<int64_t>(stream.stream_id));
871 }
872
873 // Build headers
874 auto request_headers = build_headers(method, path, headers);
875 stream.request_headers = request_headers;
876 stream.request_body = body;
877
878 // Encode headers with HPACK
879 auto encoded_headers = encoder_.encode(request_headers);
880
881 // Send HEADERS frame
882 bool end_stream = body.empty();
883 if (end_stream) stream.state = stream_state::half_closed_local;
884 headers_frame hf(stream.stream_id, std::move(encoded_headers), end_stream, true);
885
886 auto send_result = send_frame(hf);
887 if (send_result.is_err())
888 {
889 close_stream(stream.stream_id);
890 const auto& err = send_result.error();
891 if (span)
892 {
893 span->set_error(err.message);
894 }
895 return error<http2_response>(err.code, err.message,
896 "http2_client::send_request",
897 get_error_details(err));
898 }
899
900 // Send DATA frame if body exists
901 if (!body.empty())
902 {
903 stream.state = stream_state::half_closed_local;
904 data_frame df(stream.stream_id, std::vector<uint8_t>(body), true);
905 send_result = send_frame(df);
906 if (send_result.is_err())
907 {
908 close_stream(stream.stream_id);
909 const auto& err = send_result.error();
910 if (span)
911 {
912 span->set_error(err.message);
913 }
914 return error<http2_response>(err.code, err.message,
915 "http2_client::send_request",
916 get_error_details(err));
917 }
918 }
919
920 // Wait for response with timeout
921 auto future = stream.promise.get_future();
922 auto status = future.wait_for(timeout_.load());
923
924 if (status == std::future_status::timeout)
925 {
926 close_stream(stream.stream_id);
927 if (span)
928 {
929 span->set_error("Request timeout");
930 }
932 "Request timeout",
933 "http2_client::send_request");
934 }
935
936 auto response = future.get();
937 if (span)
938 {
939 span->set_attribute("http.status_code", static_cast<int64_t>(response.status_code))
940 .set_attribute("http.response.body_size", static_cast<int64_t>(response.body.size()));
941 if (response.status_code >= 400)
942 {
943 span->set_status(tracing::span_status::error, "HTTP " + std::to_string(response.status_code));
944 }
945 else
946 {
947 span->set_status(tracing::span_status::ok);
948 }
949 }
950 return ok(std::move(response));
951 }
952
953 auto http2_client::build_headers(const std::string& method,
954 const std::string& path,
955 const std::vector<http_header>& additional)
956 -> std::vector<http_header>
957 {
958 std::vector<http_header> headers;
959
960 // Pseudo-headers (must come first)
961 headers.emplace_back(":method", method);
962 headers.emplace_back(":scheme", "https");
963 headers.emplace_back(":authority", host_);
964 headers.emplace_back(":path", path.empty() ? "/" : path);
965
966 // Standard headers
967 headers.emplace_back("user-agent", "network_system/http2_client");
968
969 // Additional headers
970 for (const auto& h : additional)
971 {
972 // Skip pseudo-headers in additional
973 if (!h.name.empty() && h.name[0] != ':')
974 {
975 headers.push_back(h);
976 }
977 }
978
979 return headers;
980 }
981
982 auto http2_client::handle_headers_frame(const headers_frame& f) -> VoidResult
983 {
984 uint32_t stream_id = f.header().stream_id;
985 auto* stream = get_stream(stream_id);
986
987 if (!stream)
988 {
990 "Unknown stream ID",
991 "http2_client::handle_headers_frame");
992 }
993
994 // Decode headers
995 auto decode_result = decoder_.decode(f.header_block());
996 if (decode_result.is_err())
997 {
998 const auto& err = decode_result.error();
999 return error_void(err.code, err.message,
1000 "http2_client::handle_headers_frame",
1001 get_error_details(err));
1002 }
1003
1004 auto& decoded_headers = decode_result.value();
1005
1006 // Append to stream headers
1007 stream->response_headers.insert(stream->response_headers.end(),
1008 decoded_headers.begin(),
1009 decoded_headers.end());
1010
1011 if (f.is_end_headers())
1012 {
1013 stream->headers_complete = true;
1014
1015 // Call streaming callback if set
1016 if (stream->is_streaming && stream->on_headers)
1017 {
1018 stream->on_headers(stream->response_headers);
1019 }
1020 }
1021
1022 if (f.is_end_stream())
1023 {
1024 stream->body_complete = true;
1025 stream->state = stream_state::closed;
1026
1027 // Extract status code
1028 int status_code = 0;
1029 for (const auto& h : stream->response_headers)
1030 {
1031 if (h.name == ":status")
1032 {
1033 try
1034 {
1035 status_code = std::stoi(h.value);
1036 }
1037 catch (...)
1038 {
1039 status_code = 0;
1040 }
1041 break;
1042 }
1043 }
1044
1045 // Handle streaming vs non-streaming
1046 if (stream->is_streaming)
1047 {
1048 if (stream->on_complete)
1049 {
1050 stream->on_complete(status_code);
1051 }
1052 }
1053 else
1054 {
1055 // Build response for non-streaming
1056 http2_response response;
1057 response.headers = stream->response_headers;
1058 response.body = stream->response_body;
1059 response.status_code = status_code;
1060 stream->promise.set_value(std::move(response));
1061 }
1062 }
1063
1064 return ok();
1065 }
1066
1067 auto http2_client::handle_data_frame(const data_frame& f) -> VoidResult
1068 {
1069 uint32_t stream_id = f.header().stream_id;
1070 auto* stream = get_stream(stream_id);
1071
1072 if (!stream)
1073 {
1075 "Unknown stream ID",
1076 "http2_client::handle_data_frame");
1077 }
1078
1079 // Get data
1080 auto data = f.data();
1081
1082 // For streaming, call callback immediately; otherwise buffer
1083 if (stream->is_streaming && stream->on_data)
1084 {
1085 stream->on_data(std::vector<uint8_t>(data.begin(), data.end()));
1086 }
1087 else
1088 {
1089 stream->response_body.insert(stream->response_body.end(),
1090 data.begin(), data.end());
1091 }
1092
1093 // Update flow control window
1094 stream->window_size -= static_cast<int32_t>(data.size());
1095 connection_window_size_ -= static_cast<int32_t>(data.size());
1096
1097 // Send WINDOW_UPDATE if needed
1098 if (stream->window_size < static_cast<int32_t>(DEFAULT_WINDOW_SIZE / 2))
1099 {
1100 int32_t increment = DEFAULT_WINDOW_SIZE - stream->window_size;
1101 window_update_frame wuf(stream_id, increment);
1102 send_frame(wuf);
1103 stream->window_size += increment;
1104 }
1105
1106 if (connection_window_size_ < static_cast<int32_t>(DEFAULT_WINDOW_SIZE / 2))
1107 {
1108 int32_t increment = DEFAULT_WINDOW_SIZE - connection_window_size_;
1109 window_update_frame wuf(0, increment); // Stream 0 = connection level
1110 send_frame(wuf);
1111 connection_window_size_ += increment;
1112 }
1113
1114 if (f.is_end_stream())
1115 {
1116 stream->body_complete = true;
1117 stream->state = stream_state::closed;
1118
1119 // Extract status code
1120 int status_code = 0;
1121 for (const auto& h : stream->response_headers)
1122 {
1123 if (h.name == ":status")
1124 {
1125 try
1126 {
1127 status_code = std::stoi(h.value);
1128 }
1129 catch (...)
1130 {
1131 status_code = 0;
1132 }
1133 break;
1134 }
1135 }
1136
1137 // Handle streaming vs non-streaming
1138 if (stream->is_streaming)
1139 {
1140 if (stream->on_complete)
1141 {
1142 stream->on_complete(status_code);
1143 }
1144 }
1145 else
1146 {
1147 http2_response response;
1148 response.headers = stream->response_headers;
1149 response.body = stream->response_body;
1150 response.status_code = status_code;
1151 stream->promise.set_value(std::move(response));
1152 }
1153 }
1154
1155 return ok();
1156 }
1157
1158 auto http2_client::handle_rst_stream_frame(const rst_stream_frame& f) -> VoidResult
1159 {
1160 uint32_t stream_id = f.header().stream_id;
1161 auto* stream = get_stream(stream_id);
1162
1163 if (stream)
1164 {
1165 stream->state = stream_state::closed;
1166
1167 // Set error response
1168 http2_response response;
1169 response.status_code = 0; // Indicate error
1170 stream->promise.set_value(std::move(response));
1171 }
1172
1173 return ok();
1174 }
1175
1176 auto http2_client::handle_goaway_frame(const goaway_frame& f) -> VoidResult
1177 {
1178 goaway_received_ = true;
1179
1180 // Close all streams with ID > last_stream_id
1181 uint32_t last_stream = f.last_stream_id();
1182 {
1183 std::lock_guard<std::mutex> lock(streams_mutex_);
1184 for (auto& [id, stream] : streams_)
1185 {
1186 if (id > last_stream && stream.state != stream_state::closed)
1187 {
1188 stream.state = stream_state::closed;
1189 http2_response response;
1190 response.status_code = 0;
1191 stream.promise.set_value(std::move(response));
1192 }
1193 }
1194 }
1195
1196 return ok();
1197 }
1198
1199 auto http2_client::handle_window_update_frame(const window_update_frame& f) -> VoidResult
1200 {
1201 uint32_t stream_id = f.header().stream_id;
1202 int32_t increment = static_cast<int32_t>(f.window_size_increment());
1203
1204 if (stream_id == 0)
1205 {
1206 if (connection_window_size_ > std::numeric_limits<int32_t>::max() - increment)
1207 {
1208 // RFC 9113 section 6.9.1: reject the update before arithmetic.
1209 goaway_frame goaway(0, static_cast<uint32_t>(error_code::flow_control_error));
1210 send_frame(goaway);
1211 goaway_received_ = true;
1212 is_connected_ = false;
1213 return error_void(static_cast<int>(error_code::flow_control_error),
1214 "Connection flow-control window overflow",
1215 "http2_client::handle_window_update_frame");
1216 }
1217 connection_window_size_ += increment;
1218 }
1219 else
1220 {
1221 auto* stream = get_stream(stream_id);
1222 if (stream)
1223 {
1224 if (stream->window_size > std::numeric_limits<int32_t>::max() - increment)
1225 {
1226 rst_stream_frame reset(stream_id, static_cast<uint32_t>(error_code::flow_control_error));
1227 send_frame(reset);
1228 if (stream->state.exchange(stream_state::closed) != stream_state::closed)
1229 {
1230 if (stream->is_streaming)
1231 {
1232 if (stream->on_complete)
1233 stream->on_complete(static_cast<int>(error_code::flow_control_error));
1234 }
1235 else
1236 {
1237 http2_response response;
1238 response.status_code = 0;
1239 stream->promise.set_value(std::move(response));
1240 }
1241 }
1242 return error_void(static_cast<int>(error_code::flow_control_error),
1243 "Stream flow-control window overflow",
1244 "http2_client::handle_window_update_frame");
1245 }
1246 stream->window_size += increment;
1247 }
1248 }
1249
1250 return ok();
1251 }
1252
1253 auto http2_client::handle_ping_frame(const ping_frame& f) -> VoidResult
1254 {
1255 if (!f.is_ack())
1256 {
1257 // Send PING ACK
1258 ping_frame ack(f.opaque_data(), true);
1259 return send_frame(ack);
1260 }
1261 return ok();
1262 }
1263
1264 auto http2_client::run_io() -> void
1265 {
1266 read_next_frame();
1267 io_context_->run();
1268 }
1269
1270 auto http2_client::stop_io() -> void
1271 {
1272 {
1273 std::lock_guard<std::mutex> lock(io_submission_mutex_);
1274 is_running_ = false;
1275 is_connected_ = false;
1276 if (io_active_ && io_future_.valid())
1277 {
1278 // Cancel on the same executor as the SSL operations. Let run()
1279 // drain their handlers so every waiting writer is completed.
1280 asio::post(*io_context_, [this]() {
1281 std::error_code ec;
1282 socket_->lowest_layer().shutdown(asio::ip::tcp::socket::shutdown_both, ec);
1283 socket_->lowest_layer().close(ec);
1284 work_guard_->reset();
1285 });
1286 }
1287 else
1288 {
1289 if (socket_)
1290 {
1291 std::error_code ec;
1292 socket_->lowest_layer().close(ec);
1293 }
1294 if (work_guard_) work_guard_->reset();
1295 }
1296 }
1297
1298 if (io_future_.valid()) io_future_.wait();
1299 std::lock_guard<std::mutex> lock(io_submission_mutex_);
1300 io_future_ = {};
1301 io_active_ = false;
1302 if (io_context_) io_context_->stop();
1303 }
1304
1305} // 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
~http2_client()
Destructor - closes connection gracefully.
auto cancel_stream(uint32_t stream_id) -> VoidResult
Cancel a stream.
auto is_connected() const -> bool
Check if connected.
auto write_stream(uint32_t stream_id, const std::vector< uint8_t > &data, bool end_stream=false) -> VoidResult
Write data to an open stream.
auto put(const std::string &path, const std::string &body, const std::vector< http_header > &headers={}) -> Result< http2_response >
Perform HTTP/2 PUT request.
auto set_settings(const http2_settings &settings) -> void
Update local settings.
auto get_settings() const -> http2_settings
Get local settings.
http2_client(std::string_view client_id)
Construct HTTP/2 client.
auto send_frame(const frame &f) -> VoidResult
std::atomic< std::chrono::milliseconds > timeout_
auto post(const std::string &path, const std::string &body, const std::vector< http_header > &headers={}) -> Result< http2_response >
Perform HTTP/2 POST request.
auto close_stream_writer(uint32_t stream_id) -> VoidResult
Close the write side of a stream.
auto handle_settings_frame(const settings_frame &frame) -> VoidResult
auto disconnect() -> VoidResult
Disconnect from server.
auto get(const std::string &path, const std::vector< http_header > &headers={}) -> Result< http2_response >
Perform HTTP/2 GET request.
auto del(const std::string &path, const std::vector< http_header > &headers={}) -> Result< http2_response >
Perform HTTP/2 DELETE request.
auto get_timeout() const -> std::chrono::milliseconds
Get current timeout.
auto connect(const std::string &host, unsigned short port=443) -> VoidResult
Connect to HTTP/2 server.
auto start_stream(const std::string &path, const std::vector< http_header > &headers, std::function< void(std::vector< uint8_t >)> on_data, std::function< void(std::vector< http_header >)> on_headers, std::function< void(int)> on_complete) -> Result< uint32_t >
Start a streaming POST request.
auto set_timeout(std::chrono::milliseconds timeout) -> void
Set request timeout.
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
static auto create_span(std::string_view name) -> span
Create a new root span with a new trace context.
struct ssl_ctx_st SSL_CTX
Definition crypto.h:20
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_local
Local end closed, remote can send.
auto is_tracing_enabled() -> bool
Check if tracing is enabled.
@ ok
Operation completed successfully.
std::string get_error_details(const simple_error &err)
VoidResult error_void(int code, const std::string &message, const std::string &source="network_system", const std::string &details="")
VoidResult ok()
RAII span implementation for distributed tracing.
static auto parse(std::span< const uint8_t > data) -> Result< frame_header >
Parse frame header from raw bytes.
Definition frame.cpp:52
std::vector< uint8_t > body
Response body.
auto get_header(const std::string &name) const -> std::optional< std::string >
Get header value by name.
auto get_body_string() const -> std::string
Get body as string.
std::vector< http_header > headers
Response headers.
std::vector< uint8_t > request_body
Request body.
std::atomic< stream_state > state
Current state.
std::promise< http2_response > promise
Response promise.
std::vector< http_header > request_headers
Request headers.
std::function< void(std::vector< uint8_t >)> on_data
Callback for streaming data.
std::function< void(std::vector< http_header >)> on_headers
Callback for headers.
std::function< void(int)> on_complete
Callback when stream ends (status code)
Distributed tracing context for OpenTelemetry-compatible tracing.
Configuration structures for OpenTelemetry tracing.