GRPC C++  1.83.0
server_callback.h
Go to the documentation of this file.
1 //
2 //
3 // Copyright 2018 gRPC authors.
4 //
5 // Licensed under the Apache License, Version 2.0 (the "License");
6 // you may not use this file except in compliance with the License.
7 // You may obtain a copy of the License at
8 //
9 // http://www.apache.org/licenses/LICENSE-2.0
10 //
11 // Unless required by applicable law or agreed to in writing, software
12 // distributed under the License is distributed on an "AS IS" BASIS,
13 // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14 // See the License for the specific language governing permissions and
15 // limitations under the License.
16 //
17 //
18 
19 #ifndef GRPCPP_SUPPORT_SERVER_CALLBACK_H
20 #define GRPCPP_SUPPORT_SERVER_CALLBACK_H
21 
22 #include <grpc/impl/call.h>
23 #include <grpcpp/impl/call.h>
25 #include <grpcpp/impl/sync.h>
27 #include <grpcpp/support/config.h>
29 #include <grpcpp/support/status.h>
30 
31 #include <atomic>
32 #include <functional>
33 #include <type_traits>
34 
35 #include "absl/functional/any_invocable.h"
36 #include "absl/status/status.h"
37 
38 struct grpc_endpoint;
39 namespace grpc_core {
40 class Transport;
41 }
42 
43 namespace grpc {
44 
45 class Server;
46 
47 // Declare base class of all reactors as internal
48 namespace experimental {
49 namespace internal {
50 void BindSessionToInnerServer(grpc_call* call, grpc::Server* inner_server,
51  grpc_core::Transport** out_transport,
52  grpc_endpoint** out_endpoint);
54  grpc_core::Transport* transport, grpc_endpoint* endpoint,
55  absl::AnyInvocable<void(absl::Status)> on_shutdown);
56 } // namespace internal
57 } // namespace experimental
58 
59 namespace internal {
60 
61 // Forward declarations
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;
70 
72 
74  public:
75  virtual ~ServerReactor() = default;
76  virtual void OnDone() = 0;
77  virtual void OnCancel() = 0;
78 
79  // The following is not API. It is for internal use only and specifies whether
80  // all reactions of this Reactor can be run without extra EventEngine
81  // scheduling. This should only be used for internally-defined reactors with
82  // trivial reactions.
83  virtual bool InternalInlineable() { return false; }
84 
85  private:
86  template <class Request, class Response>
87  friend class CallbackUnaryHandler;
88  template <class Request, class Response>
90  template <class Request, class Response>
92  template <class Request, class Response>
93  friend class CallbackBidiHandler;
94 };
95 
98  public:
99  virtual ~ServerCallbackCall() {}
100 
101  // This object is responsible for tracking when it is safe to call OnDone and
102  // OnCancel. OnDone should not be called until the method handler is complete,
103  // Finish has been called, the ServerContext CompletionOp (which tracks
104  // cancellation or successful completion) has completed, and all outstanding
105  // Read/Write actions have seen their reactions. OnCancel should not be called
106  // until after the method handler is done and the RPC has completed with a
107  // cancellation. This is tracked by counting how many of these conditions have
108  // been met and calling OnCancel when none remain unmet.
109 
110  // Public versions of MaybeDone: one where we don't know the reactor in
111  // advance (used for the ServerContext CompletionOp), and one for where we
112  // know the inlineability of the OnDone reaction. You should set the inline
113  // flag to true if either the Reactor is InternalInlineable() or if this
114  // callback is already being forced to run dispatched to an EventEngine thread
115  // (typically because it contains additional work than just the MaybeDone).
116 
117  void MaybeDone() {
118  if (GPR_UNLIKELY(Unref() == 1)) {
119  ScheduleOnDone(reactor()->InternalInlineable());
120  }
121  }
122 
123  void MaybeDone(bool inline_ondone) {
124  if (GPR_UNLIKELY(Unref() == 1)) {
125  ScheduleOnDone(inline_ondone);
126  }
127  }
128 
129  // Fast version called with known reactor passed in, used from derived
130  // classes, typically in non-cancel case
132  if (GPR_UNLIKELY(UnblockCancellation())) {
133  CallOnCancel(reactor);
134  }
135  }
136 
137  // Slower version called from object that doesn't know the reactor a priori
138  // (such as the ServerContext CompletionOp which is formed before the
139  // reactor). This is used in cancel cases only, so it's ok to be slower and
140  // invoke a virtual function.
142  if (GPR_UNLIKELY(UnblockCancellation())) {
143  CallOnCancel(reactor());
144  }
145  }
146 
147  protected:
149  void Ref() { callbacks_outstanding_.fetch_add(1, std::memory_order_relaxed); }
150 
151  private:
152  virtual ServerReactor* reactor() = 0;
153 
154  virtual grpc_call* call() = 0;
155 
156  virtual void RunAsync(absl::AnyInvocable<void()> cb) {
157  grpc_call_run_in_event_engine(call(), std::move(cb));
158  }
159 
160  // CallOnDone performs the work required at completion of the RPC: invoking
161  // the OnDone function and doing all necessary cleanup. This function is only
162  // ever invoked on a fully-Unref'fed ServerCallbackCall.
163  virtual void CallOnDone() = 0;
164 
165  // If the OnDone reaction is inlineable, execute it inline. Otherwise run it
166  // async on EventEngine.
167  void ScheduleOnDone(bool inline_ondone);
168 
169  // If the OnCancel reaction is inlineable, execute it inline. Otherwise run it
170  // async on EventEngine.
171  void CallOnCancel(ServerReactor* reactor);
172 
173  // Implement the cancellation constraint counter. Return true if OnCancel
174  // should be called, false otherwise.
175  bool UnblockCancellation() {
176  return on_cancel_conditions_remaining_.fetch_sub(
177  1, std::memory_order_acq_rel) == 1;
178  }
179 
181  int Unref() {
182  return callbacks_outstanding_.fetch_sub(1, std::memory_order_acq_rel);
183  }
184 
185  std::atomic_int on_cancel_conditions_remaining_{2};
186  std::atomic_int callbacks_outstanding_{
187  3}; // reserve for start, Finish, and CompletionOp
188 };
189 
190 template <class Request, class Response>
191 class DefaultMessageHolder : public MessageHolder<Request, Response> {
192  public:
194  this->set_request(&request_obj_);
195  this->set_response(&response_obj_);
196  }
197  void Release() override {
198  // the object is allocated in the call arena.
200  }
201 
202  private:
203  Request request_obj_;
204  Response response_obj_;
205 };
206 
207 } // namespace internal
208 
209 // Forward declarations
210 class ServerUnaryReactor;
211 template <class Request>
213 template <class Response>
215 template <class Request, class Response>
217 
218 namespace experimental {
219 class ServerSessionReactor;
220 class ServerCallbackSession;
221 } // namespace experimental
222 
223 // NOTE: The actual call/stream object classes are provided as API only to
224 // support mocking. There are no implementations of these class interfaces in
225 // the API.
227  public:
228  ~ServerCallbackUnary() override {}
229  virtual void Finish(grpc::Status s) = 0;
230  virtual void SendInitialMetadata() = 0;
231 
232  protected:
233  // Use a template rather than explicitly specifying ServerUnaryReactor to
234  // delay binding and avoid a circular forward declaration issue
235  template <class Reactor>
236  void BindReactor(Reactor* reactor) {
237  reactor->InternalBindCall(this);
238  }
239 };
240 
241 namespace experimental {
243  public:
245 
246  virtual void Finish(grpc::Status s) = 0;
247  virtual void SendInitialMetadata() = 0;
248  virtual void BindInnerServer(grpc::Server* inner_server) = 0;
249  virtual void InitiateGracefulShutdown(
250  absl::AnyInvocable<void(absl::Status)> on_shutdown) = 0;
251 
252  protected:
253  template <class Reactor>
254  void BindReactor(Reactor* reactor) {
255  reactor->InternalBindSession(this);
256  }
257 };
258 } // namespace experimental
259 
260 template <class Request>
262  public:
263  ~ServerCallbackReader() override {}
264  virtual void Finish(grpc::Status s) = 0;
265  virtual void SendInitialMetadata() = 0;
266  virtual void Read(Request* msg) = 0;
267 
268  protected:
270  reactor->InternalBindReader(this);
271  }
272 };
273 
274 template <class Response>
276  public:
277  ~ServerCallbackWriter() override {}
278 
279  virtual void Finish(grpc::Status s) = 0;
280  virtual void SendInitialMetadata() = 0;
281  virtual void Write(const Response* msg, grpc::WriteOptions options) = 0;
282  virtual void WriteAndFinish(const Response* msg, grpc::WriteOptions options,
283  grpc::Status s) = 0;
284 
285  protected:
287  reactor->InternalBindWriter(this);
288  }
289 };
290 
291 template <class Request, class Response>
293  public:
295 
296  virtual void Finish(grpc::Status s) = 0;
297  virtual void SendInitialMetadata() = 0;
298  virtual void Read(Request* msg) = 0;
299  virtual void Write(const Response* msg, grpc::WriteOptions options) = 0;
300  virtual void WriteAndFinish(const Response* msg, grpc::WriteOptions options,
301  grpc::Status s) = 0;
302 
303  protected:
305  reactor->InternalBindStream(this);
306  }
307 };
308 
309 // The following classes are the reactor interfaces that are to be implemented
310 // by the user, returned as the output parameter of the method handler for a
311 // callback method. Note that none of the classes are pure; all reactions have a
312 // default empty reaction so that the user class only needs to override those
313 // reactions that it cares about. The reaction methods will be invoked by the
314 // library in response to the completion of various operations. Reactions must
315 // not include blocking operations (such as blocking I/O, starting synchronous
316 // RPCs, or waiting on condition variables). Reactions may be invoked
317 // concurrently, except that OnDone is called after all others (assuming proper
318 // API usage). The reactor may not be deleted until OnDone is called.
319 
321 template <class Request, class Response>
322 class ServerBidiReactor : public internal::ServerReactor {
323  public:
324  // NOTE: Initializing stream_ as a constructor initializer rather than a
325  // default initializer because gcc-4.x requires a copy constructor for
326  // default initializing a templated member, which isn't ok for atomic.
327  // TODO(vjpai): Switch to default constructor and default initializer when
328  // gcc-4.x is no longer supported
329  ServerBidiReactor() : stream_(nullptr) {}
330  ~ServerBidiReactor() override = default;
331 
335  void StartSendInitialMetadata() ABSL_LOCKS_EXCLUDED(stream_mu_) {
337  stream_.load(std::memory_order_acquire);
338  if (stream == nullptr) {
339  grpc::internal::MutexLock l(&stream_mu_);
340  stream = stream_.load(std::memory_order_relaxed);
341  if (stream == nullptr) {
342  backlog_.send_initial_metadata_wanted = true;
343  return;
344  }
345  }
346  stream->SendInitialMetadata();
347  }
348 
353  void StartRead(Request* req) ABSL_LOCKS_EXCLUDED(stream_mu_) {
355  stream_.load(std::memory_order_acquire);
356  if (stream == nullptr) {
357  grpc::internal::MutexLock l(&stream_mu_);
358  stream = stream_.load(std::memory_order_relaxed);
359  if (stream == nullptr) {
360  backlog_.read_wanted = req;
361  return;
362  }
363  }
364  stream->Read(req);
365  }
366 
372  void StartWrite(const Response* resp) {
374  }
375 
382  void StartWrite(const Response* resp, grpc::WriteOptions options)
383  ABSL_LOCKS_EXCLUDED(stream_mu_) {
385  stream_.load(std::memory_order_acquire);
386  if (stream == nullptr) {
387  grpc::internal::MutexLock l(&stream_mu_);
388  stream = stream_.load(std::memory_order_relaxed);
389  if (stream == nullptr) {
390  backlog_.write_wanted = resp;
391  backlog_.write_options_wanted = options;
392  return;
393  }
394  }
395  stream->Write(resp, options);
396  }
397 
411  void StartWriteAndFinish(const Response* resp, grpc::WriteOptions options,
412  grpc::Status s) ABSL_LOCKS_EXCLUDED(stream_mu_) {
414  stream_.load(std::memory_order_acquire);
415  if (stream == nullptr) {
416  grpc::internal::MutexLock l(&stream_mu_);
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);
423  return;
424  }
425  }
426  stream->WriteAndFinish(resp, options, std::move(s));
427  }
428 
437  void StartWriteLast(const Response* resp, grpc::WriteOptions options) {
438  StartWrite(resp, options.set_last_message());
439  }
440 
447  void Finish(grpc::Status s) ABSL_LOCKS_EXCLUDED(stream_mu_) {
449  stream_.load(std::memory_order_acquire);
450  if (stream == nullptr) {
451  grpc::internal::MutexLock l(&stream_mu_);
452  stream = stream_.load(std::memory_order_relaxed);
453  if (stream == nullptr) {
454  backlog_.finish_wanted = true;
455  backlog_.status_wanted = std::move(s);
456  return;
457  }
458  }
459  stream->Finish(std::move(s));
460  }
461 
468  virtual void OnSendInitialMetadataDone(bool /*ok*/) {}
469 
474  virtual void OnReadDone(bool /*ok*/) {}
475 
481  virtual void OnWriteDone(bool /*ok*/) {}
482 
486  void OnDone() override = 0;
487 
491  void OnCancel() override {}
492 
493  private:
494  friend class ServerCallbackReaderWriter<Request, Response>;
495  // May be overridden by internal implementation details. This is not a public
496  // customization point.
497  virtual void InternalBindStream(
499  grpc::internal::MutexLock l(&stream_mu_);
500 
501  if (GPR_UNLIKELY(backlog_.send_initial_metadata_wanted)) {
502  stream->SendInitialMetadata();
503  }
504  if (GPR_UNLIKELY(backlog_.read_wanted != nullptr)) {
505  stream->Read(backlog_.read_wanted);
506  }
507  if (GPR_UNLIKELY(backlog_.write_and_finish_wanted)) {
508  stream->WriteAndFinish(backlog_.write_wanted,
509  std::move(backlog_.write_options_wanted),
510  std::move(backlog_.status_wanted));
511  } else {
512  if (GPR_UNLIKELY(backlog_.write_wanted != nullptr)) {
513  stream->Write(backlog_.write_wanted,
514  std::move(backlog_.write_options_wanted));
515  }
516  if (GPR_UNLIKELY(backlog_.finish_wanted)) {
517  stream->Finish(std::move(backlog_.status_wanted));
518  }
519  }
520  // Set stream_ last so that other functions can use it lock-free
521  stream_.store(stream, std::memory_order_release);
522  }
523 
524  grpc::internal::Mutex stream_mu_;
525  // TODO(vjpai): Make stream_or_backlog_ into an std::variant once C++17 or
526  // ABSL is supported since stream and backlog are mutually exclusive in this
527  // class. Do likewise with the remaining reactor classes and their backlogs
528  // as well.
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;
536  grpc::WriteOptions write_options_wanted;
537  grpc::Status status_wanted;
538  };
539  PreBindBacklog backlog_ ABSL_GUARDED_BY(stream_mu_);
540 };
541 
543 template <class Request>
544 class ServerReadReactor : public internal::ServerReactor {
545  public:
546  ServerReadReactor() : reader_(nullptr) {}
547  ~ServerReadReactor() override = default;
548 
550  void StartSendInitialMetadata() ABSL_LOCKS_EXCLUDED(reader_mu_) {
552  reader_.load(std::memory_order_acquire);
553  if (reader == nullptr) {
554  grpc::internal::MutexLock l(&reader_mu_);
555  reader = reader_.load(std::memory_order_relaxed);
556  if (reader == nullptr) {
557  backlog_.send_initial_metadata_wanted = true;
558  return;
559  }
560  }
561  reader->SendInitialMetadata();
562  }
563  void StartRead(Request* req) ABSL_LOCKS_EXCLUDED(reader_mu_) {
565  reader_.load(std::memory_order_acquire);
566  if (reader == nullptr) {
567  grpc::internal::MutexLock l(&reader_mu_);
568  reader = reader_.load(std::memory_order_relaxed);
569  if (reader == nullptr) {
570  backlog_.read_wanted = req;
571  return;
572  }
573  }
574  reader->Read(req);
575  }
576  void Finish(grpc::Status s) ABSL_LOCKS_EXCLUDED(reader_mu_) {
578  reader_.load(std::memory_order_acquire);
579  if (reader == nullptr) {
580  grpc::internal::MutexLock l(&reader_mu_);
581  reader = reader_.load(std::memory_order_relaxed);
582  if (reader == nullptr) {
583  backlog_.finish_wanted = true;
584  backlog_.status_wanted = std::move(s);
585  return;
586  }
587  }
588  reader->Finish(std::move(s));
589  }
590 
592  virtual void OnSendInitialMetadataDone(bool /*ok*/) {}
593  virtual void OnReadDone(bool /*ok*/) {}
594  void OnDone() override = 0;
595  void OnCancel() override {}
596 
597  private:
598  friend class ServerCallbackReader<Request>;
599 
600  // May be overridden by internal implementation details. This is not a public
601  // customization point.
602  virtual void InternalBindReader(ServerCallbackReader<Request>* reader)
603  ABSL_LOCKS_EXCLUDED(reader_mu_) {
604  grpc::internal::MutexLock l(&reader_mu_);
605 
606  if (GPR_UNLIKELY(backlog_.send_initial_metadata_wanted)) {
607  reader->SendInitialMetadata();
608  }
609  if (GPR_UNLIKELY(backlog_.read_wanted != nullptr)) {
610  reader->Read(backlog_.read_wanted);
611  }
612  if (GPR_UNLIKELY(backlog_.finish_wanted)) {
613  reader->Finish(std::move(backlog_.status_wanted));
614  }
615  // Set reader_ last so that other functions can use it lock-free
616  reader_.store(reader, std::memory_order_release);
617  }
618 
619  grpc::internal::Mutex reader_mu_;
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;
625  grpc::Status status_wanted;
626  };
627  PreBindBacklog backlog_ ABSL_GUARDED_BY(reader_mu_);
628 };
629 
631 template <class Response>
632 class ServerWriteReactor : public internal::ServerReactor {
633  public:
634  ServerWriteReactor() : writer_(nullptr) {}
635  ~ServerWriteReactor() override = default;
636 
638  void StartSendInitialMetadata() ABSL_LOCKS_EXCLUDED(writer_mu_) {
640  writer_.load(std::memory_order_acquire);
641  if (writer == nullptr) {
642  grpc::internal::MutexLock l(&writer_mu_);
643  writer = writer_.load(std::memory_order_relaxed);
644  if (writer == nullptr) {
645  backlog_.send_initial_metadata_wanted = true;
646  return;
647  }
648  }
649  writer->SendInitialMetadata();
650  }
651  void StartWrite(const Response* resp) {
653  }
654  void StartWrite(const Response* resp, grpc::WriteOptions options)
655  ABSL_LOCKS_EXCLUDED(writer_mu_) {
657  writer_.load(std::memory_order_acquire);
658  if (writer == nullptr) {
659  grpc::internal::MutexLock l(&writer_mu_);
660  writer = writer_.load(std::memory_order_relaxed);
661  if (writer == nullptr) {
662  backlog_.write_wanted = resp;
663  backlog_.write_options_wanted = options;
664  return;
665  }
666  }
667  writer->Write(resp, options);
668  }
669  void StartWriteAndFinish(const Response* resp, grpc::WriteOptions options,
670  grpc::Status s) ABSL_LOCKS_EXCLUDED(writer_mu_) {
672  writer_.load(std::memory_order_acquire);
673  if (writer == nullptr) {
674  grpc::internal::MutexLock l(&writer_mu_);
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);
681  return;
682  }
683  }
684  writer->WriteAndFinish(resp, options, std::move(s));
685  }
686  void StartWriteLast(const Response* resp, grpc::WriteOptions options) {
687  StartWrite(resp, options.set_last_message());
688  }
689  void Finish(grpc::Status s) ABSL_LOCKS_EXCLUDED(writer_mu_) {
691  writer_.load(std::memory_order_acquire);
692  if (writer == nullptr) {
693  grpc::internal::MutexLock l(&writer_mu_);
694  writer = writer_.load(std::memory_order_relaxed);
695  if (writer == nullptr) {
696  backlog_.finish_wanted = true;
697  backlog_.status_wanted = std::move(s);
698  return;
699  }
700  }
701  writer->Finish(std::move(s));
702  }
703 
705  virtual void OnSendInitialMetadataDone(bool /*ok*/) {}
706  virtual void OnWriteDone(bool /*ok*/) {}
707  void OnDone() override = 0;
708  void OnCancel() override {}
709 
710  private:
711  friend class ServerCallbackWriter<Response>;
712  // May be overridden by internal implementation details. This is not a public
713  // customization point.
714  virtual void InternalBindWriter(ServerCallbackWriter<Response>* writer)
715  ABSL_LOCKS_EXCLUDED(writer_mu_) {
716  grpc::internal::MutexLock l(&writer_mu_);
717 
718  if (GPR_UNLIKELY(backlog_.send_initial_metadata_wanted)) {
719  writer->SendInitialMetadata();
720  }
721  if (GPR_UNLIKELY(backlog_.write_and_finish_wanted)) {
722  writer->WriteAndFinish(backlog_.write_wanted,
723  std::move(backlog_.write_options_wanted),
724  std::move(backlog_.status_wanted));
725  } else {
726  if (GPR_UNLIKELY(backlog_.write_wanted != nullptr)) {
727  writer->Write(backlog_.write_wanted,
728  std::move(backlog_.write_options_wanted));
729  }
730  if (GPR_UNLIKELY(backlog_.finish_wanted)) {
731  writer->Finish(std::move(backlog_.status_wanted));
732  }
733  }
734  // Set writer_ last so that other functions can use it lock-free
735  writer_.store(writer, std::memory_order_release);
736  }
737 
738  grpc::internal::Mutex writer_mu_;
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;
745  grpc::WriteOptions write_options_wanted;
746  grpc::Status status_wanted;
747  };
748  PreBindBacklog backlog_ ABSL_GUARDED_BY(writer_mu_);
749 };
750 
752  public:
753  ServerUnaryReactor() : call_(nullptr) {}
754  ~ServerUnaryReactor() override = default;
755 
757  void StartSendInitialMetadata() ABSL_LOCKS_EXCLUDED(call_mu_) {
758  ServerCallbackUnary* call = call_.load(std::memory_order_acquire);
759  if (call == nullptr) {
760  grpc::internal::MutexLock l(&call_mu_);
761  call = call_.load(std::memory_order_relaxed);
762  if (call == nullptr) {
763  backlog_.send_initial_metadata_wanted = true;
764  return;
765  }
766  }
767  call->SendInitialMetadata();
768  }
772  void Finish(grpc::Status s) ABSL_LOCKS_EXCLUDED(call_mu_) {
773  ServerCallbackUnary* call = call_.load(std::memory_order_acquire);
774  if (call == nullptr) {
775  grpc::internal::MutexLock l(&call_mu_);
776  call = call_.load(std::memory_order_relaxed);
777  if (call == nullptr) {
778  backlog_.finish_wanted = true;
779  backlog_.status_wanted = std::move(s);
780  return;
781  }
782  }
783  call->Finish(std::move(s));
784  }
785 
787  virtual void OnSendInitialMetadataDone(bool /*ok*/) {}
788  void OnDone() override = 0;
789  void OnCancel() override {}
790 
791  private:
792  friend class ServerCallbackUnary;
793  // May be overridden by internal implementation details. This is not a public
794  // customization point.
795  virtual void InternalBindCall(ServerCallbackUnary* call)
796  ABSL_LOCKS_EXCLUDED(call_mu_) {
797  grpc::internal::MutexLock l(&call_mu_);
798 
799  if (GPR_UNLIKELY(backlog_.send_initial_metadata_wanted)) {
800  call->SendInitialMetadata();
801  }
802  if (GPR_UNLIKELY(backlog_.finish_wanted)) {
803  call->Finish(std::move(backlog_.status_wanted));
804  }
805  // Set call_ last so that other functions can use it lock-free
806  call_.store(call, std::memory_order_release);
807  }
808 
809  grpc::internal::Mutex call_mu_;
810  std::atomic<ServerCallbackUnary*> call_{nullptr};
811  struct PreBindBacklog {
812  bool send_initial_metadata_wanted = false;
813  bool finish_wanted = false;
814  grpc::Status status_wanted;
815  };
816  PreBindBacklog backlog_ ABSL_GUARDED_BY(call_mu_);
817 };
818 
819 namespace experimental {
821  public:
822  ServerSessionReactor() : session_(nullptr) {}
823  ~ServerSessionReactor() override = default;
824 
827  void StartVirtualRPCs() ABSL_LOCKS_EXCLUDED(session_mu_) {
828  ServerCallbackSession* session = session_.load(std::memory_order_acquire);
829  if (session == nullptr) {
830  grpc::internal::MutexLock l(&session_mu_);
831  session = session_.load(std::memory_order_relaxed);
832  if (session == nullptr) {
833  backlog_.send_initial_metadata_wanted = true;
834  return;
835  }
836  }
837  session->SendInitialMetadata();
838  }
839 
843  void Finish(grpc::Status s) ABSL_LOCKS_EXCLUDED(session_mu_) {
844  ServerCallbackSession* session = session_.load(std::memory_order_acquire);
845  if (session == nullptr) {
846  grpc::internal::MutexLock l(&session_mu_);
847  session = session_.load(std::memory_order_relaxed);
848  if (session == nullptr) {
849  backlog_.finish_wanted = true;
850  backlog_.status_wanted = std::move(s);
851  return;
852  }
853  }
854  session->Finish(std::move(s));
855  }
856 
864  absl::AnyInvocable<void(absl::Status)> on_shutdown)
865  ABSL_LOCKS_EXCLUDED(session_mu_) {
866  ServerCallbackSession* session = session_.load(std::memory_order_acquire);
867  if (session == nullptr) {
868  grpc::internal::MutexLock l(&session_mu_);
869  session = session_.load(std::memory_order_relaxed);
870  if (session == nullptr) {
871  backlog_.graceful_shutdown_wanted_callback = std::move(on_shutdown);
872  return;
873  }
874  }
875  session->InitiateGracefulShutdown(std::move(on_shutdown));
876  }
877 
879  virtual void OnSendInitialMetadataDone(bool /*ok*/) {}
880  void OnDone() override = 0;
881  void OnCancel() override {}
882 
883  private:
884  friend class ServerCallbackSession;
885  // May be overridden by internal implementation details. This is not a public
886  // customization point.
887  virtual void InternalBindSession(ServerCallbackSession* session)
888  ABSL_LOCKS_EXCLUDED(session_mu_) {
889  grpc::internal::MutexLock l(&session_mu_);
890 
891  if (GPR_UNLIKELY(backlog_.send_initial_metadata_wanted)) {
892  session->SendInitialMetadata();
893  }
894  if (GPR_UNLIKELY(backlog_.graceful_shutdown_wanted_callback != nullptr)) {
895  session->InitiateGracefulShutdown(
896  std::move(backlog_.graceful_shutdown_wanted_callback));
897  }
898  if (GPR_UNLIKELY(backlog_.finish_wanted)) {
899  session->Finish(std::move(backlog_.status_wanted));
900  }
901  // Set session_ last so that other functions can use it lock-free
902  session_.store(session, std::memory_order_release);
903  }
904 
905  grpc::internal::Mutex session_mu_;
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;
911  grpc::Status status_wanted;
912  };
913 
914  // Resolves a race condition where application code calls reactor methods
915  // before gRPC binds the session to the reactor.
916  // These early operations are stored here and executed during
917  // InternalBindSession().
918  PreBindBacklog backlog_ ABSL_GUARDED_BY(session_mu_);
919 };
920 } // namespace experimental
921 
922 namespace internal {
923 
924 template <class Base>
925 class FinishOnlyReactor : public Base {
926  public:
927  explicit FinishOnlyReactor(grpc::Status s) { this->Finish(std::move(s)); }
928  void OnDone() override { this->~FinishOnlyReactor(); }
929 };
930 
932 template <class Request>
934 template <class Response>
937 template <class Request, class Response>
940 
941 } // namespace internal
942 
943 namespace experimental {
944 namespace internal {
947 } // namespace internal
948 } // namespace experimental
949 
950 // TODO(vjpai): Remove namespace experimental when last known users are migrated
951 // off.
952 namespace experimental {
953 
954 template <class Request, class Response>
956 
957 } // namespace experimental
958 
959 } // namespace grpc
960 
961 #endif // GRPCPP_SUPPORT_SERVER_CALLBACK_H
grpc::ServerReadReactor
ServerReadReactor is the interface for a client-streaming RPC.
Definition: server_callback.h:212
grpc::ServerCallbackReaderWriter::SendInitialMetadata
virtual void SendInitialMetadata()=0
grpc::ServerCallbackReader::~ServerCallbackReader
~ServerCallbackReader() override
Definition: server_callback.h:263
grpc::MessageHolder< Request, Response >::set_response
void set_response(Response *response)
Definition: message_allocator.h:50
grpc::ServerReadReactor::OnCancel
void OnCancel() override
Definition: server_callback.h:595
grpc::ServerBidiReactor::StartSendInitialMetadata
void StartSendInitialMetadata() ABSL_LOCKS_EXCLUDED(stream_mu_)
Send any initial metadata stored in the RPC context.
Definition: server_callback.h:335
grpc::experimental::ServerSessionReactor::StartVirtualRPCs
void StartVirtualRPCs() ABSL_LOCKS_EXCLUDED(session_mu_)
StartVirtualRPCs is exactly like ServerBidiReactor's StartSendInitialMetadata.
Definition: server_callback.h:827
grpc::MessageHolder< Request, Response >::set_request
void set_request(Request *request)
Definition: message_allocator.h:49
message_allocator.h
grpc::ServerReadReactor::ServerReadReactor
ServerReadReactor()
Definition: server_callback.h:546
grpc::ServerReadReactor::OnDone
void OnDone() override=0
grpc::ServerCallbackReaderWriter::Finish
virtual void Finish(grpc::Status s)=0
grpc::ServerWriteReactor::StartWrite
void StartWrite(const Response *resp)
Definition: server_callback.h:651
grpc::Server
Represents a gRPC server.
Definition: server.h:58
grpc::ServerCallbackWriter::Finish
virtual void Finish(grpc::Status s)=0
grpc
An Alarm posts the user-provided tag to its associated completion queue or invokes the user-provided ...
Definition: alarm.h:33
grpc::ServerUnaryReactor::OnDone
void OnDone() override=0
grpc::ServerBidiReactor::StartWriteAndFinish
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
grpc::WriteOptions::set_last_message
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
grpc_core
Definition: context_types.h:20
grpc::internal::ServerReactor::~ServerReactor
virtual ~ServerReactor()=default
grpc::internal::ServerCallbackCall::MaybeCallOnCancel
void MaybeCallOnCancel()
Definition: server_callback.h:141
grpc::ServerWriteReactor::StartWrite
void StartWrite(const Response *resp, grpc::WriteOptions options) ABSL_LOCKS_EXCLUDED(writer_mu_)
Definition: server_callback.h:654
grpc::ServerWriteReactor::OnWriteDone
virtual void OnWriteDone(bool)
Definition: server_callback.h:706
grpc::ServerBidiReactor::ServerBidiReactor
ServerBidiReactor()
Definition: server_callback.h:329
grpc::ServerWriteReactor::OnSendInitialMetadataDone
virtual void OnSendInitialMetadataDone(bool)
The following notifications are exactly like ServerBidiReactor.
Definition: server_callback.h:705
grpc::ServerCallbackReader::Finish
virtual void Finish(grpc::Status s)=0
grpc::ServerReadReactor::StartRead
void StartRead(Request *req) ABSL_LOCKS_EXCLUDED(reader_mu_)
Definition: server_callback.h:563
grpc::internal::DefaultMessageHolder::DefaultMessageHolder
DefaultMessageHolder()
Definition: server_callback.h:193
grpc::experimental::ServerSessionReactor::~ServerSessionReactor
~ServerSessionReactor() override=default
grpc::ServerCallbackReaderWriter::~ServerCallbackReaderWriter
~ServerCallbackReaderWriter() override
Definition: server_callback.h:294
status.h
grpc::ServerCallbackUnary::BindReactor
void BindReactor(Reactor *reactor)
Definition: server_callback.h:236
grpc::internal::ServerCallbackCall::MaybeCallOnCancel
void MaybeCallOnCancel(ServerReactor *reactor)
Definition: server_callback.h:131
grpc::ServerWriteReactor::StartWriteLast
void StartWriteLast(const Response *resp, grpc::WriteOptions options)
Definition: server_callback.h:686
grpc::ServerReadReactor::~ServerReadReactor
~ServerReadReactor() override=default
grpc::experimental::ServerCallbackSession::InitiateGracefulShutdown
virtual void InitiateGracefulShutdown(absl::AnyInvocable< void(absl::Status)> on_shutdown)=0
grpc::ServerCallbackReaderWriter
Definition: server_callback.h:292
grpc::ServerBidiReactor::StartWrite
void StartWrite(const Response *resp)
Initiate a write operation.
Definition: server_callback.h:372
grpc::experimental::ServerCallbackSession::BindInnerServer
virtual void BindInnerServer(grpc::Server *inner_server)=0
grpc::ServerBidiReactor::OnCancel
void OnCancel() override
Notifies the application that this RPC has been cancelled.
Definition: server_callback.h:491
grpc::ServerCallbackReaderWriter::Write
virtual void Write(const Response *msg, grpc::WriteOptions options)=0
grpc::experimental::ServerCallbackSession::~ServerCallbackSession
~ServerCallbackSession() override
Definition: server_callback.h:244
grpc::ServerBidiReactor::StartRead
void StartRead(Request *req) ABSL_LOCKS_EXCLUDED(stream_mu_)
Initiate a read operation.
Definition: server_callback.h:353
grpc::Status
Did it work? If it didn't, why?
Definition: status.h:34
GPR_UNLIKELY
#define GPR_UNLIKELY(x)
Definition: port_platform.h:861
grpc::ServerWriteReactor::OnDone
void OnDone() override=0
grpc::ServerCallbackUnary::~ServerCallbackUnary
~ServerCallbackUnary() override
Definition: server_callback.h:228
grpc::ServerUnaryReactor::OnSendInitialMetadataDone
virtual void OnSendInitialMetadataDone(bool)
The following notifications are exactly like ServerBidiReactor.
Definition: server_callback.h:787
grpc::experimental::ServerSessionReactor::ServerSessionReactor
ServerSessionReactor()
Definition: server_callback.h:822
grpc::experimental::ServerSessionReactor
Definition: server_callback.h:820
grpc::internal::ServerCallbackCall
The base class of ServerCallbackUnary etc.
Definition: server_callback.h:97
grpc::experimental::ServerSessionReactor::OnSendInitialMetadataDone
virtual void OnSendInitialMetadataDone(bool)
The following notifications are exactly like ServerBidiReactor.
Definition: server_callback.h:879
grpc::ServerBidiReactor::OnWriteDone
virtual void OnWriteDone(bool)
Notifies the application that a StartWrite (or StartWriteLast) operation completed.
Definition: server_callback.h:481
grpc::ServerBidiReactor::StartWrite
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
callback_common.h
grpc::experimental::ServerSessionReactor::Finish
void Finish(grpc::Status s) ABSL_LOCKS_EXCLUDED(session_mu_)
Finish is similar to ServerBidiReactor except for one detail.
Definition: server_callback.h:843
grpc::ServerCallbackReaderWriter::WriteAndFinish
virtual void WriteAndFinish(const Response *msg, grpc::WriteOptions options, grpc::Status s)=0
grpc::ServerBidiReactor::OnDone
void OnDone() override=0
Notifies the application that all operations associated with this RPC have completed.
grpc_call
struct grpc_call grpc_call
A Call represents an RPC.
Definition: grpc_types.h:68
grpc::ServerBidiReactor
ServerBidiReactor is the interface for a bidirectional streaming RPC.
Definition: server_callback.h:216
grpc::internal::CallbackUnaryHandler
Definition: server_callback_handlers.h:38
grpc::ServerCallbackWriter::SendInitialMetadata
virtual void SendInitialMetadata()=0
grpc::ServerBidiReactor::~ServerBidiReactor
~ServerBidiReactor() override=default
grpc::internal::ReturnPreexistingErrors
bool ReturnPreexistingErrors()
grpc::internal::DefaultMessageHolder
Definition: server_callback.h:191
grpc::internal::FinishOnlyReactor::FinishOnlyReactor
FinishOnlyReactor(grpc::Status s)
Definition: server_callback.h:927
grpc::ServerBidiReactor::OnReadDone
virtual void OnReadDone(bool)
Notifies the application that a StartRead operation completed.
Definition: server_callback.h:474
grpc::internal::ServerReactor::OnDone
virtual void OnDone()=0
grpc::ServerWriteReactor::StartWriteAndFinish
void StartWriteAndFinish(const Response *resp, grpc::WriteOptions options, grpc::Status s) ABSL_LOCKS_EXCLUDED(writer_mu_)
Definition: server_callback.h:669
grpc::internal::CallbackServerStreamingHandler
Definition: server_callback_handlers.h:460
grpc::experimental::ServerCallbackSession
Definition: server_callback.h:242
grpc::experimental::internal::InitiateSessionGracefulShutdown
void InitiateSessionGracefulShutdown(grpc_core::Transport *transport, grpc_endpoint *endpoint, absl::AnyInvocable< void(absl::Status)> on_shutdown)
grpc::ServerCallbackReaderWriter::Read
virtual void Read(Request *msg)=0
grpc::ServerBidiReactor::Finish
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
grpc::ServerCallbackReader::BindReactor
void BindReactor(ServerReadReactor< Request > *reactor)
Definition: server_callback.h:269
grpc::ServerCallbackWriter::BindReactor
void BindReactor(ServerWriteReactor< Response > *reactor)
Definition: server_callback.h:286
grpc::internal::FinishOnlyReactor
Definition: server_context.h:90
grpc::ServerCallbackWriter::~ServerCallbackWriter
~ServerCallbackWriter() override
Definition: server_callback.h:277
grpc_call_run_in_event_engine
void grpc_call_run_in_event_engine(const grpc_call *call, absl::AnyInvocable< void()> cb)
grpc::ServerReadReactor::StartSendInitialMetadata
void StartSendInitialMetadata() ABSL_LOCKS_EXCLUDED(reader_mu_)
The following operation initiations are exactly like ServerBidiReactor.
Definition: server_callback.h:550
grpc::internal::ServerCallbackCall::~ServerCallbackCall
virtual ~ServerCallbackCall()
Definition: server_callback.h:99
grpc::WriteOptions
Per-message write options.
Definition: call_op_set.h:79
grpc::experimental::ServerCallbackSession::BindReactor
void BindReactor(Reactor *reactor)
Definition: server_callback.h:254
grpc::internal::ServerCallbackCall::MaybeDone
void MaybeDone(bool inline_ondone)
Definition: server_callback.h:123
grpc::ServerWriteReactor
ServerWriteReactor is the interface for a server-streaming RPC.
Definition: server_callback.h:214
grpc::ServerReadReactor::OnReadDone
virtual void OnReadDone(bool)
Definition: server_callback.h:593
grpc::ServerCallbackReader::Read
virtual void Read(Request *msg)=0
grpc::internal::CallbackBidiHandler
Definition: server_callback_handlers.h:695
grpc::internal::MutexLock
Definition: sync.h:80
grpc::ServerWriteReactor::OnCancel
void OnCancel() override
Definition: server_callback.h:708
grpc::ServerWriteReactor::~ServerWriteReactor
~ServerWriteReactor() override=default
grpc::ServerCallbackReaderWriter::BindReactor
void BindReactor(ServerBidiReactor< Request, Response > *reactor)
Definition: server_callback.h:304
config.h
grpc::ServerCallbackWriter
Definition: server_callback.h:275
call.h
grpc::experimental::ServerSessionReactor::OnDone
void OnDone() override=0
grpc::ServerWriteReactor::StartSendInitialMetadata
void StartSendInitialMetadata() ABSL_LOCKS_EXCLUDED(writer_mu_)
The following operation initiations are exactly like ServerBidiReactor.
Definition: server_callback.h:638
call_op_set.h
call.h
grpc::experimental::internal::UnimplementedSessionReactor
grpc::internal::FinishOnlyReactor< ServerSessionReactor > UnimplementedSessionReactor
Definition: server_callback.h:946
grpc::internal::ServerReactor::InternalInlineable
virtual bool InternalInlineable()
Definition: server_callback.h:83
grpc::ServerCallbackUnary::SendInitialMetadata
virtual void SendInitialMetadata()=0
grpc::ServerReadReactor::OnSendInitialMetadataDone
virtual void OnSendInitialMetadataDone(bool)
The following notifications are exactly like ServerBidiReactor.
Definition: server_callback.h:592
grpc::ServerBidiReactor::StartWriteLast
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
grpc::internal::FinishOnlyReactor::OnDone
void OnDone() override
Definition: server_callback.h:928
grpc::internal::ServerReactor
Definition: server_callback.h:73
grpc::experimental::ServerCallbackSession::SendInitialMetadata
virtual void SendInitialMetadata()=0
grpc::internal::CallbackClientStreamingHandler
Definition: server_callback_handlers.h:262
grpc::ServerUnaryReactor::OnCancel
void OnCancel() override
Definition: server_callback.h:789
grpc::internal::ServerCallbackCall::Ref
void Ref()
Increases the reference count.
Definition: server_callback.h:149
grpc::internal::ServerReactor::OnCancel
virtual void OnCancel()=0
grpc::ServerCallbackReader
Definition: server_callback.h:261
grpc::ServerReadReactor::Finish
void Finish(grpc::Status s) ABSL_LOCKS_EXCLUDED(reader_mu_)
Definition: server_callback.h:576
grpc::ServerCallbackWriter::Write
virtual void Write(const Response *msg, grpc::WriteOptions options)=0
grpc::MessageHolder
Definition: message_allocator.h:39
grpc::internal::Mutex
Definition: sync.h:57
grpc::ServerUnaryReactor::StartSendInitialMetadata
void StartSendInitialMetadata() ABSL_LOCKS_EXCLUDED(call_mu_)
StartSendInitialMetadata is exactly like ServerBidiReactor.
Definition: server_callback.h:757
grpc::ServerWriteReactor::Finish
void Finish(grpc::Status s) ABSL_LOCKS_EXCLUDED(writer_mu_)
Definition: server_callback.h:689
grpc::experimental::internal::BindSessionToInnerServer
void BindSessionToInnerServer(grpc_call *call, grpc::Server *inner_server, grpc_core::Transport **out_transport, grpc_endpoint **out_endpoint)
grpc::experimental::ServerSessionReactor::InitiateGracefulShutdown
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
grpc::ServerUnaryReactor::Finish
void Finish(grpc::Status s) ABSL_LOCKS_EXCLUDED(call_mu_)
Finish is similar to ServerBidiReactor except for one detail.
Definition: server_callback.h:772
grpc::ServerBidiReactor::OnSendInitialMetadataDone
virtual void OnSendInitialMetadataDone(bool)
Notifies the application that an explicit StartSendInitialMetadata operation completed.
Definition: server_callback.h:468
grpc::ServerUnaryReactor::~ServerUnaryReactor
~ServerUnaryReactor() override=default
grpc::ServerCallbackUnary
Definition: server_callback.h:226
grpc::ServerUnaryReactor::ServerUnaryReactor
ServerUnaryReactor()
Definition: server_callback.h:753
grpc::experimental::ServerCallbackSession::Finish
virtual void Finish(grpc::Status s)=0
grpc::ServerWriteReactor::ServerWriteReactor
ServerWriteReactor()
Definition: server_callback.h:634
grpc::ServerUnaryReactor
Definition: server_callback.h:751
grpc::protobuf::util::Status
::absl::Status Status
Definition: config_protobuf.h:107
grpc::ServerCallbackWriter::WriteAndFinish
virtual void WriteAndFinish(const Response *msg, grpc::WriteOptions options, grpc::Status s)=0
sync.h
grpc::ServerCallbackUnary::Finish
virtual void Finish(grpc::Status s)=0
grpc::ServerCallbackReader::SendInitialMetadata
virtual void SendInitialMetadata()=0
grpc::internal::ServerCallbackCall::MaybeDone
void MaybeDone()
Definition: server_callback.h:117
grpc::experimental::ServerSessionReactor::OnCancel
void OnCancel() override
Definition: server_callback.h:881
grpc::internal::DefaultMessageHolder::Release
void Release() override
Definition: server_callback.h:197
grpc::experimental::ServerBidiReactor
::grpc::ServerBidiReactor< Request, Response > ServerBidiReactor
Definition: server_callback.h:955