Network System 0.1.1
High-performance modular networking library for scalable client-server applications
Loading...
Searching...
No Matches
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
8
12
13#include <atomic>
14#include <charconv>
15#include <condition_variable>
16#include <mutex>
17#include <thread>
18
19#if NETWORK_GRPC_OFFICIAL
20#include <grpcpp/grpcpp.h>
21#include <grpcpp/generic/generic_stub.h>
22#include <grpcpp/support/byte_buffer.h>
23#else
26#endif
27
29{
30
31#if NETWORK_GRPC_OFFICIAL
32
33// ============================================================================
34// Official gRPC Library Client Implementation
35// ============================================================================
36
37namespace detail
38{
39
40auto vector_to_byte_buffer(const std::vector<uint8_t>& data) -> ::grpc::ByteBuffer
41{
42 ::grpc::Slice slice(data.data(), data.size());
43 return ::grpc::ByteBuffer(&slice, 1);
44}
45
46auto byte_buffer_to_vector(const ::grpc::ByteBuffer& buffer) -> std::vector<uint8_t>
47{
48 std::vector<::grpc::Slice> slices;
49 buffer.Dump(&slices);
50
51 std::vector<uint8_t> result;
52 result.reserve(buffer.Length());
53
54 for (const auto& slice : slices)
55 {
56 const auto* begin = reinterpret_cast<const uint8_t*>(slice.begin());
57 result.insert(result.end(), begin, begin + slice.size());
58 }
59
60 return result;
61}
62
63} // namespace detail
64
65// Server stream reader implementation for official gRPC
66class official_server_stream_reader : public grpc_client::server_stream_reader
67{
68public:
69 official_server_stream_reader(
70 std::unique_ptr<::grpc::ClientContext> ctx,
71 std::unique_ptr<::grpc::GenericClientAsyncReader> reader,
72 ::grpc::CompletionQueue* cq)
73 : ctx_(std::move(ctx))
74 , reader_(std::move(reader))
75 , cq_(cq)
76 , has_more_(true)
77 {
78 }
79
80 auto read() -> Result<grpc_message> override
81 {
82 if (!has_more_)
83 {
84 return error<grpc_message>(
85 static_cast<int>(status_code::ok),
86 "End of stream",
87 "grpc::client::server_stream_reader");
88 }
89
90 ::grpc::ByteBuffer buffer;
91 void* tag = nullptr;
92 bool ok = false;
93
94 reader_->Read(&buffer, &tag);
95
96 if (!cq_->Next(&tag, &ok) || !ok)
97 {
98 has_more_ = false;
99 return error<grpc_message>(
100 static_cast<int>(status_code::ok),
101 "End of stream",
102 "grpc::client::server_stream_reader");
103 }
104
105 auto data = detail::byte_buffer_to_vector(buffer);
106 return kcenon::network::protocols::grpc::ok(grpc_message{std::move(data)});
107 }
108
109 auto has_more() const -> bool override
110 {
111 return has_more_;
112 }
113
114 auto finish() -> grpc_status override
115 {
116 ::grpc::Status status;
117 void* tag = nullptr;
118 bool ok = false;
119
120 reader_->Finish(&status, &tag);
121 cq_->Next(&tag, &ok);
122
123 return from_grpc_status(status);
124 }
125
126private:
127 std::unique_ptr<::grpc::ClientContext> ctx_;
128 std::unique_ptr<::grpc::GenericClientAsyncReader> reader_;
129 ::grpc::CompletionQueue* cq_;
130 bool has_more_;
131};
132
133// Client stream writer implementation for official gRPC
134class official_client_stream_writer : public grpc_client::client_stream_writer
135{
136public:
137 official_client_stream_writer(
138 std::unique_ptr<::grpc::ClientContext> ctx,
139 std::unique_ptr<::grpc::GenericClientAsyncWriter> writer,
140 ::grpc::CompletionQueue* cq)
141 : ctx_(std::move(ctx))
142 , writer_(std::move(writer))
143 , cq_(cq)
144 , writes_done_(false)
145 {
146 }
147
148 auto write(const std::vector<uint8_t>& message) -> VoidResult override
149 {
150 if (writes_done_)
151 {
152 return error_void(
153 error_codes::common_errors::internal_error,
154 "Stream writes already done",
155 "grpc::client::client_stream_writer");
156 }
157
158 ::grpc::ByteBuffer buffer = detail::vector_to_byte_buffer(message);
159 void* tag = nullptr;
160 bool ok = false;
161
162 writer_->Write(buffer, &tag);
163
164 if (!cq_->Next(&tag, &ok) || !ok)
165 {
166 return error_void(
167 error_codes::network_system::connection_failed,
168 "Failed to write message",
169 "grpc::client::client_stream_writer");
170 }
171
173 }
174
175 auto writes_done() -> VoidResult override
176 {
177 if (writes_done_)
178 {
180 }
181
182 void* tag = nullptr;
183 bool ok = false;
184
185 writer_->WritesDone(&tag);
186 cq_->Next(&tag, &ok);
187
188 writes_done_ = true;
190 }
191
192 auto finish() -> Result<grpc_message> override
193 {
194 if (!writes_done_)
195 {
196 writes_done();
197 }
198
199 ::grpc::ByteBuffer response;
200 ::grpc::Status status;
201 void* tag = nullptr;
202 bool ok = false;
203
204 writer_->Finish(&status, &tag);
205 cq_->Next(&tag, &ok);
206
207 if (!status.ok())
208 {
209 auto grpc_st = from_grpc_status(status);
210 return error<grpc_message>(
211 static_cast<int>(grpc_st.code),
212 grpc_st.message,
213 "grpc::client::client_stream_writer");
214 }
215
216 // Note: Response is not available from writer, need to handle differently
217 return kcenon::network::protocols::grpc::ok(grpc_message{});
218 }
219
220private:
221 std::unique_ptr<::grpc::ClientContext> ctx_;
222 std::unique_ptr<::grpc::GenericClientAsyncWriter> writer_;
223 ::grpc::CompletionQueue* cq_;
224 bool writes_done_;
225};
226
227// Bidirectional stream implementation for official gRPC
228class official_bidi_stream : public grpc_client::bidi_stream
229{
230public:
231 official_bidi_stream(
232 std::unique_ptr<::grpc::ClientContext> ctx,
233 std::unique_ptr<::grpc::GenericClientAsyncReaderWriter> stream,
234 ::grpc::CompletionQueue* cq)
235 : ctx_(std::move(ctx))
236 , stream_(std::move(stream))
237 , cq_(cq)
238 , writes_done_(false)
239 {
240 }
241
242 auto write(const std::vector<uint8_t>& message) -> VoidResult override
243 {
244 if (writes_done_)
245 {
246 return error_void(
247 error_codes::common_errors::internal_error,
248 "Stream writes already done",
249 "grpc::client::bidi_stream");
250 }
251
252 ::grpc::ByteBuffer buffer = detail::vector_to_byte_buffer(message);
253 void* tag = nullptr;
254 bool ok = false;
255
256 stream_->Write(buffer, &tag);
257
258 if (!cq_->Next(&tag, &ok) || !ok)
259 {
260 return error_void(
261 error_codes::network_system::connection_failed,
262 "Failed to write message",
263 "grpc::client::bidi_stream");
264 }
265
267 }
268
269 auto read() -> Result<grpc_message> override
270 {
271 ::grpc::ByteBuffer buffer;
272 void* tag = nullptr;
273 bool ok = false;
274
275 stream_->Read(&buffer, &tag);
276
277 if (!cq_->Next(&tag, &ok) || !ok)
278 {
279 return error<grpc_message>(
280 static_cast<int>(status_code::ok),
281 "End of stream",
282 "grpc::client::bidi_stream");
283 }
284
285 auto data = detail::byte_buffer_to_vector(buffer);
286 return kcenon::network::protocols::grpc::ok(grpc_message{std::move(data)});
287 }
288
289 auto writes_done() -> VoidResult override
290 {
291 if (writes_done_)
292 {
294 }
295
296 void* tag = nullptr;
297 bool ok = false;
298
299 stream_->WritesDone(&tag);
300 cq_->Next(&tag, &ok);
301
302 writes_done_ = true;
304 }
305
306 auto finish() -> grpc_status override
307 {
308 if (!writes_done_)
309 {
310 writes_done();
311 }
312
313 ::grpc::Status status;
314 void* tag = nullptr;
315 bool ok = false;
316
317 stream_->Finish(&status, &tag);
318 cq_->Next(&tag, &ok);
319
320 return from_grpc_status(status);
321 }
322
323private:
324 std::unique_ptr<::grpc::ClientContext> ctx_;
325 std::unique_ptr<::grpc::GenericClientAsyncReaderWriter> stream_;
326 ::grpc::CompletionQueue* cq_;
327 bool writes_done_;
328};
329
330// Implementation class using official gRPC
331class grpc_client::impl : public std::enable_shared_from_this<grpc_client::impl>
332{
333public:
334 explicit impl(std::string target, grpc_channel_config config)
335 : target_(std::move(target))
336 , config_(std::move(config))
337 , connected_(false)
338 {
339 }
340
341 ~impl()
342 {
343 disconnect();
344 }
345
346 auto connect() -> VoidResult
347 {
348 auto span = tracing::is_tracing_enabled()
349 ? std::make_optional(tracing::trace_context::create_span("grpc.client.connect"))
350 : std::nullopt;
351 if (span)
352 {
353 span->set_attribute("rpc.system", "grpc")
354 .set_attribute("rpc.grpc.target", target_)
355 .set_attribute("net.transport", "tcp")
356 .set_attribute("rpc.grpc.use_tls", config_.use_tls);
357 }
358
359 std::lock_guard<std::mutex> lock(mutex_);
360
361 if (connected_.load())
362 {
363 if (span)
364 {
365 span->set_attribute("grpc.client.already_connected", true);
366 }
368 }
369
370 // Create channel credentials
371 channel_credentials_config creds_config;
372 creds_config.insecure = !config_.use_tls;
373
374 if (config_.use_tls)
375 {
376 creds_config.root_certificates = config_.root_certificates;
377 creds_config.client_certificate = config_.client_certificate;
378 creds_config.client_key = config_.client_key;
379 }
380
381 channel_ = create_channel(target_, creds_config);
382
383 if (!channel_)
384 {
385 if (span)
386 {
387 span->set_error("Failed to create gRPC channel");
388 }
389 return error_void(
391 "Failed to create gRPC channel",
392 "grpc::client");
393 }
394
395 // Wait for channel to be ready
396 if (!wait_for_channel_ready(channel_, config_.default_timeout))
397 {
398 if (span)
399 {
400 span->set_error("Failed to connect to gRPC server");
401 }
402 return error_void(
404 "Failed to connect to gRPC server",
405 "grpc::client",
406 target_);
407 }
408
409 stub_ = std::make_unique<::grpc::GenericStub>(channel_);
410 connected_.store(true);
411
413 }
414
415 auto disconnect() -> void
416 {
417 std::lock_guard<std::mutex> lock(mutex_);
418
419 stub_.reset();
420 channel_.reset();
421 connected_.store(false);
422 }
423
424 auto is_connected() const -> bool
425 {
426 if (!connected_.load() || !channel_)
427 {
428 return false;
429 }
430
431 auto state = channel_->GetState(false);
432 return state == GRPC_CHANNEL_READY || state == GRPC_CHANNEL_IDLE;
433 }
434
435 auto wait_for_connected(std::chrono::milliseconds timeout) -> bool
436 {
437 if (!channel_)
438 {
439 return false;
440 }
441
442 return wait_for_channel_ready(channel_, timeout);
443 }
444
445 auto target() const -> const std::string&
446 {
447 return target_;
448 }
449
450 auto call_raw(const std::string& method,
451 const std::vector<uint8_t>& request,
452 const call_options& options) -> Result<grpc_message>
453 {
454 auto span = tracing::is_tracing_enabled()
455 ? std::make_optional(tracing::trace_context::create_span("grpc.client.call"))
456 : std::nullopt;
457 if (span)
458 {
459 span->set_attribute("rpc.system", "grpc")
460 .set_attribute("rpc.method", method)
461 .set_attribute("rpc.grpc.target", target_)
462 .set_attribute("rpc.request.size", static_cast<int64_t>(request.size()));
463 }
464
465 if (!is_connected())
466 {
467 if (span)
468 {
469 span->set_error("Not connected to server");
470 }
471 return error<grpc_message>(
473 "Not connected to server",
474 "grpc::client");
475 }
476
477 ::grpc::ClientContext ctx;
478
479 // Set deadline if provided
480 if (options.deadline.has_value())
481 {
482 set_deadline(&ctx, options.deadline.value());
483 }
484 else
485 {
486 // Use default timeout
488 }
489
490 // Add metadata
491 for (const auto& [key, value] : options.metadata)
492 {
493 ctx.AddMetadata(key, value);
494 }
495
496 // Set wait for ready
497 if (options.wait_for_ready)
498 {
499 ctx.set_wait_for_ready(true);
500 }
501
502 // Prepare request
503 ::grpc::ByteBuffer request_buffer = detail::vector_to_byte_buffer(request);
504 ::grpc::ByteBuffer response_buffer;
505
506 // Make unary call
507 ::grpc::Status status = stub_->UnaryCall(&ctx, method, request_buffer, &response_buffer);
508
509 if (!status.ok())
510 {
511 auto grpc_st = from_grpc_status(status);
512 if (span)
513 {
514 span->set_error(grpc_st.message)
515 .set_attribute("rpc.grpc.status_code", static_cast<int64_t>(grpc_st.code));
516 }
517 return error<grpc_message>(
518 static_cast<int>(grpc_st.code),
519 grpc_st.message,
520 "grpc::client");
521 }
522
523 auto response_data = detail::byte_buffer_to_vector(response_buffer);
524 if (span)
525 {
526 span->set_attribute("rpc.response.size", static_cast<int64_t>(response_data.size()))
527 .set_attribute("rpc.grpc.status_code", static_cast<int64_t>(0));
528 }
529 return kcenon::network::protocols::grpc::ok(grpc_message{std::move(response_data)});
530 }
531
532 auto call_raw_async(const std::string& method,
533 const std::vector<uint8_t>& request,
534 std::function<void(Result<grpc_message>)> callback,
535 const call_options& options) -> void
536 {
537 // Create async call in separate thread
538 std::thread([self = shared_from_this(), method, request, callback, options]() {
539 auto result = self->call_raw(method, request, options);
540 if (callback)
541 {
542 callback(std::move(result));
543 }
544 }).detach();
545 }
546
547 auto server_stream_raw(const std::string& method,
548 const std::vector<uint8_t>& request,
549 const call_options& options)
550 -> Result<std::unique_ptr<grpc_client::server_stream_reader>>
551 {
552 if (!is_connected())
553 {
556 "Not connected to server",
557 "grpc::client");
558 }
559
560 auto ctx = std::make_unique<::grpc::ClientContext>();
561
562 // Set deadline if provided
563 if (options.deadline.has_value())
564 {
565 set_deadline(ctx.get(), options.deadline.value());
566 }
567 else
568 {
570 }
571
572 // Add metadata
573 for (const auto& [key, value] : options.metadata)
574 {
575 ctx->AddMetadata(key, value);
576 }
577
578 // Note: GenericStub doesn't have PrepareUnaryCall with streaming
579 // For full streaming support, we would need to use async API
580 // This is a simplified synchronous implementation
581
583 static_cast<int>(status_code::unimplemented),
584 "Server streaming not fully implemented for official gRPC wrapper",
585 "grpc::client");
586 }
587
588 auto client_stream_raw(const std::string& method,
589 const call_options& options)
590 -> Result<std::unique_ptr<grpc_client::client_stream_writer>>
591 {
592 if (!is_connected())
593 {
596 "Not connected to server",
597 "grpc::client");
598 }
599
601 static_cast<int>(status_code::unimplemented),
602 "Client streaming not fully implemented for official gRPC wrapper",
603 "grpc::client");
604 }
605
606 auto bidi_stream_raw(const std::string& method,
607 const call_options& options)
608 -> Result<std::unique_ptr<grpc_client::bidi_stream>>
609 {
610 if (!is_connected())
611 {
614 "Not connected to server",
615 "grpc::client");
616 }
617
619 static_cast<int>(status_code::unimplemented),
620 "Bidirectional streaming not fully implemented for official gRPC wrapper",
621 "grpc::client");
622 }
623
624private:
625 std::string target_;
626 grpc_channel_config config_;
627 std::shared_ptr<::grpc::Channel> channel_;
628 std::unique_ptr<::grpc::GenericStub> stub_;
629 std::atomic<bool> connected_;
630 mutable std::mutex mutex_;
631};
632
633#else // !NETWORK_GRPC_OFFICIAL
634
635// ============================================================================
636// Prototype Implementation (HTTP/2 Transport)
637// ============================================================================
638
639// Server stream reader implementation
641{
642public:
643 server_stream_reader_impl(std::shared_ptr<http2::http2_client> http2_client,
644 uint32_t stream_id)
645 : http2_client_(std::move(http2_client))
646 , stream_id_(stream_id)
647 , has_more_(true)
648 {
649 }
650
651 auto read() -> Result<grpc_message> override
652 {
653 std::unique_lock<std::mutex> lock(mutex_);
654
655 // Wait for data or end of stream
656 cv_.wait(lock, [this] { return !buffer_.empty() || !has_more_; });
657
658 if (buffer_.empty())
659 {
660 if (!has_more_)
661 {
662 return error<grpc_message>(
663 static_cast<int>(status_code::ok),
664 "End of stream",
665 "server_stream_reader::read");
666 }
667 }
668
669 // Parse gRPC message from buffer
670 auto parse_result = grpc_message::parse(buffer_);
671 if (parse_result.is_err())
672 {
673 return parse_result;
674 }
675
676 // Remove parsed data from buffer
677 auto& msg = parse_result.value();
678 size_t msg_size = 5 + msg.data.size(); // 1 byte compression + 4 bytes length + data
679 if (buffer_.size() >= msg_size)
680 {
681 buffer_.erase(buffer_.begin(), buffer_.begin() + static_cast<ptrdiff_t>(msg_size));
682 }
683
684 return ok(std::move(msg));
685 }
686
687 auto has_more() const -> bool override
688 {
689 std::lock_guard<std::mutex> lock(mutex_);
690 return has_more_ || !buffer_.empty();
691 }
692
693 auto finish() -> grpc_status override
694 {
695 std::unique_lock<std::mutex> lock(mutex_);
696 cv_.wait(lock, [this] { return !has_more_; });
697 return final_status_;
698 }
699
700 void on_data(const std::vector<uint8_t>& data)
701 {
702 std::lock_guard<std::mutex> lock(mutex_);
703 buffer_.insert(buffer_.end(), data.begin(), data.end());
704 cv_.notify_one();
705 }
706
707 void on_headers(const std::vector<http2::http_header>& headers)
708 {
709 std::lock_guard<std::mutex> lock(mutex_);
710 for (const auto& h : headers)
711 {
712 if (h.name == "grpc-status")
713 {
714 try { final_status_.code = static_cast<status_code>(std::stoi(h.value)); }
715 catch (...) {}
716 }
717 else if (h.name == "grpc-message")
718 {
719 final_status_.message = h.value;
720 }
721 }
722 }
723
725 {
726 std::lock_guard<std::mutex> lock(mutex_);
727 has_more_ = false;
728 if (status_code != 200)
729 {
731 final_status_.message = "HTTP error: " + std::to_string(status_code);
732 }
733 cv_.notify_all();
734 }
735
736private:
737 std::shared_ptr<http2::http2_client> http2_client_;
738 [[maybe_unused]] uint32_t stream_id_; // Reserved for future stream operations
739 mutable std::mutex mutex_;
740 std::condition_variable cv_;
741 std::vector<uint8_t> buffer_;
744};
745
746// Client stream writer implementation
748{
749public:
750 client_stream_writer_impl(std::shared_ptr<http2::http2_client> http2_client,
751 uint32_t stream_id)
752 : http2_client_(std::move(http2_client))
753 , stream_id_(stream_id)
754 , writes_done_(false)
755 {
756 }
757
758 void set_stream_id(uint32_t stream_id) { stream_id_ = stream_id; }
759
760 auto write(const std::vector<uint8_t>& message) -> VoidResult override
761 {
762 if (writes_done_)
763 {
764 return error_void(
766 "Stream writes already done",
767 "client_stream_writer::write");
768 }
769
770 // Serialize with gRPC framing
771 grpc_message msg{std::vector<uint8_t>(message)};
772 auto serialized = msg.serialize();
773
774 return http2_client_->write_stream(stream_id_, serialized, false);
775 }
776
777 auto writes_done() -> VoidResult override
778 {
779 if (writes_done_)
780 {
781 return ok();
782 }
783
784 writes_done_ = true;
785 return http2_client_->close_stream_writer(stream_id_);
786 }
787
788 auto finish() -> Result<grpc_message> override
789 {
790 // Close write side if not already done
791 if (!writes_done_)
792 {
793 auto wd_result = writes_done();
794 if (wd_result.is_err())
795 {
796 const auto& err = wd_result.error();
797 return error<grpc_message>(err.code, err.message, "client_stream_writer::finish");
798 }
799 }
800
801 // Wait for response
802 std::unique_lock<std::mutex> lock(mutex_);
803 cv_.wait(lock, [this] { return response_ready_; });
804
806 {
807 return ok(std::move(response_));
808 }
809
810 return error<grpc_message>(
811 static_cast<int>(final_status_.code),
813 "client_stream_writer::finish");
814 }
815
816 void on_data(const std::vector<uint8_t>& data)
817 {
818 std::lock_guard<std::mutex> lock(mutex_);
819 response_buffer_.insert(response_buffer_.end(), data.begin(), data.end());
820 }
821
822 void on_headers(const std::vector<http2::http_header>& headers)
823 {
824 std::lock_guard<std::mutex> lock(mutex_);
825 for (const auto& h : headers)
826 {
827 if (h.name == "grpc-status")
828 {
829 try { final_status_.code = static_cast<status_code>(std::stoi(h.value)); }
830 catch (...) {}
831 }
832 else if (h.name == "grpc-message")
833 {
834 final_status_.message = h.value;
835 }
836 }
837 }
838
840 {
841 std::lock_guard<std::mutex> lock(mutex_);
842
843 if (status_code != 200)
844 {
846 final_status_.message = "HTTP error: " + std::to_string(status_code);
847 }
848
849 // Parse response
850 if (!response_buffer_.empty())
851 {
852 auto parse_result = grpc_message::parse(response_buffer_);
853 if (parse_result.is_ok())
854 {
855 response_ = std::move(parse_result.value());
856 }
857 }
858
859 response_ready_ = true;
860 cv_.notify_all();
861 }
862
863private:
864 std::shared_ptr<http2::http2_client> http2_client_;
865 uint32_t stream_id_;
867 mutable std::mutex mutex_;
868 std::condition_variable cv_;
869 std::vector<uint8_t> response_buffer_;
872 bool response_ready_ = false;
873};
874
875// Bidirectional stream implementation
877{
878public:
879 bidi_stream_impl(std::shared_ptr<http2::http2_client> http2_client,
880 uint32_t stream_id)
881 : http2_client_(std::move(http2_client))
882 , stream_id_(stream_id)
883 , writes_done_(false)
884 , stream_ended_(false)
885 {
886 }
887
888 void set_stream_id(uint32_t stream_id) { stream_id_ = stream_id; }
889
890 auto write(const std::vector<uint8_t>& message) -> VoidResult override
891 {
892 if (writes_done_)
893 {
894 return error_void(
896 "Stream writes already done",
897 "bidi_stream::write");
898 }
899
900 grpc_message msg{std::vector<uint8_t>(message)};
901 auto serialized = msg.serialize();
902
903 return http2_client_->write_stream(stream_id_, serialized, false);
904 }
905
906 auto read() -> Result<grpc_message> override
907 {
908 std::unique_lock<std::mutex> lock(mutex_);
909
910 cv_.wait(lock, [this] { return !buffer_.empty() || stream_ended_; });
911
912 if (buffer_.empty())
913 {
914 if (stream_ended_)
915 {
916 return error<grpc_message>(
917 static_cast<int>(status_code::ok),
918 "End of stream",
919 "bidi_stream::read");
920 }
921 }
922
923 auto parse_result = grpc_message::parse(buffer_);
924 if (parse_result.is_err())
925 {
926 return parse_result;
927 }
928
929 auto& msg = parse_result.value();
930 size_t msg_size = 5 + msg.data.size();
931 if (buffer_.size() >= msg_size)
932 {
933 buffer_.erase(buffer_.begin(), buffer_.begin() + static_cast<ptrdiff_t>(msg_size));
934 }
935
936 return ok(std::move(msg));
937 }
938
939 auto writes_done() -> VoidResult override
940 {
941 if (writes_done_)
942 {
943 return ok();
944 }
945
946 writes_done_ = true;
947 return http2_client_->close_stream_writer(stream_id_);
948 }
949
950 auto finish() -> grpc_status override
951 {
952 if (!writes_done_)
953 {
954 writes_done();
955 }
956
957 std::unique_lock<std::mutex> lock(mutex_);
958 cv_.wait(lock, [this] { return stream_ended_; });
959 return final_status_;
960 }
961
962 void on_data(const std::vector<uint8_t>& data)
963 {
964 std::lock_guard<std::mutex> lock(mutex_);
965 buffer_.insert(buffer_.end(), data.begin(), data.end());
966 cv_.notify_one();
967 }
968
969 void on_headers(const std::vector<http2::http_header>& headers)
970 {
971 std::lock_guard<std::mutex> lock(mutex_);
972 for (const auto& h : headers)
973 {
974 if (h.name == "grpc-status")
975 {
976 try { final_status_.code = static_cast<status_code>(std::stoi(h.value)); }
977 catch (...) {}
978 }
979 else if (h.name == "grpc-message")
980 {
981 final_status_.message = h.value;
982 }
983 }
984 }
985
987 {
988 std::lock_guard<std::mutex> lock(mutex_);
989 stream_ended_ = true;
990 if (status_code != 200)
991 {
993 final_status_.message = "HTTP error: " + std::to_string(status_code);
994 }
995 cv_.notify_all();
996 }
997
998private:
999 std::shared_ptr<http2::http2_client> http2_client_;
1000 uint32_t stream_id_;
1003 mutable std::mutex mutex_;
1004 std::condition_variable cv_;
1005 std::vector<uint8_t> buffer_;
1007};
1008
1009// Implementation class using HTTP/2 transport
1010class grpc_client::impl : public std::enable_shared_from_this<grpc_client::impl>
1011{
1012public:
1013 explicit impl(std::string target, grpc_channel_config config)
1014 : target_(std::move(target))
1015 , config_(std::move(config))
1016 , connected_(false)
1017 {
1018 }
1019
1021 {
1022 disconnect();
1023 }
1024
1026 {
1027 auto span = tracing::is_tracing_enabled()
1028 ? std::make_optional(tracing::trace_context::create_span("grpc.client.connect"))
1029 : std::nullopt;
1030 if (span)
1031 {
1032 span->set_attribute("rpc.system", "grpc")
1033 .set_attribute("rpc.grpc.target", target_)
1034 .set_attribute("net.transport", "tcp");
1035 }
1036
1037 std::lock_guard<std::mutex> lock(mutex_);
1038
1039 if (connected_.load())
1040 {
1041 if (span)
1042 {
1043 span->set_attribute("grpc.client.already_connected", true);
1044 }
1045 return ok();
1046 }
1047
1048 // Parse target address
1049 auto colon_pos = target_.find(':');
1050 if (colon_pos == std::string::npos)
1051 {
1052 if (span)
1053 {
1054 span->set_error("Invalid target address format");
1055 }
1056 return error_void(
1058 "Invalid target address format",
1059 "grpc::client",
1060 "Expected format: host:port");
1061 }
1062
1063 host_ = target_.substr(0, colon_pos);
1064 auto port_str = target_.substr(colon_pos + 1);
1065
1066 unsigned short port = 0;
1067 auto [ptr, ec] = std::from_chars(
1068 port_str.data(),
1069 port_str.data() + port_str.size(),
1070 port);
1071
1072 if (ec != std::errc())
1073 {
1074 if (span)
1075 {
1076 span->set_error("Invalid port number");
1077 }
1078 return error_void(
1080 "Invalid port number",
1081 "grpc::client");
1082 }
1083
1084 if (span)
1085 {
1086 span->set_attribute("net.peer.name", host_)
1087 .set_attribute("net.peer.port", static_cast<int64_t>(port));
1088 }
1089
1090 // Create HTTP/2 client
1091 http2_client_ = std::make_shared<http2::http2_client>("grpc-client");
1093
1094 // Connect using HTTP/2
1095 auto result = http2_client_->connect(host_, port);
1096 if (result.is_err())
1097 {
1098 const auto& err = result.error();
1099 if (span)
1100 {
1101 span->set_error(err.message);
1102 }
1103 return error_void(err.code, err.message, "grpc::client",
1104 get_error_details(err));
1105 }
1106
1107 connected_.store(true);
1108 return ok();
1109 }
1110
1111 auto disconnect() -> void
1112 {
1113 std::lock_guard<std::mutex> lock(mutex_);
1114
1115 if (http2_client_ && connected_.load())
1116 {
1117 http2_client_->disconnect();
1118 }
1119
1120 connected_.store(false);
1121 http2_client_.reset();
1122 }
1123
1124 auto is_connected() const -> bool
1125 {
1126 return connected_.load() && http2_client_ && http2_client_->is_connected();
1127 }
1128
1129 auto wait_for_connected(std::chrono::milliseconds timeout) -> bool
1130 {
1131 auto deadline = std::chrono::steady_clock::now() + timeout;
1132
1133 while (std::chrono::steady_clock::now() < deadline)
1134 {
1135 if (is_connected())
1136 {
1137 return true;
1138 }
1139 std::this_thread::sleep_for(std::chrono::milliseconds(10));
1140 }
1141
1142 return is_connected();
1143 }
1144
1145 auto target() const -> const std::string&
1146 {
1147 return target_;
1148 }
1149
1150 auto call_raw(const std::string& method,
1151 const std::vector<uint8_t>& request,
1152 const call_options& options) -> Result<grpc_message>
1153 {
1154 auto span = tracing::is_tracing_enabled()
1155 ? std::make_optional(tracing::trace_context::create_span("grpc.client.call"))
1156 : std::nullopt;
1157 if (span)
1158 {
1159 span->set_attribute("rpc.system", "grpc")
1160 .set_attribute("rpc.method", method)
1161 .set_attribute("rpc.grpc.target", target_)
1162 .set_attribute("rpc.request.size", static_cast<int64_t>(request.size()));
1163 }
1164
1165 if (!is_connected())
1166 {
1167 if (span)
1168 {
1169 span->set_error("Not connected to server");
1170 }
1171 return error<grpc_message>(
1173 "Not connected to server",
1174 "grpc::client");
1175 }
1176
1177 // Validate method format
1178 if (method.empty() || method[0] != '/')
1179 {
1180 if (span)
1181 {
1182 span->set_error("Invalid method format");
1183 }
1184 return error<grpc_message>(
1186 "Invalid method format",
1187 "grpc::client",
1188 "Method must start with '/'");
1189 }
1190
1191 // Check deadline
1192 if (options.deadline.has_value())
1193 {
1194 if (std::chrono::system_clock::now() > options.deadline.value())
1195 {
1196 if (span)
1197 {
1198 span->set_error("Deadline exceeded before call started")
1199 .set_attribute("rpc.grpc.status_code", static_cast<int64_t>(status_code::deadline_exceeded));
1200 }
1201 return error<grpc_message>(
1202 static_cast<int>(status_code::deadline_exceeded),
1203 "Deadline exceeded before call started",
1204 "grpc::client");
1205 }
1206 }
1207
1208 // Build gRPC headers
1209 std::vector<http2::http_header> headers;
1210 headers.emplace_back(header_names::content_type, grpc_content_type);
1211 headers.emplace_back(header_names::te, "trailers");
1212 headers.emplace_back(header_names::grpc_accept_encoding,
1213 std::string(compression::identity) + "," +
1215
1216 // Add timeout header if deadline is set
1217 if (options.deadline.has_value())
1218 {
1219 auto now = std::chrono::system_clock::now();
1220 auto remaining = std::chrono::duration_cast<std::chrono::milliseconds>(
1221 options.deadline.value() - now);
1222 if (remaining.count() > 0)
1223 {
1224 headers.emplace_back(header_names::grpc_timeout,
1225 format_timeout(static_cast<uint64_t>(remaining.count())));
1226 }
1227 }
1228
1229 // Add custom metadata
1230 for (const auto& [key, value] : options.metadata)
1231 {
1232 headers.emplace_back(key, value);
1233 }
1234
1235 // Serialize request with gRPC framing
1236 grpc_message request_msg{std::vector<uint8_t>{request.begin(), request.end()}};
1237 auto serialized_request = request_msg.serialize();
1238
1239 // Send HTTP/2 POST request
1240 auto response_result = http2_client_->post(method, serialized_request, headers);
1241 if (response_result.is_err())
1242 {
1243 const auto& err = response_result.error();
1244 if (span)
1245 {
1246 span->set_error(err.message);
1247 }
1248 return error<grpc_message>(err.code, err.message, "grpc::client",
1249 get_error_details(err));
1250 }
1251
1252 const auto& response = response_result.value();
1253
1254 // Check HTTP status - gRPC always uses 200 OK for successful transport
1255 if (response.status_code != 200)
1256 {
1257 if (span)
1258 {
1259 span->set_error("HTTP error: " + std::to_string(response.status_code))
1260 .set_attribute("http.status_code", static_cast<int64_t>(response.status_code));
1261 }
1262 return error<grpc_message>(
1263 static_cast<int>(status_code::unavailable),
1264 "HTTP error: " + std::to_string(response.status_code),
1265 "grpc::client");
1266 }
1267
1268 // Extract gRPC status from trailers/headers
1270 std::string grpc_message_str;
1271
1272 for (const auto& header : response.headers)
1273 {
1274 if (header.name == trailer_names::grpc_status)
1275 {
1276 int status_int = 0;
1277 auto [ptr, ec] = std::from_chars(
1278 header.value.data(),
1279 header.value.data() + header.value.size(),
1280 status_int);
1281 if (ec == std::errc())
1282 {
1283 grpc_status = static_cast<status_code>(status_int);
1284 }
1285 }
1286 else if (header.name == trailer_names::grpc_message)
1287 {
1288 grpc_message_str = header.value;
1289 }
1290 }
1291
1292 // Check gRPC status
1294 {
1295 if (span)
1296 {
1297 span->set_error(grpc_message_str.empty() ?
1298 std::string(status_code_to_string(grpc_status)) : grpc_message_str)
1299 .set_attribute("rpc.grpc.status_code", static_cast<int64_t>(grpc_status));
1300 }
1301 return error<grpc_message>(
1302 static_cast<int>(grpc_status),
1303 grpc_message_str.empty() ?
1304 std::string(status_code_to_string(grpc_status)) : grpc_message_str,
1305 "grpc::client");
1306 }
1307
1308 // Parse response message
1309 if (response.body.empty())
1310 {
1311 // Empty response is valid for some RPCs
1312 if (span)
1313 {
1314 span->set_attribute("rpc.response.size", static_cast<int64_t>(0))
1315 .set_attribute("rpc.grpc.status_code", static_cast<int64_t>(0));
1316 }
1317 return ok(grpc_message{});
1318 }
1319
1320 auto parse_result = grpc_message::parse(response.body);
1321 if (parse_result.is_err())
1322 {
1323 const auto& err = parse_result.error();
1324 if (span)
1325 {
1326 span->set_error(err.message);
1327 }
1328 return error<grpc_message>(err.code, err.message, "grpc::client",
1329 get_error_details(err));
1330 }
1331
1332 if (span)
1333 {
1334 span->set_attribute("rpc.response.size", static_cast<int64_t>(parse_result.value().data.size()))
1335 .set_attribute("rpc.grpc.status_code", static_cast<int64_t>(0));
1336 }
1337 return ok(std::move(parse_result.value()));
1338 }
1339
1340 auto call_raw_async(const std::string& method,
1341 const std::vector<uint8_t>& request,
1342 std::function<void(Result<grpc_message>)> callback,
1343 const call_options& options) -> void
1344 {
1345 // Execute asynchronously using thread pool
1347 [self = shared_from_this(), method, request, callback, options]() {
1348 auto result = self->call_raw(method, request, options);
1349 if (callback)
1350 {
1351 callback(std::move(result));
1352 }
1353 });
1354 }
1355
1356 auto server_stream_raw(const std::string& method,
1357 const std::vector<uint8_t>& request,
1358 const call_options& options)
1360 {
1361 if (!is_connected())
1362 {
1365 "Not connected to server",
1366 "grpc::client");
1367 }
1368
1369 // Validate method format
1370 if (method.empty() || method[0] != '/')
1371 {
1374 "Invalid method format",
1375 "grpc::client",
1376 "Method must start with '/'");
1377 }
1378
1379 // Build gRPC headers
1380 std::vector<http2::http_header> headers;
1381 headers.emplace_back("content-type", std::string(grpc_content_type));
1382 headers.emplace_back("te", "trailers");
1383 headers.emplace_back("grpc-accept-encoding",
1384 std::string(compression::identity) + "," +
1386
1387 // Add timeout header if deadline is set
1388 if (options.deadline.has_value())
1389 {
1390 auto now = std::chrono::system_clock::now();
1391 auto remaining = std::chrono::duration_cast<std::chrono::milliseconds>(
1392 options.deadline.value() - now);
1393 if (remaining.count() > 0)
1394 {
1395 headers.emplace_back("grpc-timeout",
1396 format_timeout(static_cast<uint64_t>(remaining.count())));
1397 }
1398 }
1399
1400 // Add custom metadata
1401 for (const auto& [key, value] : options.metadata)
1402 {
1403 headers.emplace_back(key, value);
1404 }
1405
1406 // Callbacks observe handles weakly: a handle owns its transport, so
1407 // transport -> callback -> handle must not create an ownership cycle.
1408 // Keep the handle alive locally until it is returned to the caller.
1409 auto reader = std::make_shared<server_stream_reader_impl>(http2_client_, 0);
1410
1411 // Start streaming request
1412 auto stream_result = http2_client_->start_stream(
1413 method,
1414 headers,
1415 [weak = std::weak_ptr{reader}](std::vector<uint8_t> data) {
1416 if (auto stream = weak.lock()) stream->on_data(data);
1417 },
1418 [weak = std::weak_ptr{reader}](std::vector<http2::http_header> hdrs) {
1419 if (auto stream = weak.lock()) stream->on_headers(hdrs);
1420 },
1421 [weak = std::weak_ptr{reader}](int status) {
1422 if (auto stream = weak.lock()) stream->on_complete(status);
1423 });
1424
1425 if (stream_result.is_err())
1426 {
1427 const auto& err = stream_result.error();
1429 err.code, err.message, "grpc::client", get_error_details(err));
1430 }
1431
1432 // Send the request with gRPC framing
1433 grpc_message request_msg{std::vector<uint8_t>(request)};
1434 auto serialized = request_msg.serialize();
1435 auto write_result = http2_client_->write_stream(stream_result.value(), serialized, true);
1436
1437 if (write_result.is_err())
1438 {
1439 const auto& err = write_result.error();
1441 err.code, err.message, "grpc::client", get_error_details(err));
1442 }
1443
1444 // Return the reader wrapped in a shared_ptr-owning unique_ptr wrapper
1445 // Use a custom deleter that captures the shared_ptr to extend lifetime
1446 struct shared_holder : public grpc_client::server_stream_reader {
1447 std::shared_ptr<server_stream_reader_impl> impl;
1448 explicit shared_holder(std::shared_ptr<server_stream_reader_impl> p) : impl(std::move(p)) {}
1449 auto read() -> Result<grpc_message> override { return impl->read(); }
1450 auto has_more() const -> bool override { return impl->has_more(); }
1451 auto finish() -> grpc_status override { return impl->finish(); }
1452 };
1453
1454 return ok(std::unique_ptr<grpc_client::server_stream_reader>(
1455 new shared_holder(reader)));
1456 }
1457
1458 auto client_stream_raw(const std::string& method,
1459 const call_options& options)
1461 {
1462 if (!is_connected())
1463 {
1466 "Not connected to server",
1467 "grpc::client");
1468 }
1469
1470 // Validate method format
1471 if (method.empty() || method[0] != '/')
1472 {
1475 "Invalid method format",
1476 "grpc::client",
1477 "Method must start with '/'");
1478 }
1479
1480 // Build gRPC headers
1481 std::vector<http2::http_header> headers;
1482 headers.emplace_back("content-type", std::string(grpc_content_type));
1483 headers.emplace_back("te", "trailers");
1484 headers.emplace_back("grpc-accept-encoding",
1485 std::string(compression::identity) + "," +
1487
1488 // Add timeout header if deadline is set
1489 if (options.deadline.has_value())
1490 {
1491 auto now = std::chrono::system_clock::now();
1492 auto remaining = std::chrono::duration_cast<std::chrono::milliseconds>(
1493 options.deadline.value() - now);
1494 if (remaining.count() > 0)
1495 {
1496 headers.emplace_back("grpc-timeout",
1497 format_timeout(static_cast<uint64_t>(remaining.count())));
1498 }
1499 }
1500
1501 // Add custom metadata
1502 for (const auto& [key, value] : options.metadata)
1503 {
1504 headers.emplace_back(key, value);
1505 }
1506
1507 // Weak callback captures avoid a transport/stream ownership cycle.
1508 auto writer = std::make_shared<client_stream_writer_impl>(http2_client_, 0);
1509
1510 // Start streaming request
1511 auto stream_result = http2_client_->start_stream(
1512 method,
1513 headers,
1514 [weak = std::weak_ptr{writer}](std::vector<uint8_t> data) {
1515 if (auto stream = weak.lock()) stream->on_data(data);
1516 },
1517 [weak = std::weak_ptr{writer}](std::vector<http2::http_header> hdrs) {
1518 if (auto stream = weak.lock()) stream->on_headers(hdrs);
1519 },
1520 [weak = std::weak_ptr{writer}](int status) {
1521 if (auto stream = weak.lock()) stream->on_complete(status);
1522 });
1523
1524 if (stream_result.is_err())
1525 {
1526 const auto& err = stream_result.error();
1528 err.code, err.message, "grpc::client", get_error_details(err));
1529 }
1530
1531 // Update the writer with actual stream ID
1532 writer->set_stream_id(stream_result.value());
1533
1534 struct shared_writer_holder : public grpc_client::client_stream_writer {
1535 std::shared_ptr<client_stream_writer_impl> impl;
1536 explicit shared_writer_holder(std::shared_ptr<client_stream_writer_impl> p) : impl(std::move(p)) {}
1537 auto write(const std::vector<uint8_t>& message) -> VoidResult override { return impl->write(message); }
1538 auto writes_done() -> VoidResult override { return impl->writes_done(); }
1539 auto finish() -> Result<grpc_message> override { return impl->finish(); }
1540 };
1541
1542 return ok(std::unique_ptr<grpc_client::client_stream_writer>(
1543 new shared_writer_holder(writer)));
1544 }
1545
1546 auto bidi_stream_raw(const std::string& method,
1547 const call_options& options)
1549 {
1550 if (!is_connected())
1551 {
1554 "Not connected to server",
1555 "grpc::client");
1556 }
1557
1558 // Validate method format
1559 if (method.empty() || method[0] != '/')
1560 {
1563 "Invalid method format",
1564 "grpc::client",
1565 "Method must start with '/'");
1566 }
1567
1568 // Build gRPC headers
1569 std::vector<http2::http_header> headers;
1570 headers.emplace_back("content-type", std::string(grpc_content_type));
1571 headers.emplace_back("te", "trailers");
1572 headers.emplace_back("grpc-accept-encoding",
1573 std::string(compression::identity) + "," +
1575
1576 // Add timeout header if deadline is set
1577 if (options.deadline.has_value())
1578 {
1579 auto now = std::chrono::system_clock::now();
1580 auto remaining = std::chrono::duration_cast<std::chrono::milliseconds>(
1581 options.deadline.value() - now);
1582 if (remaining.count() > 0)
1583 {
1584 headers.emplace_back("grpc-timeout",
1585 format_timeout(static_cast<uint64_t>(remaining.count())));
1586 }
1587 }
1588
1589 // Add custom metadata
1590 for (const auto& [key, value] : options.metadata)
1591 {
1592 headers.emplace_back(key, value);
1593 }
1594
1595 // Weak callback captures avoid a transport/stream ownership cycle.
1596 auto bidi = std::make_shared<bidi_stream_impl>(http2_client_, 0);
1597
1598 // Start streaming request
1599 auto stream_result = http2_client_->start_stream(
1600 method,
1601 headers,
1602 [weak = std::weak_ptr{bidi}](std::vector<uint8_t> data) {
1603 if (auto stream = weak.lock()) stream->on_data(data);
1604 },
1605 [weak = std::weak_ptr{bidi}](std::vector<http2::http_header> hdrs) {
1606 if (auto stream = weak.lock()) stream->on_headers(hdrs);
1607 },
1608 [weak = std::weak_ptr{bidi}](int status) {
1609 if (auto stream = weak.lock()) stream->on_complete(status);
1610 });
1611
1612 if (stream_result.is_err())
1613 {
1614 const auto& err = stream_result.error();
1616 err.code, err.message, "grpc::client", get_error_details(err));
1617 }
1618
1619 // Return the same object captured by the response callbacks.
1620 bidi->set_stream_id(stream_result.value());
1621
1622 struct shared_bidi_holder : public grpc_client::bidi_stream {
1623 std::shared_ptr<bidi_stream_impl> impl;
1624 explicit shared_bidi_holder(std::shared_ptr<bidi_stream_impl> p) : impl(std::move(p)) {}
1625 auto write(const std::vector<uint8_t>& message) -> VoidResult override { return impl->write(message); }
1626 auto read() -> Result<grpc_message> override { return impl->read(); }
1627 auto writes_done() -> VoidResult override { return impl->writes_done(); }
1628 auto finish() -> grpc_status override { return impl->finish(); }
1629 };
1630
1631 return ok(std::unique_ptr<grpc_client::bidi_stream>(
1632 new shared_bidi_holder(bidi)));
1633 }
1634
1635private:
1636 std::string target_;
1637 std::string host_;
1639 std::shared_ptr<http2::http2_client> http2_client_;
1640 std::atomic<bool> connected_;
1641 mutable std::mutex mutex_;
1642};
1643
1644#endif // !NETWORK_GRPC_OFFICIAL
1645
1646// grpc_client implementation
1647
1648grpc_client::grpc_client(const std::string& target,
1649 const grpc_channel_config& config)
1650 : impl_(std::make_shared<impl>(target, config))
1651{
1652}
1653
1654grpc_client::~grpc_client() = default;
1655
1656grpc_client::grpc_client(grpc_client&&) noexcept = default;
1657grpc_client& grpc_client::operator=(grpc_client&&) noexcept = default;
1658
1659auto grpc_client::connect() -> VoidResult
1660{
1661 return impl_->connect();
1662}
1663
1665{
1666 impl_->disconnect();
1667}
1668
1669auto grpc_client::is_connected() const -> bool
1670{
1671 return impl_->is_connected();
1672}
1673
1674auto grpc_client::wait_for_connected(std::chrono::milliseconds timeout) -> bool
1675{
1676 return impl_->wait_for_connected(timeout);
1677}
1678
1679auto grpc_client::target() const -> const std::string&
1680{
1681 return impl_->target();
1682}
1683
1684auto grpc_client::call_raw(const std::string& method,
1685 const std::vector<uint8_t>& request,
1686 const call_options& options) -> Result<grpc_message>
1687{
1688 return impl_->call_raw(method, request, options);
1689}
1690
1691auto grpc_client::call_raw_async(const std::string& method,
1692 const std::vector<uint8_t>& request,
1693 std::function<void(Result<grpc_message>)> callback,
1694 const call_options& options) -> void
1695{
1696 impl_->call_raw_async(method, request, std::move(callback), options);
1697}
1698
1699auto grpc_client::server_stream_raw(const std::string& method,
1700 const std::vector<uint8_t>& request,
1701 const call_options& options)
1703{
1704 return impl_->server_stream_raw(method, request, options);
1705}
1706
1707auto grpc_client::client_stream_raw(const std::string& method,
1708 const call_options& options)
1710{
1711 return impl_->client_stream_raw(method, options);
1712}
1713
1714auto grpc_client::bidi_stream_raw(const std::string& method,
1715 const call_options& options)
1717{
1718 return impl_->bidi_stream_raw(method, options);
1719}
1720
1721} // namespace kcenon::network::protocols::grpc
static thread_integration_manager & instance()
Get the singleton instance.
std::future< void > submit_task(std::function< void()> task)
Submit a task to the thread pool.
auto write(const std::vector< uint8_t > &message) -> VoidResult override
Write message to stream.
Definition client.cpp:890
void on_headers(const std::vector< http2::http_header > &headers)
Definition client.cpp:969
auto read() -> Result< grpc_message > override
Read next message from stream.
Definition client.cpp:906
void on_data(const std::vector< uint8_t > &data)
Definition client.cpp:962
auto finish() -> grpc_status override
Finish the call and get final status.
Definition client.cpp:950
std::shared_ptr< http2::http2_client > http2_client_
Definition client.cpp:999
auto writes_done() -> VoidResult override
Signal that writing is done.
Definition client.cpp:939
bidi_stream_impl(std::shared_ptr< http2::http2_client > http2_client, uint32_t stream_id)
Definition client.cpp:879
auto finish() -> Result< grpc_message > override
Finish the call and get response.
Definition client.cpp:788
void on_headers(const std::vector< http2::http_header > &headers)
Definition client.cpp:822
std::shared_ptr< http2::http2_client > http2_client_
Definition client.cpp:864
void on_data(const std::vector< uint8_t > &data)
Definition client.cpp:816
auto write(const std::vector< uint8_t > &message) -> VoidResult override
Write message to stream.
Definition client.cpp:760
auto writes_done() -> VoidResult override
Signal that writing is done.
Definition client.cpp:777
client_stream_writer_impl(std::shared_ptr< http2::http2_client > http2_client, uint32_t stream_id)
Definition client.cpp:750
auto call_raw(const std::string &method, const std::vector< uint8_t > &request, const call_options &options) -> Result< grpc_message >
Definition client.cpp:1150
auto bidi_stream_raw(const std::string &method, const call_options &options) -> Result< std::unique_ptr< grpc_client::bidi_stream > >
Definition client.cpp:1546
std::shared_ptr< http2::http2_client > http2_client_
Definition client.cpp:1639
impl(std::string target, grpc_channel_config config)
Definition client.cpp:1013
auto client_stream_raw(const std::string &method, const call_options &options) -> Result< std::unique_ptr< grpc_client::client_stream_writer > >
Definition client.cpp:1458
auto wait_for_connected(std::chrono::milliseconds timeout) -> bool
Definition client.cpp:1129
auto server_stream_raw(const std::string &method, const std::vector< uint8_t > &request, const call_options &options) -> Result< std::unique_ptr< grpc_client::server_stream_reader > >
Definition client.cpp:1356
auto target() const -> const std::string &
Definition client.cpp:1145
auto call_raw_async(const std::string &method, const std::vector< uint8_t > &request, std::function< void(Result< grpc_message >)> callback, const call_options &options) -> void
Definition client.cpp:1340
gRPC client for making RPC calls
Definition client.h:115
auto is_connected() const -> bool
Check if connected.
Definition client.cpp:1669
auto client_stream_raw(const std::string &method, const call_options &options={}) -> Result< std::unique_ptr< client_stream_writer > >
Start a client streaming RPC call.
Definition client.cpp:1707
auto target() const -> const std::string &
Get the target address.
Definition client.cpp:1679
auto call_raw_async(const std::string &method, const std::vector< uint8_t > &request, std::function< void(Result< grpc_message >)> callback, const call_options &options={}) -> void
Make an async unary RPC call.
Definition client.cpp:1691
auto wait_for_connected(std::chrono::milliseconds timeout) -> bool
Wait for connection to be ready.
Definition client.cpp:1674
auto disconnect() -> void
Disconnect from the server.
Definition client.cpp:1664
auto call_raw(const std::string &method, const std::vector< uint8_t > &request, const call_options &options={}) -> Result< grpc_message >
Make a unary RPC call.
Definition client.cpp:1684
auto server_stream_raw(const std::string &method, const std::vector< uint8_t > &request, const call_options &options={}) -> Result< std::unique_ptr< server_stream_reader > >
Start a server streaming RPC call.
Definition client.cpp:1699
auto bidi_stream_raw(const std::string &method, const call_options &options={}) -> Result< std::unique_ptr< bidi_stream > >
Start a bidirectional streaming RPC call.
Definition client.cpp:1714
grpc_client(const std::string &target, const grpc_channel_config &config={})
Construct gRPC client.
Definition client.cpp:1648
auto has_more() const -> bool override
Check if stream has more messages.
Definition client.cpp:687
auto finish() -> grpc_status override
Get final status after stream ends.
Definition client.cpp:693
auto read() -> Result< grpc_message > override
Read next message from stream.
Definition client.cpp:651
std::shared_ptr< http2::http2_client > http2_client_
Definition client.cpp:737
void on_headers(const std::vector< http2::http_header > &headers)
Definition client.cpp:707
server_stream_reader_impl(std::shared_ptr< http2::http2_client > http2_client, uint32_t stream_id)
Definition client.cpp:643
void on_data(const std::vector< uint8_t > &data)
Definition client.cpp:700
static auto create_span(std::string_view name) -> span
Create a new root span with a new trace context.
gRPC client channel configuration and connection.
void set_timeout(D duration)
Set a timeout using any duration type.
tracing_config config
Definition exporters.cpp:29
Official gRPC library wrapper interfaces.
gRPC message framing and serialization.
constexpr const char * grpc_accept_encoding
Definition frame.h:134
gRPC protocol implementation
Definition client.h:34
auto format_timeout(uint64_t timeout_ms) -> std::string
Format timeout as gRPC timeout string.
Definition frame.cpp:131
constexpr const char * grpc_content_type
gRPC content-type header value
Definition frame.h:109
status_code
gRPC status codes (as defined in grpc/status.h)
Definition status.h:36
@ deadline_exceeded
Deadline expired before operation completed.
@ unimplemented
Operation not implemented.
@ ok
Not an error; returned on success.
constexpr auto status_code_to_string(status_code code) -> std::string_view
Convert status code to string.
Definition status.h:61
auto is_tracing_enabled() -> bool
Check if tracing is enabled.
std::string get_error_details(const simple_error &err)
Result< std::monostate > VoidResult
VoidResult error_void(int code, const std::string &message, const std::string &source="network_system", const std::string &details="")
RAII span implementation for distributed tracing.
Options for individual RPC calls.
Definition client.h:79
std::string root_certificates
Root certificates for TLS (PEM format)
Definition client.h:53
std::optional< std::string > client_key
Client private key for mutual TLS (PEM format)
Definition client.h:59
std::optional< std::string > client_certificate
Client certificate for mutual TLS (PEM format)
Definition client.h:56
std::chrono::milliseconds default_timeout
Default timeout for RPC calls.
Definition client.h:47
gRPC message with compression flag and payload
Definition frame.h:50
static auto parse(std::span< const uint8_t > input) -> Result< grpc_message >
Parse gRPC message from raw bytes.
Definition frame.cpp:14
auto serialize() const -> std::vector< uint8_t >
Serialize message to bytes with length prefix.
Definition frame.cpp:67
std::vector< uint8_t > data
Message payload.
Definition frame.h:52
gRPC status with code, message, and optional details
Definition status.h:95
Thread system integration interface for network_system.
Distributed tracing context for OpenTelemetry-compatible tracing.
Configuration structures for OpenTelemetry tracing.