18 #ifndef GRPCPP_IMPL_SERVER_CALLBACK_HANDLERS_H
19 #define GRPCPP_IMPL_SERVER_CALLBACK_HANDLERS_H
31 #include "absl/log/absl_check.h"
32 #include "absl/status/status.h"
37 template <
class RequestType,
class ResponseType>
42 const RequestType*, ResponseType*)>
44 : get_reactor_(
std::move(get_reactor)) {}
48 allocator_ = allocator;
54 auto* allocator_state =
59 sizeof(ServerCallbackUnaryImpl)))
60 ServerCallbackUnaryImpl(
62 param.call, allocator_state, param.call_requester);
63 param.server_context->BeginCompletionOp(
64 param.call, [call](
bool) { call->MaybeDone(); }, call);
71 }
else if (param.status.ok()) {
72 reactor = grpc::internal::CatchingReactorGetter<ServerUnaryReactor>(
78 if (reactor ==
nullptr) {
87 call->SetupReactor(reactor);
94 RequestType* request =
nullptr;
96 if (allocator_ !=
nullptr) {
103 *handler_data = allocator_state;
104 request = allocator_state->
request();
115 const RequestType*, ResponseType*)>
132 reactor_.load(std::memory_order_relaxed)->InternalInlineable());
135 finish_ops_.set_core_cq_tag(&finish_tag_);
137 if (!ctx_->sent_initial_metadata_) {
139 ctx_->initial_metadata_flags());
140 if (ctx_->compression_level_set()) {
141 finish_ops_.set_compression_level(ctx_->compression_level());
143 ctx_->MarkInitialMetadataSent();
147 finish_ops_.ServerSendStatus(
148 &ctx_->trailing_metadata_,
149 finish_ops_.SendMessagePtr(response(), ctx_->memory_allocator()));
151 finish_ops_.ServerSendStatus(&ctx_->trailing_metadata_, s);
153 finish_ops_.set_core_cq_tag(&finish_tag_);
154 finish_ops_.FillOps(&call_);
157 void SendInitialMetadata()
override {
158 ABSL_CHECK(!ctx_->sent_initial_metadata_);
168 ServerUnaryReactor* reactor =
169 reactor_.load(std::memory_order_relaxed);
170 reactor->OnSendInitialMetadataDone(ok);
171 this->MaybeDone(true);
174 meta_ops_.SendInitialMetadata(&ctx_->initial_metadata_,
175 ctx_->initial_metadata_flags());
176 if (ctx_->compression_level_set()) {
177 meta_ops_.set_compression_level(ctx_->compression_level());
179 ctx_->MarkInitialMetadataSent();
180 meta_ops_.set_core_cq_tag(&meta_tag_);
181 meta_ops_.FillOps(&call_);
187 ServerCallbackUnaryImpl(
189 MessageHolder<RequestType, ResponseType>* allocator_state,
190 std::function<
void()> call_requester)
193 allocator_state_(allocator_state),
194 call_requester_(
std::move(call_requester)) {
195 ctx_->set_message_allocator_state(allocator_state);
204 void SetupReactor(ServerUnaryReactor* reactor) {
205 reactor_.store(reactor, std::memory_order_relaxed);
208 this->
MaybeDone(reactor->InternalInlineable());
211 const RequestType* request() {
return allocator_state_->request(); }
212 ResponseType* response() {
return allocator_state_->response(); }
214 void CallOnDone()
override {
215 reactor_.load(std::memory_order_relaxed)->OnDone();
217 auto call_requester = std::move(call_requester_);
218 allocator_state_->Release();
219 if (ctx_->context_allocator() !=
nullptr) {
220 ctx_->context_allocator()->Release(ctx_);
222 this->~ServerCallbackUnaryImpl();
227 ServerReactor* reactor()
override {
228 return reactor_.load(std::memory_order_relaxed);
242 MessageHolder<RequestType, ResponseType>*
const allocator_state_;
243 std::function<void()> call_requester_;
254 std::atomic<ServerUnaryReactor*> reactor_;
256 std::atomic<intptr_t> callbacks_outstanding_{
261 template <
class RequestType,
class ResponseType>
268 : get_reactor_(
std::move(get_reactor)) {}
274 sizeof(ServerCallbackReaderImpl)))
275 ServerCallbackReaderImpl(
277 param.call, param.call_requester);
281 param.server_context->BeginCompletionOp(
283 [reader](
bool) { reader->MaybeDone(
false); },
291 }
else if (param.status.ok()) {
293 grpc::internal::CatchingReactorGetter<ServerReadReactor<RequestType>>(
299 if (reactor ==
nullptr) {
307 reader->SetupReactor(reactor);
327 this->MaybeDone(false);
330 if (!ctx_->sent_initial_metadata_) {
332 ctx_->initial_metadata_flags());
333 if (ctx_->compression_level_set()) {
334 finish_ops_.set_compression_level(ctx_->compression_level());
336 ctx_->MarkInitialMetadataSent();
340 finish_ops_.ServerSendStatus(
341 &ctx_->trailing_metadata_,
342 finish_ops_.SendMessagePtr(&resp_, ctx_->memory_allocator()));
344 finish_ops_.ServerSendStatus(&ctx_->trailing_metadata_, s);
346 finish_ops_.set_core_cq_tag(&finish_tag_);
347 finish_ops_.FillOps(&call_);
350 void SendInitialMetadata()
override {
351 ABSL_CHECK(!ctx_->sent_initial_metadata_);
359 ServerReadReactor<RequestType>* reactor =
360 reactor_.load(std::memory_order_relaxed);
361 reactor->OnSendInitialMetadataDone(ok);
362 this->MaybeDone(true);
365 meta_ops_.SendInitialMetadata(&ctx_->initial_metadata_,
366 ctx_->initial_metadata_flags());
367 if (ctx_->compression_level_set()) {
368 meta_ops_.set_compression_level(ctx_->compression_level());
370 ctx_->MarkInitialMetadataSent();
371 meta_ops_.set_core_cq_tag(&meta_tag_);
372 meta_ops_.FillOps(&call_);
375 void Read(RequestType* req)
override {
377 read_ops_.RecvMessage(req);
378 read_ops_.FillOps(&call_);
386 std::function<
void()> call_requester)
387 : ctx_(ctx), call_(*call), call_requester_(
std::move(call_requester)) {}
391 void SetupReactor(ServerReadReactor<RequestType>* reactor) {
392 reactor_.store(reactor, std::memory_order_relaxed);
398 [
this, reactor](
bool ok) {
399 if (GPR_UNLIKELY(!ok)) {
400 ctx_->MaybeMarkCancelledOnRead();
402 reactor->OnReadDone(ok);
403 this->MaybeDone(
true);
406 read_ops_.set_core_cq_tag(&read_tag_);
415 ~ServerCallbackReaderImpl() {}
417 ResponseType* response() {
return &resp_; }
419 void CallOnDone()
override {
420 reactor_.load(std::memory_order_relaxed)->OnDone();
422 auto call_requester = std::move(call_requester_);
423 if (ctx_->context_allocator() !=
nullptr) {
424 ctx_->context_allocator()->Release(ctx_);
426 this->~ServerCallbackReaderImpl();
431 ServerReactor* reactor()
override {
432 return reactor_.load(std::memory_order_relaxed);
450 std::function<void()> call_requester_;
452 std::atomic<ServerReadReactor<RequestType>*> reactor_;
454 std::atomic<intptr_t> callbacks_outstanding_{
459 template <
class RequestType,
class ResponseType>
466 : get_reactor_(
std::move(get_reactor)) {}
472 sizeof(ServerCallbackWriterImpl)))
473 ServerCallbackWriterImpl(
475 param.call,
static_cast<RequestType*
>(param.request),
476 param.call_requester);
480 param.server_context->BeginCompletionOp(
482 [writer](
bool) { writer->MaybeDone(
false); },
490 }
else if (param.status.ok()) {
497 if (reactor ==
nullptr) {
505 writer->SetupReactor(reactor);
519 request->~RequestType();
540 this->MaybeDone(false);
543 finish_ops_.set_core_cq_tag(&finish_tag_);
545 if (!ctx_->sent_initial_metadata_) {
547 ctx_->initial_metadata_flags());
548 if (ctx_->compression_level_set()) {
549 finish_ops_.set_compression_level(ctx_->compression_level());
551 ctx_->MarkInitialMetadataSent();
553 finish_ops_.ServerSendStatus(&ctx_->trailing_metadata_, s);
554 finish_ops_.FillOps(&call_);
557 void SendInitialMetadata()
override {
558 ABSL_CHECK(!ctx_->sent_initial_metadata_);
566 ServerWriteReactor<ResponseType>* reactor =
567 reactor_.load(std::memory_order_relaxed);
568 reactor->OnSendInitialMetadataDone(ok);
569 this->MaybeDone(true);
572 meta_ops_.SendInitialMetadata(&ctx_->initial_metadata_,
573 ctx_->initial_metadata_flags());
574 if (ctx_->compression_level_set()) {
575 meta_ops_.set_compression_level(ctx_->compression_level());
577 ctx_->MarkInitialMetadataSent();
578 meta_ops_.set_core_cq_tag(&meta_tag_);
579 meta_ops_.FillOps(&call_);
587 if (!ctx_->sent_initial_metadata_) {
588 write_ops_.SendInitialMetadata(&ctx_->initial_metadata_,
589 ctx_->initial_metadata_flags());
590 if (ctx_->compression_level_set()) {
591 write_ops_.set_compression_level(ctx_->compression_level());
593 ctx_->MarkInitialMetadataSent();
597 write_ops_.SendMessagePtr(resp, options, ctx_->memory_allocator())
599 write_ops_.FillOps(&call_);
607 finish_ops_.SendMessagePtr(resp, options, ctx_->memory_allocator())
609 Finish(std::move(s));
613 friend class CallbackServerStreamingHandler<RequestType, ResponseType>;
617 std::function<
void()> call_requester)
621 call_requester_(
std::move(call_requester)) {}
625 void SetupReactor(ServerWriteReactor<ResponseType>* reactor) {
626 reactor_.store(reactor, std::memory_order_relaxed);
632 [
this, reactor](
bool ok) {
633 reactor->OnWriteDone(ok);
634 this->MaybeDone(true);
637 write_ops_.set_core_cq_tag(&write_tag_);
638 this->BindReactor(reactor);
639 this->MaybeCallOnCancel(reactor);
643 this->MaybeDone(
false);
645 ~ServerCallbackWriterImpl() {
646 if (req_ !=
nullptr) {
647 req_->~RequestType();
651 const RequestType* request() {
return req_; }
653 void CallOnDone()
override {
654 reactor_.load(std::memory_order_relaxed)->OnDone();
656 auto call_requester = std::move(call_requester_);
657 if (ctx_->context_allocator() !=
nullptr) {
658 ctx_->context_allocator()->Release(ctx_);
660 this->~ServerCallbackWriterImpl();
665 ServerReactor* reactor()
override {
666 return reactor_.load(std::memory_order_relaxed);
684 const RequestType* req_;
685 std::function<void()> call_requester_;
687 std::atomic<ServerWriteReactor<ResponseType>*> reactor_;
689 std::atomic<intptr_t> callbacks_outstanding_{
694 template <
class RequestType,
class ResponseType>
701 : get_reactor_(
std::move(get_reactor)) {}
706 param.call->call(),
sizeof(ServerCallbackReaderWriterImpl)))
707 ServerCallbackReaderWriterImpl(
709 param.call, param.call_requester);
713 param.server_context->BeginCompletionOp(
715 [stream](
bool) { stream->MaybeDone(
false); },
724 }
else if (param.status.ok()) {
731 if (reactor ==
nullptr) {
740 stream->SetupReactor(reactor);
744 std::function<ServerBidiReactor<RequestType, ResponseType>*(
748 class ServerCallbackReaderWriterImpl
761 this->MaybeDone(false);
764 finish_ops_.set_core_cq_tag(&finish_tag_);
766 if (!ctx_->sent_initial_metadata_) {
768 ctx_->initial_metadata_flags());
769 if (ctx_->compression_level_set()) {
770 finish_ops_.set_compression_level(ctx_->compression_level());
772 ctx_->MarkInitialMetadataSent();
774 finish_ops_.ServerSendStatus(&ctx_->trailing_metadata_, s);
775 finish_ops_.FillOps(&call_);
778 void SendInitialMetadata()
override {
779 ABSL_CHECK(!ctx_->sent_initial_metadata_);
787 ServerBidiReactor<RequestType, ResponseType>* reactor =
788 reactor_.load(std::memory_order_relaxed);
789 reactor->OnSendInitialMetadataDone(ok);
790 this->MaybeDone(true);
793 meta_ops_.SendInitialMetadata(&ctx_->initial_metadata_,
794 ctx_->initial_metadata_flags());
795 if (ctx_->compression_level_set()) {
796 meta_ops_.set_compression_level(ctx_->compression_level());
798 ctx_->MarkInitialMetadataSent();
799 meta_ops_.set_core_cq_tag(&meta_tag_);
800 meta_ops_.FillOps(&call_);
808 if (!ctx_->sent_initial_metadata_) {
809 write_ops_.SendInitialMetadata(&ctx_->initial_metadata_,
810 ctx_->initial_metadata_flags());
811 if (ctx_->compression_level_set()) {
812 write_ops_.set_compression_level(ctx_->compression_level());
814 ctx_->MarkInitialMetadataSent();
818 write_ops_.SendMessagePtr(resp, options, ctx_->memory_allocator())
820 write_ops_.FillOps(&call_);
827 finish_ops_.SendMessagePtr(resp, options, ctx_->memory_allocator())
829 Finish(std::move(s));
832 void Read(RequestType* req)
override {
834 read_ops_.RecvMessage(req);
835 read_ops_.FillOps(&call_);
839 friend class CallbackBidiHandler<RequestType, ResponseType>;
843 std::function<
void()> call_requester)
844 : ctx_(ctx), call_(*call), call_requester_(
std::move(call_requester)) {}
848 void SetupReactor(ServerBidiReactor<RequestType, ResponseType>* reactor) {
849 reactor_.store(reactor, std::memory_order_relaxed);
855 [
this, reactor](
bool ok) {
856 reactor->OnWriteDone(ok);
857 this->MaybeDone(true);
860 write_ops_.set_core_cq_tag(&write_tag_);
863 [
this, reactor](
bool ok) {
864 if (GPR_UNLIKELY(!ok)) {
865 ctx_->MaybeMarkCancelledOnRead();
867 reactor->OnReadDone(ok);
868 this->MaybeDone(
true);
871 read_ops_.set_core_cq_tag(&read_tag_);
872 this->BindReactor(reactor);
873 this->MaybeCallOnCancel(reactor);
877 this->MaybeDone(
false);
880 void CallOnDone()
override {
881 reactor_.load(std::memory_order_relaxed)->OnDone();
883 auto call_requester = std::move(call_requester_);
884 if (ctx_->context_allocator() !=
nullptr) {
885 ctx_->context_allocator()->Release(ctx_);
887 this->~ServerCallbackReaderWriterImpl();
892 ServerReactor* reactor()
override {
893 return reactor_.load(std::memory_order_relaxed);
914 std::function<void()> call_requester_;
916 std::atomic<ServerBidiReactor<RequestType, ResponseType>*> reactor_;
918 std::atomic<intptr_t> callbacks_outstanding_{
925 namespace experimental {
928 template <
class RequestType>
936 : get_reactor_(
std::move(get_reactor)), service_(service) {
937 ABSL_CHECK(service_ !=
nullptr && service_->is_virtual_service_);
943 auto* allocator_state =
945 param.internal_data);
948 inner_server =
static_cast<grpc::Server*
>(service_->server_);
949 ABSL_CHECK(inner_server !=
nullptr);
952 sizeof(ServerCallbackSessionImpl)))
953 ServerCallbackSessionImpl(
955 param.call, allocator_state, param.call_requester, inner_server);
957 param.server_context->BeginCompletionOp(
958 param.call, [call](
bool) { call->MaybeDone(); }, call);
961 if (param.status.ok()) {
969 if (reactor ==
nullptr) {
978 call->SetupReactor(reactor);
985 RequestType* request =
nullptr;
991 *handler_data = allocator_state;
992 request = allocator_state->
request();
1007 class ServerCallbackSessionImpl
1011 if (ctx_->IsCancelled()) {
1013 reactor_.load(std::memory_order_relaxed)->InternalInlineable());
1026 reactor_.load(std::memory_order_relaxed)->InternalInlineable());
1028 &finish_ops_,
true);
1029 finish_ops_.set_core_cq_tag(&finish_tag_);
1031 bool is_first_metadata = !ctx_->sent_initial_metadata_;
1032 if (is_first_metadata) {
1033 finish_ops_.SendInitialMetadata(&ctx_->initial_metadata_,
1034 ctx_->initial_metadata_flags());
1035 if (ctx_->compression_level_set()) {
1036 finish_ops_.set_compression_level(ctx_->compression_level());
1038 ctx_->MarkInitialMetadataSent();
1040 finish_ops_.ServerSendStatus(&ctx_->trailing_metadata_, s);
1041 finish_ops_.set_core_cq_tag(&finish_tag_);
1042 finish_ops_.FillOps(&call_);
1045 void SendInitialMetadata()
override {
1046 ABSL_CHECK(!ctx_->sent_initial_metadata_);
1056 grpc::experimental::ServerSessionReactor* reactor =
1057 reactor_.load(std::memory_order_relaxed);
1058 reactor->OnSendInitialMetadataDone(ok);
1059 this->MaybeDone(true);
1062 meta_ops_.SendInitialMetadata(&ctx_->initial_metadata_,
1063 ctx_->initial_metadata_flags());
1064 if (ctx_->compression_level_set()) {
1065 meta_ops_.set_compression_level(ctx_->compression_level());
1067 ctx_->MarkInitialMetadataSent();
1068 meta_ops_.set_core_cq_tag(&meta_tag_);
1069 meta_ops_.FillOps(&call_);
1073 BindInnerServer(inner_server_);
1076 void BindInnerServer(
grpc::Server* inner_server)
override {
1078 call_.call(), inner_server, &transport_, &endpoint_);
1081 void InitiateGracefulShutdown(
1082 absl::AnyInvocable<
void(
absl::Status)> on_shutdown)
override {
1084 transport_, endpoint_, std::move(on_shutdown));
1088 friend class CallbackSessionHandler<RequestType>;
1090 ServerCallbackSessionImpl(
1092 MessageHolder<RequestType, grpc::ByteBuffer>* allocator_state,
1093 std::function<
void()> call_requester,
grpc::Server* inner_server)
1096 allocator_state_(allocator_state),
1097 call_requester_(
std::move(call_requester)),
1098 inner_server_(inner_server) {
1099 ABSL_CHECK(inner_server_ !=
nullptr);
1100 ctx_->set_message_allocator_state(allocator_state);
1110 reactor_.store(reactor, std::memory_order_relaxed);
1111 this->BindReactor(reactor);
1112 this->MaybeCallOnCancel(reactor);
1116 const RequestType* request() {
return allocator_state_->request(); }
1118 void CallOnDone()
override {
1119 reactor_.load(std::memory_order_relaxed)->OnDone();
1121 auto call_requester = std::move(call_requester_);
1122 allocator_state_->Release();
1123 if (ctx_->context_allocator() !=
nullptr) {
1124 ctx_->context_allocator()->Release(ctx_);
1126 this->~ServerCallbackSessionImpl();
1132 return reactor_.load(std::memory_order_relaxed);
1145 grpc_core::Transport* transport_ =
nullptr;
1146 grpc_endpoint* endpoint_ =
nullptr;
1147 MessageHolder<RequestType, grpc::ByteBuffer>*
const allocator_state_;
1148 std::function<void()> call_requester_;
1151 std::atomic<grpc::experimental::ServerSessionReactor*> reactor_;
1153 std::atomic<intptr_t> callbacks_outstanding_{
1162 #endif // GRPCPP_IMPL_SERVER_CALLBACK_HANDLERS_H