Go to the documentation of this file.
19 #ifndef GRPCPP_SUPPORT_SERVER_CALLBACK_H
20 #define GRPCPP_SUPPORT_SERVER_CALLBACK_H
33 #include <type_traits>
35 #include "absl/functional/any_invocable.h"
36 #include "absl/status/status.h"
48 namespace experimental {
51 grpc_core::Transport** out_transport,
52 grpc_endpoint** out_endpoint);
54 grpc_core::Transport* transport, grpc_endpoint* endpoint,
62 template <
class Request,
class Response>
63 class CallbackUnaryHandler;
64 template <
class Request,
class Response>
65 class CallbackClientStreamingHandler;
66 template <
class Request,
class Response>
67 class CallbackServerStreamingHandler;
68 template <
class Request,
class Response>
69 class CallbackBidiHandler;
86 template <
class Request,
class Response>
88 template <
class Request,
class Response>
90 template <
class Request,
class Response>
92 template <
class Request,
class Response>
119 ScheduleOnDone(reactor()->InternalInlineable());
125 ScheduleOnDone(inline_ondone);
133 CallOnCancel(reactor);
143 CallOnCancel(reactor());
149 void Ref() { callbacks_outstanding_.fetch_add(1, std::memory_order_relaxed); }
156 virtual void RunAsync(absl::AnyInvocable<
void()> cb) {
163 virtual void CallOnDone() = 0;
167 void ScheduleOnDone(
bool inline_ondone);
171 void CallOnCancel(ServerReactor* reactor);
175 bool UnblockCancellation() {
176 return on_cancel_conditions_remaining_.fetch_sub(
177 1, std::memory_order_acq_rel) == 1;
182 return callbacks_outstanding_.fetch_sub(1, std::memory_order_acq_rel);
185 std::atomic_int on_cancel_conditions_remaining_{2};
186 std::atomic_int callbacks_outstanding_{
190 template <
class Request,
class Response>
203 Request request_obj_;
204 Response response_obj_;
210 class ServerUnaryReactor;
211 template <
class Request>
213 template <
class Response>
215 template <
class Request,
class Response>
218 namespace experimental {
219 class ServerSessionReactor;
220 class ServerCallbackSession;
235 template <
class Reactor>
237 reactor->InternalBindCall(
this);
241 namespace experimental {
250 absl::AnyInvocable<
void(
absl::Status)> on_shutdown) = 0;
253 template <
class Reactor>
255 reactor->InternalBindSession(
this);
260 template <
class Request>
266 virtual void Read(Request* msg) = 0;
270 reactor->InternalBindReader(
this);
274 template <
class Response>
287 reactor->InternalBindWriter(
this);
291 template <
class Request,
class Response>
298 virtual void Read(Request* msg) = 0;
305 reactor->InternalBindStream(
this);
321 template <
class Request,
class Response>
337 stream_.load(std::memory_order_acquire);
338 if (stream ==
nullptr) {
340 stream = stream_.load(std::memory_order_relaxed);
341 if (stream ==
nullptr) {
342 backlog_.send_initial_metadata_wanted =
true;
353 void StartRead(Request* req) ABSL_LOCKS_EXCLUDED(stream_mu_) {
355 stream_.load(std::memory_order_acquire);
356 if (stream ==
nullptr) {
358 stream = stream_.load(std::memory_order_relaxed);
359 if (stream ==
nullptr) {
360 backlog_.read_wanted = req;
383 ABSL_LOCKS_EXCLUDED(stream_mu_) {
385 stream_.load(std::memory_order_acquire);
386 if (stream ==
nullptr) {
388 stream = stream_.load(std::memory_order_relaxed);
389 if (stream ==
nullptr) {
390 backlog_.write_wanted = resp;
391 backlog_.write_options_wanted = options;
395 stream->
Write(resp, options);
414 stream_.load(std::memory_order_acquire);
415 if (stream ==
nullptr) {
417 stream = stream_.load(std::memory_order_relaxed);
418 if (stream ==
nullptr) {
419 backlog_.write_and_finish_wanted =
true;
420 backlog_.write_wanted = resp;
421 backlog_.write_options_wanted = options;
422 backlog_.status_wanted = std::move(s);
449 stream_.load(std::memory_order_acquire);
450 if (stream ==
nullptr) {
452 stream = stream_.load(std::memory_order_relaxed);
453 if (stream ==
nullptr) {
454 backlog_.finish_wanted =
true;
455 backlog_.status_wanted = std::move(s);
459 stream->
Finish(std::move(s));
486 void OnDone()
override = 0;
497 virtual void InternalBindStream(
501 if (
GPR_UNLIKELY(backlog_.send_initial_metadata_wanted)) {
505 stream->
Read(backlog_.read_wanted);
509 std::move(backlog_.write_options_wanted),
510 std::move(backlog_.status_wanted));
513 stream->
Write(backlog_.write_wanted,
514 std::move(backlog_.write_options_wanted));
517 stream->
Finish(std::move(backlog_.status_wanted));
521 stream_.store(stream, std::memory_order_release);
529 std::atomic<ServerCallbackReaderWriter<Request, Response>*> stream_{
nullptr};
530 struct PreBindBacklog {
531 bool send_initial_metadata_wanted =
false;
532 bool write_and_finish_wanted =
false;
533 bool finish_wanted =
false;
534 Request* read_wanted =
nullptr;
535 const Response* write_wanted =
nullptr;
539 PreBindBacklog backlog_ ABSL_GUARDED_BY(stream_mu_);
543 template <
class Request>
544 class ServerReadReactor :
public internal::ServerReactor {
552 reader_.load(std::memory_order_acquire);
553 if (reader ==
nullptr) {
555 reader = reader_.load(std::memory_order_relaxed);
556 if (reader ==
nullptr) {
557 backlog_.send_initial_metadata_wanted =
true;
563 void StartRead(Request* req) ABSL_LOCKS_EXCLUDED(reader_mu_) {
565 reader_.load(std::memory_order_acquire);
566 if (reader ==
nullptr) {
568 reader = reader_.load(std::memory_order_relaxed);
569 if (reader ==
nullptr) {
570 backlog_.read_wanted = req;
578 reader_.load(std::memory_order_acquire);
579 if (reader ==
nullptr) {
581 reader = reader_.load(std::memory_order_relaxed);
582 if (reader ==
nullptr) {
583 backlog_.finish_wanted =
true;
584 backlog_.status_wanted = std::move(s);
588 reader->
Finish(std::move(s));
594 void OnDone()
override = 0;
603 ABSL_LOCKS_EXCLUDED(reader_mu_) {
606 if (
GPR_UNLIKELY(backlog_.send_initial_metadata_wanted)) {
607 reader->SendInitialMetadata();
610 reader->Read(backlog_.read_wanted);
613 reader->Finish(std::move(backlog_.status_wanted));
616 reader_.store(reader, std::memory_order_release);
620 std::atomic<ServerCallbackReader<Request>*> reader_{
nullptr};
621 struct PreBindBacklog {
622 bool send_initial_metadata_wanted =
false;
623 bool finish_wanted =
false;
624 Request* read_wanted =
nullptr;
627 PreBindBacklog backlog_ ABSL_GUARDED_BY(reader_mu_);
631 template <
class Response>
632 class ServerWriteReactor :
public internal::ServerReactor {
640 writer_.load(std::memory_order_acquire);
641 if (writer ==
nullptr) {
643 writer = writer_.load(std::memory_order_relaxed);
644 if (writer ==
nullptr) {
645 backlog_.send_initial_metadata_wanted =
true;
655 ABSL_LOCKS_EXCLUDED(writer_mu_) {
657 writer_.load(std::memory_order_acquire);
658 if (writer ==
nullptr) {
660 writer = writer_.load(std::memory_order_relaxed);
661 if (writer ==
nullptr) {
662 backlog_.write_wanted = resp;
663 backlog_.write_options_wanted = options;
667 writer->
Write(resp, options);
672 writer_.load(std::memory_order_acquire);
673 if (writer ==
nullptr) {
675 writer = writer_.load(std::memory_order_relaxed);
676 if (writer ==
nullptr) {
677 backlog_.write_and_finish_wanted =
true;
678 backlog_.write_wanted = resp;
679 backlog_.write_options_wanted = options;
680 backlog_.status_wanted = std::move(s);
691 writer_.load(std::memory_order_acquire);
692 if (writer ==
nullptr) {
694 writer = writer_.load(std::memory_order_relaxed);
695 if (writer ==
nullptr) {
696 backlog_.finish_wanted =
true;
697 backlog_.status_wanted = std::move(s);
701 writer->
Finish(std::move(s));
707 void OnDone()
override = 0;
715 ABSL_LOCKS_EXCLUDED(writer_mu_) {
718 if (
GPR_UNLIKELY(backlog_.send_initial_metadata_wanted)) {
719 writer->SendInitialMetadata();
722 writer->WriteAndFinish(backlog_.write_wanted,
723 std::move(backlog_.write_options_wanted),
724 std::move(backlog_.status_wanted));
727 writer->Write(backlog_.write_wanted,
728 std::move(backlog_.write_options_wanted));
731 writer->Finish(std::move(backlog_.status_wanted));
735 writer_.store(writer, std::memory_order_release);
739 std::atomic<ServerCallbackWriter<Response>*> writer_{
nullptr};
740 struct PreBindBacklog {
741 bool send_initial_metadata_wanted =
false;
742 bool write_and_finish_wanted =
false;
743 bool finish_wanted =
false;
744 const Response* write_wanted =
nullptr;
748 PreBindBacklog backlog_ ABSL_GUARDED_BY(writer_mu_);
759 if (call ==
nullptr) {
761 call = call_.load(std::memory_order_relaxed);
762 if (call ==
nullptr) {
763 backlog_.send_initial_metadata_wanted =
true;
774 if (call ==
nullptr) {
776 call = call_.load(std::memory_order_relaxed);
777 if (call ==
nullptr) {
778 backlog_.finish_wanted =
true;
779 backlog_.status_wanted = std::move(s);
783 call->
Finish(std::move(s));
788 void OnDone()
override = 0;
796 ABSL_LOCKS_EXCLUDED(call_mu_) {
799 if (
GPR_UNLIKELY(backlog_.send_initial_metadata_wanted)) {
800 call->SendInitialMetadata();
803 call->Finish(std::move(backlog_.status_wanted));
806 call_.store(call, std::memory_order_release);
810 std::atomic<ServerCallbackUnary*> call_{
nullptr};
811 struct PreBindBacklog {
812 bool send_initial_metadata_wanted =
false;
813 bool finish_wanted =
false;
816 PreBindBacklog backlog_ ABSL_GUARDED_BY(call_mu_);
819 namespace experimental {
829 if (session ==
nullptr) {
831 session = session_.load(std::memory_order_relaxed);
832 if (session ==
nullptr) {
833 backlog_.send_initial_metadata_wanted =
true;
845 if (session ==
nullptr) {
847 session = session_.load(std::memory_order_relaxed);
848 if (session ==
nullptr) {
849 backlog_.finish_wanted =
true;
850 backlog_.status_wanted = std::move(s);
854 session->
Finish(std::move(s));
865 ABSL_LOCKS_EXCLUDED(session_mu_) {
867 if (session ==
nullptr) {
869 session = session_.load(std::memory_order_relaxed);
870 if (session ==
nullptr) {
871 backlog_.graceful_shutdown_wanted_callback = std::move(on_shutdown);
880 void OnDone()
override = 0;
888 ABSL_LOCKS_EXCLUDED(session_mu_) {
891 if (
GPR_UNLIKELY(backlog_.send_initial_metadata_wanted)) {
892 session->SendInitialMetadata();
894 if (
GPR_UNLIKELY(backlog_.graceful_shutdown_wanted_callback !=
nullptr)) {
895 session->InitiateGracefulShutdown(
896 std::move(backlog_.graceful_shutdown_wanted_callback));
899 session->Finish(std::move(backlog_.status_wanted));
902 session_.store(session, std::memory_order_release);
906 std::atomic<ServerCallbackSession*> session_{
nullptr};
907 struct PreBindBacklog {
908 bool send_initial_metadata_wanted =
false;
909 bool finish_wanted =
false;
910 absl::AnyInvocable<void(
absl::Status)> graceful_shutdown_wanted_callback;
918 PreBindBacklog backlog_ ABSL_GUARDED_BY(session_mu_);
924 template <
class Base>
925 class FinishOnlyReactor :
public Base {
932 template <
class Request>
934 template <
class Response>
937 template <
class Request,
class Response>
943 namespace experimental {
952 namespace experimental {
954 template <
class Request,
class Response>
961 #endif // GRPCPP_SUPPORT_SERVER_CALLBACK_H
ServerReadReactor is the interface for a client-streaming RPC.
Definition: server_callback.h:212
virtual void SendInitialMetadata()=0
~ServerCallbackReader() override
Definition: server_callback.h:263
void set_response(Response *response)
Definition: message_allocator.h:50
void OnCancel() override
Definition: server_callback.h:595
void StartSendInitialMetadata() ABSL_LOCKS_EXCLUDED(stream_mu_)
Send any initial metadata stored in the RPC context.
Definition: server_callback.h:335
void StartVirtualRPCs() ABSL_LOCKS_EXCLUDED(session_mu_)
StartVirtualRPCs is exactly like ServerBidiReactor's StartSendInitialMetadata.
Definition: server_callback.h:827
void set_request(Request *request)
Definition: message_allocator.h:49
ServerReadReactor()
Definition: server_callback.h:546
virtual void Finish(grpc::Status s)=0
void StartWrite(const Response *resp)
Definition: server_callback.h:651
Represents a gRPC server.
Definition: server.h:58
virtual void Finish(grpc::Status s)=0
An Alarm posts the user-provided tag to its associated completion queue or invokes the user-provided ...
Definition: alarm.h:33
void StartWriteAndFinish(const Response *resp, grpc::WriteOptions options, grpc::Status s) ABSL_LOCKS_EXCLUDED(stream_mu_)
Initiate a write operation with specified options and final RPC Status, which also causes any trailin...
Definition: server_callback.h:411
WriteOptions & set_last_message()
last-message bit: indicates this is the last message in a stream client-side: makes Write the equival...
Definition: call_op_set.h:156
Definition: context_types.h:20
virtual ~ServerReactor()=default
void MaybeCallOnCancel()
Definition: server_callback.h:141
void StartWrite(const Response *resp, grpc::WriteOptions options) ABSL_LOCKS_EXCLUDED(writer_mu_)
Definition: server_callback.h:654
virtual void OnWriteDone(bool)
Definition: server_callback.h:706
ServerBidiReactor()
Definition: server_callback.h:329
virtual void OnSendInitialMetadataDone(bool)
The following notifications are exactly like ServerBidiReactor.
Definition: server_callback.h:705
virtual void Finish(grpc::Status s)=0
void StartRead(Request *req) ABSL_LOCKS_EXCLUDED(reader_mu_)
Definition: server_callback.h:563
DefaultMessageHolder()
Definition: server_callback.h:193
~ServerSessionReactor() override=default
~ServerCallbackReaderWriter() override
Definition: server_callback.h:294
void BindReactor(Reactor *reactor)
Definition: server_callback.h:236
void MaybeCallOnCancel(ServerReactor *reactor)
Definition: server_callback.h:131
void StartWriteLast(const Response *resp, grpc::WriteOptions options)
Definition: server_callback.h:686
~ServerReadReactor() override=default
virtual void InitiateGracefulShutdown(absl::AnyInvocable< void(absl::Status)> on_shutdown)=0
Definition: server_callback.h:292
void StartWrite(const Response *resp)
Initiate a write operation.
Definition: server_callback.h:372
virtual void BindInnerServer(grpc::Server *inner_server)=0
void OnCancel() override
Notifies the application that this RPC has been cancelled.
Definition: server_callback.h:491
virtual void Write(const Response *msg, grpc::WriteOptions options)=0
~ServerCallbackSession() override
Definition: server_callback.h:244
void StartRead(Request *req) ABSL_LOCKS_EXCLUDED(stream_mu_)
Initiate a read operation.
Definition: server_callback.h:353
Did it work? If it didn't, why?
Definition: status.h:34
~ServerCallbackUnary() override
Definition: server_callback.h:228
virtual void OnSendInitialMetadataDone(bool)
The following notifications are exactly like ServerBidiReactor.
Definition: server_callback.h:787
ServerSessionReactor()
Definition: server_callback.h:822
Definition: server_callback.h:820
The base class of ServerCallbackUnary etc.
Definition: server_callback.h:97
virtual void OnSendInitialMetadataDone(bool)
The following notifications are exactly like ServerBidiReactor.
Definition: server_callback.h:879
virtual void OnWriteDone(bool)
Notifies the application that a StartWrite (or StartWriteLast) operation completed.
Definition: server_callback.h:481
void StartWrite(const Response *resp, grpc::WriteOptions options) ABSL_LOCKS_EXCLUDED(stream_mu_)
Initiate a write operation with specified options.
Definition: server_callback.h:382
void Finish(grpc::Status s) ABSL_LOCKS_EXCLUDED(session_mu_)
Finish is similar to ServerBidiReactor except for one detail.
Definition: server_callback.h:843
virtual void WriteAndFinish(const Response *msg, grpc::WriteOptions options, grpc::Status s)=0
void OnDone() override=0
Notifies the application that all operations associated with this RPC have completed.
struct grpc_call grpc_call
A Call represents an RPC.
Definition: grpc_types.h:68
ServerBidiReactor is the interface for a bidirectional streaming RPC.
Definition: server_callback.h:216
Definition: server_callback_handlers.h:38
virtual void SendInitialMetadata()=0
~ServerBidiReactor() override=default
bool ReturnPreexistingErrors()
Definition: server_callback.h:191
FinishOnlyReactor(grpc::Status s)
Definition: server_callback.h:927
virtual void OnReadDone(bool)
Notifies the application that a StartRead operation completed.
Definition: server_callback.h:474
void StartWriteAndFinish(const Response *resp, grpc::WriteOptions options, grpc::Status s) ABSL_LOCKS_EXCLUDED(writer_mu_)
Definition: server_callback.h:669
Definition: server_callback_handlers.h:460
Definition: server_callback.h:242
void InitiateSessionGracefulShutdown(grpc_core::Transport *transport, grpc_endpoint *endpoint, absl::AnyInvocable< void(absl::Status)> on_shutdown)
virtual void Read(Request *msg)=0
void Finish(grpc::Status s) ABSL_LOCKS_EXCLUDED(stream_mu_)
Indicate that the stream is to be finished and the trailing metadata and RPC status are to be sent.
Definition: server_callback.h:447
void BindReactor(ServerReadReactor< Request > *reactor)
Definition: server_callback.h:269
void BindReactor(ServerWriteReactor< Response > *reactor)
Definition: server_callback.h:286
Definition: server_context.h:90
~ServerCallbackWriter() override
Definition: server_callback.h:277
void grpc_call_run_in_event_engine(const grpc_call *call, absl::AnyInvocable< void()> cb)
void StartSendInitialMetadata() ABSL_LOCKS_EXCLUDED(reader_mu_)
The following operation initiations are exactly like ServerBidiReactor.
Definition: server_callback.h:550
virtual ~ServerCallbackCall()
Definition: server_callback.h:99
Per-message write options.
Definition: call_op_set.h:79
void BindReactor(Reactor *reactor)
Definition: server_callback.h:254
void MaybeDone(bool inline_ondone)
Definition: server_callback.h:123
ServerWriteReactor is the interface for a server-streaming RPC.
Definition: server_callback.h:214
virtual void OnReadDone(bool)
Definition: server_callback.h:593
virtual void Read(Request *msg)=0
Definition: server_callback_handlers.h:695
void OnCancel() override
Definition: server_callback.h:708
~ServerWriteReactor() override=default
void BindReactor(ServerBidiReactor< Request, Response > *reactor)
Definition: server_callback.h:304
Definition: server_callback.h:275
void StartSendInitialMetadata() ABSL_LOCKS_EXCLUDED(writer_mu_)
The following operation initiations are exactly like ServerBidiReactor.
Definition: server_callback.h:638
grpc::internal::FinishOnlyReactor< ServerSessionReactor > UnimplementedSessionReactor
Definition: server_callback.h:946
virtual bool InternalInlineable()
Definition: server_callback.h:83
virtual void SendInitialMetadata()=0
virtual void OnSendInitialMetadataDone(bool)
The following notifications are exactly like ServerBidiReactor.
Definition: server_callback.h:592
void StartWriteLast(const Response *resp, grpc::WriteOptions options)
Inform system of a planned write operation with specified options, but allow the library to schedule ...
Definition: server_callback.h:437
void OnDone() override
Definition: server_callback.h:928
Definition: server_callback.h:73
virtual void SendInitialMetadata()=0
Definition: server_callback_handlers.h:262
void OnCancel() override
Definition: server_callback.h:789
void Ref()
Increases the reference count.
Definition: server_callback.h:149
virtual void OnCancel()=0
Definition: server_callback.h:261
void Finish(grpc::Status s) ABSL_LOCKS_EXCLUDED(reader_mu_)
Definition: server_callback.h:576
virtual void Write(const Response *msg, grpc::WriteOptions options)=0
Definition: message_allocator.h:39
void StartSendInitialMetadata() ABSL_LOCKS_EXCLUDED(call_mu_)
StartSendInitialMetadata is exactly like ServerBidiReactor.
Definition: server_callback.h:757
void Finish(grpc::Status s) ABSL_LOCKS_EXCLUDED(writer_mu_)
Definition: server_callback.h:689
void BindSessionToInnerServer(grpc_call *call, grpc::Server *inner_server, grpc_core::Transport **out_transport, grpc_endpoint **out_endpoint)
void InitiateGracefulShutdown(absl::AnyInvocable< void(absl::Status)> on_shutdown) ABSL_LOCKS_EXCLUDED(session_mu_)
Initiate a graceful shutdown of a SESSION_RPC.
Definition: server_callback.h:863
void Finish(grpc::Status s) ABSL_LOCKS_EXCLUDED(call_mu_)
Finish is similar to ServerBidiReactor except for one detail.
Definition: server_callback.h:772
virtual void OnSendInitialMetadataDone(bool)
Notifies the application that an explicit StartSendInitialMetadata operation completed.
Definition: server_callback.h:468
~ServerUnaryReactor() override=default
Definition: server_callback.h:226
ServerUnaryReactor()
Definition: server_callback.h:753
virtual void Finish(grpc::Status s)=0
ServerWriteReactor()
Definition: server_callback.h:634
Definition: server_callback.h:751
::absl::Status Status
Definition: config_protobuf.h:107
virtual void WriteAndFinish(const Response *msg, grpc::WriteOptions options, grpc::Status s)=0
virtual void Finish(grpc::Status s)=0
virtual void SendInitialMetadata()=0
void MaybeDone()
Definition: server_callback.h:117
void OnCancel() override
Definition: server_callback.h:881
void Release() override
Definition: server_callback.h:197
::grpc::ServerBidiReactor< Request, Response > ServerBidiReactor
Definition: server_callback.h:955