GRPC C++  1.83.0
server_callback_handlers.h
Go to the documentation of this file.
1 //
2 //
3 // Copyright 2019 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 #ifndef GRPCPP_IMPL_SERVER_CALLBACK_HANDLERS_H
19 #define GRPCPP_IMPL_SERVER_CALLBACK_HANDLERS_H
20 
21 #include <grpc/grpc.h>
22 #include <grpc/impl/call.h>
25 #include <grpcpp/server.h>
26 #include <grpcpp/server_context.h>
29 #include <grpcpp/support/status.h>
30 
31 #include "absl/log/absl_check.h"
32 #include "absl/status/status.h"
33 
34 namespace grpc {
35 namespace internal {
36 
37 template <class RequestType, class ResponseType>
39  public:
42  const RequestType*, ResponseType*)>
43  get_reactor)
44  : get_reactor_(std::move(get_reactor)) {}
45 
48  allocator_ = allocator;
49  }
50 
51  void RunHandler(const HandlerParameter& param) final {
52  // Arena allocate a controller structure (that includes request/response)
53  grpc_call_ref(param.call->call());
54  auto* allocator_state =
56  param.internal_data);
57 
58  auto* call = new (grpc_call_arena_alloc(param.call->call(),
59  sizeof(ServerCallbackUnaryImpl)))
60  ServerCallbackUnaryImpl(
61  static_cast<grpc::CallbackServerContext*>(param.server_context),
62  param.call, allocator_state, param.call_requester);
63  param.server_context->BeginCompletionOp(
64  param.call, [call](bool) { call->MaybeDone(); }, call);
65 
66  ServerUnaryReactor* reactor = nullptr;
67  if (ReturnPreexistingErrors() && !param.status.ok()) {
68  reactor = new (grpc_call_arena_alloc(param.call->call(),
70  UnimplementedUnaryReactor(param.status);
71  } else if (param.status.ok()) {
72  reactor = grpc::internal::CatchingReactorGetter<ServerUnaryReactor>(
73  get_reactor_,
74  static_cast<grpc::CallbackServerContext*>(param.server_context),
75  call->request(), call->response());
76  }
77 
78  if (reactor == nullptr) {
79  // if deserialization or reactor creator failed, we need to fail the call
80  reactor = new (grpc_call_arena_alloc(param.call->call(),
84  }
85 
87  call->SetupReactor(reactor);
88  }
89 
91  grpc::Status* status, void** handler_data) final {
92  grpc::ByteBuffer buf;
93  buf.set_buffer(req);
94  RequestType* request = nullptr;
96  if (allocator_ != nullptr) {
97  allocator_state = allocator_->AllocateMessages();
98  } else {
99  allocator_state = new (grpc_call_arena_alloc(
102  }
103  *handler_data = allocator_state;
104  request = allocator_state->request();
105  *status = grpc::Deserialize(&buf, request);
106  buf.Release();
107  if (status->ok()) {
108  return request;
109  }
110  return nullptr;
111  }
112 
113  private:
115  const RequestType*, ResponseType*)>
116  get_reactor_;
117  MessageAllocator<RequestType, ResponseType>* allocator_ = nullptr;
118 
119  class ServerCallbackUnaryImpl : public ServerCallbackUnary {
120  public:
121  void Finish(grpc::Status s) override {
122  // A callback that only contains a call to MaybeDone can be run as an
123  // inline callback regardless of whether or not OnDone is inlineable
124  // because if the actual OnDone callback needs to be scheduled, MaybeDone
125  // is responsible for dispatching to an EventEngine thread if needed.
126  // Thus, when setting up the finish_tag_, we can set its own callback to
127  // inlineable.
128  finish_tag_.Set(
129  call_.call(),
130  [this](bool) {
131  this->MaybeDone(
132  reactor_.load(std::memory_order_relaxed)->InternalInlineable());
133  },
134  &finish_ops_, /*can_inline=*/true);
135  finish_ops_.set_core_cq_tag(&finish_tag_);
136 
137  if (!ctx_->sent_initial_metadata_) {
138  finish_ops_.SendInitialMetadata(&ctx_->initial_metadata_,
139  ctx_->initial_metadata_flags());
140  if (ctx_->compression_level_set()) {
141  finish_ops_.set_compression_level(ctx_->compression_level());
142  }
143  ctx_->MarkInitialMetadataSent();
144  }
145  // The response is dropped if the status is not OK.
146  if (s.ok()) {
147  finish_ops_.ServerSendStatus(
148  &ctx_->trailing_metadata_,
149  finish_ops_.SendMessagePtr(response(), ctx_->memory_allocator()));
150  } else {
151  finish_ops_.ServerSendStatus(&ctx_->trailing_metadata_, s);
152  }
153  finish_ops_.set_core_cq_tag(&finish_tag_);
154  finish_ops_.FillOps(&call_);
155  }
156 
157  void SendInitialMetadata() override {
158  ABSL_CHECK(!ctx_->sent_initial_metadata_);
159  this->Ref();
160  // The callback for this function should not be marked inline because it
161  // is directly invoking a user-controlled reaction
162  // (OnSendInitialMetadataDone). Thus it must be dispatched to an
163  // EventEngine thread. However, any OnDone needed after that can be
164  // inlined because it is already running on an EventEngine thread.
165  meta_tag_.Set(
166  call_.call(),
167  [this](bool ok) {
168  ServerUnaryReactor* reactor =
169  reactor_.load(std::memory_order_relaxed);
170  reactor->OnSendInitialMetadataDone(ok);
171  this->MaybeDone(/*inlineable_ondone=*/true);
172  },
173  &meta_ops_, /*can_inline=*/false);
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());
178  }
179  ctx_->MarkInitialMetadataSent();
180  meta_ops_.set_core_cq_tag(&meta_tag_);
181  meta_ops_.FillOps(&call_);
182  }
183 
184  private:
185  friend class CallbackUnaryHandler<RequestType, ResponseType>;
186 
187  ServerCallbackUnaryImpl(
189  MessageHolder<RequestType, ResponseType>* allocator_state,
190  std::function<void()> call_requester)
191  : ctx_(ctx),
192  call_(*call),
193  allocator_state_(allocator_state),
194  call_requester_(std::move(call_requester)) {
195  ctx_->set_message_allocator_state(allocator_state);
196  }
197 
198  grpc_call* call() override { return call_.call(); }
199 
204  void SetupReactor(ServerUnaryReactor* reactor) {
205  reactor_.store(reactor, std::memory_order_relaxed);
206  this->BindReactor(reactor);
207  this->MaybeCallOnCancel(reactor);
208  this->MaybeDone(reactor->InternalInlineable());
209  }
210 
211  const RequestType* request() { return allocator_state_->request(); }
212  ResponseType* response() { return allocator_state_->response(); }
213 
214  void CallOnDone() override {
215  reactor_.load(std::memory_order_relaxed)->OnDone();
216  grpc_call* call = call_.call();
217  auto call_requester = std::move(call_requester_);
218  allocator_state_->Release();
219  if (ctx_->context_allocator() != nullptr) {
220  ctx_->context_allocator()->Release(ctx_);
221  }
222  this->~ServerCallbackUnaryImpl(); // explicitly call destructor
223  grpc_call_unref(call);
224  call_requester();
225  }
226 
227  ServerReactor* reactor() override {
228  return reactor_.load(std::memory_order_relaxed);
229  }
230 
232  meta_ops_;
237  finish_ops_;
239 
240  grpc::CallbackServerContext* const ctx_;
241  grpc::internal::Call call_;
242  MessageHolder<RequestType, ResponseType>* const allocator_state_;
243  std::function<void()> call_requester_;
244  // reactor_ can always be loaded/stored with relaxed memory ordering because
245  // its value is only set once, independently of other data in the object,
246  // and the loads that use it will always actually come provably later even
247  // though they are from different threads since they are triggered by
248  // actions initiated only by the setting up of the reactor_ variable. In
249  // a sense, it's a delayed "const": it gets its value from the SetupReactor
250  // method (not the constructor, so it's not a true const), but it doesn't
251  // change after that and it only gets used by actions caused, directly or
252  // indirectly, by that setup. This comment also applies to the reactor_
253  // variables of the other streaming objects in this file.
254  std::atomic<ServerUnaryReactor*> reactor_;
255  // callbacks_outstanding_ follows a refcount pattern
256  std::atomic<intptr_t> callbacks_outstanding_{
257  3}; // reserve for start, Finish, and CompletionOp
258  };
259 };
260 
261 template <class RequestType, class ResponseType>
263  public:
265  std::function<ServerReadReactor<RequestType>*(
266  grpc::CallbackServerContext*, ResponseType*)>
267  get_reactor)
268  : get_reactor_(std::move(get_reactor)) {}
269  void RunHandler(const HandlerParameter& param) final {
270  // Arena allocate a reader structure (that includes response)
271  grpc_call_ref(param.call->call());
272 
273  auto* reader = new (grpc_call_arena_alloc(param.call->call(),
274  sizeof(ServerCallbackReaderImpl)))
275  ServerCallbackReaderImpl(
276  static_cast<grpc::CallbackServerContext*>(param.server_context),
277  param.call, param.call_requester);
278  // Inlineable OnDone can be false in the CompletionOp callback because there
279  // is no read reactor that has an inlineable OnDone; this only applies to
280  // the DefaultReactor (which is unary).
281  param.server_context->BeginCompletionOp(
282  param.call,
283  [reader](bool) { reader->MaybeDone(/*inlineable_ondone=*/false); },
284  reader);
285 
286  ServerReadReactor<RequestType>* reactor = nullptr;
287  if (ReturnPreexistingErrors() && !param.status.ok()) {
288  reactor = new (grpc_call_arena_alloc(
289  param.call->call(), sizeof(UnimplementedReadReactor<RequestType>)))
291  } else if (param.status.ok()) {
292  reactor =
293  grpc::internal::CatchingReactorGetter<ServerReadReactor<RequestType>>(
294  get_reactor_,
295  static_cast<grpc::CallbackServerContext*>(param.server_context),
296  reader->response());
297  }
298 
299  if (reactor == nullptr) {
300  // if deserialization or reactor creator failed, we need to fail the call
301  reactor = new (grpc_call_arena_alloc(
302  param.call->call(), sizeof(UnimplementedReadReactor<RequestType>)))
305  }
306 
307  reader->SetupReactor(reactor);
308  }
309 
310  private:
311  std::function<ServerReadReactor<RequestType>*(grpc::CallbackServerContext*,
312  ResponseType*)>
313  get_reactor_;
314 
315  class ServerCallbackReaderImpl : public ServerCallbackReader<RequestType> {
316  public:
317  void Finish(grpc::Status s) override {
318  // A finish tag with only MaybeDone can have its callback inlined
319  // regardless even if OnDone is not inlineable because this callback just
320  // checks a ref and then decides whether or not to dispatch OnDone.
321  finish_tag_.Set(
322  call_.call(),
323  [this](bool) {
324  // Inlineable OnDone can be false here because there is
325  // no read reactor that has an inlineable OnDone; this
326  // only applies to the DefaultReactor (which is unary).
327  this->MaybeDone(/*inlineable_ondone=*/false);
328  },
329  &finish_ops_, /*can_inline=*/true);
330  if (!ctx_->sent_initial_metadata_) {
331  finish_ops_.SendInitialMetadata(&ctx_->initial_metadata_,
332  ctx_->initial_metadata_flags());
333  if (ctx_->compression_level_set()) {
334  finish_ops_.set_compression_level(ctx_->compression_level());
335  }
336  ctx_->MarkInitialMetadataSent();
337  }
338  // The response is dropped if the status is not OK.
339  if (s.ok()) {
340  finish_ops_.ServerSendStatus(
341  &ctx_->trailing_metadata_,
342  finish_ops_.SendMessagePtr(&resp_, ctx_->memory_allocator()));
343  } else {
344  finish_ops_.ServerSendStatus(&ctx_->trailing_metadata_, s);
345  }
346  finish_ops_.set_core_cq_tag(&finish_tag_);
347  finish_ops_.FillOps(&call_);
348  }
349 
350  void SendInitialMetadata() override {
351  ABSL_CHECK(!ctx_->sent_initial_metadata_);
352  this->Ref();
353  // The callback for this function should not be inlined because it invokes
354  // a user-controlled reaction, but any resulting OnDone can be inlined in
355  // the EventEngine thread to which this callback is dispatched.
356  meta_tag_.Set(
357  call_.call(),
358  [this](bool ok) {
359  ServerReadReactor<RequestType>* reactor =
360  reactor_.load(std::memory_order_relaxed);
361  reactor->OnSendInitialMetadataDone(ok);
362  this->MaybeDone(/*inlineable_ondone=*/true);
363  },
364  &meta_ops_, /*can_inline=*/false);
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());
369  }
370  ctx_->MarkInitialMetadataSent();
371  meta_ops_.set_core_cq_tag(&meta_tag_);
372  meta_ops_.FillOps(&call_);
373  }
374 
375  void Read(RequestType* req) override {
376  this->Ref();
377  read_ops_.RecvMessage(req);
378  read_ops_.FillOps(&call_);
379  }
380 
381  private:
382  friend class CallbackClientStreamingHandler<RequestType, ResponseType>;
383 
384  ServerCallbackReaderImpl(grpc::CallbackServerContext* ctx,
385  grpc::internal::Call* call,
386  std::function<void()> call_requester)
387  : ctx_(ctx), call_(*call), call_requester_(std::move(call_requester)) {}
388 
389  grpc_call* call() override { return call_.call(); }
390 
391  void SetupReactor(ServerReadReactor<RequestType>* reactor) {
392  reactor_.store(reactor, std::memory_order_relaxed);
393  // The callback for this function should not be inlined because it invokes
394  // a user-controlled reaction, but any resulting OnDone can be inlined in
395  // the EventEngine thread to which this callback is dispatched.
396  read_tag_.Set(
397  call_.call(),
398  [this, reactor](bool ok) {
399  if (GPR_UNLIKELY(!ok)) {
400  ctx_->MaybeMarkCancelledOnRead();
401  }
402  reactor->OnReadDone(ok);
403  this->MaybeDone(/*inlineable_ondone=*/true);
404  },
405  &read_ops_, /*can_inline=*/false);
406  read_ops_.set_core_cq_tag(&read_tag_);
407  this->BindReactor(reactor);
408  this->MaybeCallOnCancel(reactor);
409  // Inlineable OnDone can be false here because there is no read
410  // reactor that has an inlineable OnDone; this only applies to the
411  // DefaultReactor (which is unary).
412  this->MaybeDone(/*inlineable_ondone=*/false);
413  }
414 
415  ~ServerCallbackReaderImpl() {}
416 
417  ResponseType* response() { return &resp_; }
418 
419  void CallOnDone() override {
420  reactor_.load(std::memory_order_relaxed)->OnDone();
421  grpc_call* call = call_.call();
422  auto call_requester = std::move(call_requester_);
423  if (ctx_->context_allocator() != nullptr) {
424  ctx_->context_allocator()->Release(ctx_);
425  }
426  this->~ServerCallbackReaderImpl(); // explicitly call destructor
427  grpc_call_unref(call);
428  call_requester();
429  }
430 
431  ServerReactor* reactor() override {
432  return reactor_.load(std::memory_order_relaxed);
433  }
434 
436  meta_ops_;
441  finish_ops_;
444  read_ops_;
446 
447  grpc::CallbackServerContext* const ctx_;
448  grpc::internal::Call call_;
449  ResponseType resp_;
450  std::function<void()> call_requester_;
451  // The memory ordering of reactor_ follows ServerCallbackUnaryImpl.
452  std::atomic<ServerReadReactor<RequestType>*> reactor_;
453  // callbacks_outstanding_ follows a refcount pattern
454  std::atomic<intptr_t> callbacks_outstanding_{
455  3}; // reserve for OnStarted, Finish, and CompletionOp
456  };
457 };
458 
459 template <class RequestType, class ResponseType>
461  public:
463  std::function<ServerWriteReactor<ResponseType>*(
464  grpc::CallbackServerContext*, const RequestType*)>
465  get_reactor)
466  : get_reactor_(std::move(get_reactor)) {}
467  void RunHandler(const HandlerParameter& param) final {
468  // Arena allocate a writer structure
469  grpc_call_ref(param.call->call());
470 
471  auto* writer = new (grpc_call_arena_alloc(param.call->call(),
472  sizeof(ServerCallbackWriterImpl)))
473  ServerCallbackWriterImpl(
474  static_cast<grpc::CallbackServerContext*>(param.server_context),
475  param.call, static_cast<RequestType*>(param.request),
476  param.call_requester);
477  // Inlineable OnDone can be false in the CompletionOp callback because there
478  // is no write reactor that has an inlineable OnDone; this only applies to
479  // the DefaultReactor (which is unary).
480  param.server_context->BeginCompletionOp(
481  param.call,
482  [writer](bool) { writer->MaybeDone(/*inlineable_ondone=*/false); },
483  writer);
484 
485  ServerWriteReactor<ResponseType>* reactor = nullptr;
486  if (ReturnPreexistingErrors() && !param.status.ok()) {
487  reactor = new (grpc_call_arena_alloc(
488  param.call->call(), sizeof(UnimplementedWriteReactor<ResponseType>)))
490  } else if (param.status.ok()) {
493  get_reactor_,
494  static_cast<grpc::CallbackServerContext*>(param.server_context),
495  writer->request());
496  }
497  if (reactor == nullptr) {
498  // if deserialization or reactor creator failed, we need to fail the call
499  reactor = new (grpc_call_arena_alloc(
500  param.call->call(), sizeof(UnimplementedWriteReactor<ResponseType>)))
503  }
504 
505  writer->SetupReactor(reactor);
506  }
507 
509  grpc::Status* status, void** /*handler_data*/) final {
510  grpc::ByteBuffer buf;
511  buf.set_buffer(req);
512  auto* request =
513  new (grpc_call_arena_alloc(call, sizeof(RequestType))) RequestType();
514  *status = grpc::Deserialize(&buf, request);
515  buf.Release();
516  if (status->ok()) {
517  return request;
518  }
519  request->~RequestType();
520  return nullptr;
521  }
522 
523  private:
524  std::function<ServerWriteReactor<ResponseType>*(grpc::CallbackServerContext*,
525  const RequestType*)>
526  get_reactor_;
527 
528  class ServerCallbackWriterImpl : public ServerCallbackWriter<ResponseType> {
529  public:
530  void Finish(grpc::Status s) override {
531  // A finish tag with only MaybeDone can have its callback inlined
532  // regardless even if OnDone is not inlineable because this callback just
533  // checks a ref and then decides whether or not to dispatch OnDone.
534  finish_tag_.Set(
535  call_.call(),
536  [this](bool) {
537  // Inlineable OnDone can be false here because there is
538  // no write reactor that has an inlineable OnDone; this
539  // only applies to the DefaultReactor (which is unary).
540  this->MaybeDone(/*inlineable_ondone=*/false);
541  },
542  &finish_ops_, /*can_inline=*/true);
543  finish_ops_.set_core_cq_tag(&finish_tag_);
544 
545  if (!ctx_->sent_initial_metadata_) {
546  finish_ops_.SendInitialMetadata(&ctx_->initial_metadata_,
547  ctx_->initial_metadata_flags());
548  if (ctx_->compression_level_set()) {
549  finish_ops_.set_compression_level(ctx_->compression_level());
550  }
551  ctx_->MarkInitialMetadataSent();
552  }
553  finish_ops_.ServerSendStatus(&ctx_->trailing_metadata_, s);
554  finish_ops_.FillOps(&call_);
555  }
556 
557  void SendInitialMetadata() override {
558  ABSL_CHECK(!ctx_->sent_initial_metadata_);
559  this->Ref();
560  // The callback for this function should not be inlined because it invokes
561  // a user-controlled reaction, but any resulting OnDone can be inlined in
562  // the EventEngine thread to which this callback is dispatched.
563  meta_tag_.Set(
564  call_.call(),
565  [this](bool ok) {
566  ServerWriteReactor<ResponseType>* reactor =
567  reactor_.load(std::memory_order_relaxed);
568  reactor->OnSendInitialMetadataDone(ok);
569  this->MaybeDone(/*inlineable_ondone=*/true);
570  },
571  &meta_ops_, /*can_inline=*/false);
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());
576  }
577  ctx_->MarkInitialMetadataSent();
578  meta_ops_.set_core_cq_tag(&meta_tag_);
579  meta_ops_.FillOps(&call_);
580  }
581 
582  void Write(const ResponseType* resp, grpc::WriteOptions options) override {
583  this->Ref();
584  if (options.is_last_message()) {
585  options.set_buffer_hint();
586  }
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());
592  }
593  ctx_->MarkInitialMetadataSent();
594  }
595  // TODO(vjpai): don't assert
596  ABSL_CHECK(
597  write_ops_.SendMessagePtr(resp, options, ctx_->memory_allocator())
598  .ok());
599  write_ops_.FillOps(&call_);
600  }
601 
602  void WriteAndFinish(const ResponseType* resp, grpc::WriteOptions options,
603  grpc::Status s) override {
604  // This combines the write into the finish callback
605  // TODO(vjpai): don't assert
606  ABSL_CHECK(
607  finish_ops_.SendMessagePtr(resp, options, ctx_->memory_allocator())
608  .ok());
609  Finish(std::move(s));
610  }
611 
612  private:
613  friend class CallbackServerStreamingHandler<RequestType, ResponseType>;
614 
615  ServerCallbackWriterImpl(grpc::CallbackServerContext* ctx,
616  grpc::internal::Call* call, const RequestType* req,
617  std::function<void()> call_requester)
618  : ctx_(ctx),
619  call_(*call),
620  req_(req),
621  call_requester_(std::move(call_requester)) {}
622 
623  grpc_call* call() override { return call_.call(); }
624 
625  void SetupReactor(ServerWriteReactor<ResponseType>* reactor) {
626  reactor_.store(reactor, std::memory_order_relaxed);
627  // The callback for this function should not be inlined because it invokes
628  // a user-controlled reaction, but any resulting OnDone can be inlined in
629  // the EventEngine thread to which this callback is dispatched.
630  write_tag_.Set(
631  call_.call(),
632  [this, reactor](bool ok) {
633  reactor->OnWriteDone(ok);
634  this->MaybeDone(/*inlineable_ondone=*/true);
635  },
636  &write_ops_, /*can_inline=*/false);
637  write_ops_.set_core_cq_tag(&write_tag_);
638  this->BindReactor(reactor);
639  this->MaybeCallOnCancel(reactor);
640  // Inlineable OnDone can be false here because there is no write
641  // reactor that has an inlineable OnDone; this only applies to the
642  // DefaultReactor (which is unary).
643  this->MaybeDone(/*inlineable_ondone=*/false);
644  }
645  ~ServerCallbackWriterImpl() {
646  if (req_ != nullptr) {
647  req_->~RequestType();
648  }
649  }
650 
651  const RequestType* request() { return req_; }
652 
653  void CallOnDone() override {
654  reactor_.load(std::memory_order_relaxed)->OnDone();
655  grpc_call* call = call_.call();
656  auto call_requester = std::move(call_requester_);
657  if (ctx_->context_allocator() != nullptr) {
658  ctx_->context_allocator()->Release(ctx_);
659  }
660  this->~ServerCallbackWriterImpl(); // explicitly call destructor
661  grpc_call_unref(call);
662  call_requester();
663  }
664 
665  ServerReactor* reactor() override {
666  return reactor_.load(std::memory_order_relaxed);
667  }
668 
670  meta_ops_;
675  finish_ops_;
679  write_ops_;
681 
682  grpc::CallbackServerContext* const ctx_;
683  grpc::internal::Call call_;
684  const RequestType* req_;
685  std::function<void()> call_requester_;
686  // The memory ordering of reactor_ follows ServerCallbackUnaryImpl.
687  std::atomic<ServerWriteReactor<ResponseType>*> reactor_;
688  // callbacks_outstanding_ follows a refcount pattern
689  std::atomic<intptr_t> callbacks_outstanding_{
690  3}; // reserve for OnStarted, Finish, and CompletionOp
691  };
692 };
693 
694 template <class RequestType, class ResponseType>
696  public:
700  get_reactor)
701  : get_reactor_(std::move(get_reactor)) {}
702  void RunHandler(const HandlerParameter& param) final {
703  grpc_call_ref(param.call->call());
704 
705  auto* stream = new (grpc_call_arena_alloc(
706  param.call->call(), sizeof(ServerCallbackReaderWriterImpl)))
707  ServerCallbackReaderWriterImpl(
708  static_cast<grpc::CallbackServerContext*>(param.server_context),
709  param.call, param.call_requester);
710  // Inlineable OnDone can be false in the CompletionOp callback because there
711  // is no bidi reactor that has an inlineable OnDone; this only applies to
712  // the DefaultReactor (which is unary).
713  param.server_context->BeginCompletionOp(
714  param.call,
715  [stream](bool) { stream->MaybeDone(/*inlineable_ondone=*/false); },
716  stream);
717 
719  if (ReturnPreexistingErrors() && !param.status.ok()) {
720  reactor = new (grpc_call_arena_alloc(
721  param.call->call(),
724  } else if (param.status.ok()) {
727  get_reactor_,
728  static_cast<grpc::CallbackServerContext*>(param.server_context));
729  }
730 
731  if (reactor == nullptr) {
732  // if deserialization or reactor creator failed, we need to fail the call
733  reactor = new (grpc_call_arena_alloc(
734  param.call->call(),
738  }
739 
740  stream->SetupReactor(reactor);
741  }
742 
743  private:
744  std::function<ServerBidiReactor<RequestType, ResponseType>*(
746  get_reactor_;
747 
748  class ServerCallbackReaderWriterImpl
749  : public ServerCallbackReaderWriter<RequestType, ResponseType> {
750  public:
751  void Finish(grpc::Status s) override {
752  // A finish tag with only MaybeDone can have its callback inlined
753  // regardless even if OnDone is not inlineable because this callback just
754  // checks a ref and then decides whether or not to dispatch OnDone.
755  finish_tag_.Set(
756  call_.call(),
757  [this](bool) {
758  // Inlineable OnDone can be false here because there is
759  // no bidi reactor that has an inlineable OnDone; this
760  // only applies to the DefaultReactor (which is unary).
761  this->MaybeDone(/*inlineable_ondone=*/false);
762  },
763  &finish_ops_, /*can_inline=*/true);
764  finish_ops_.set_core_cq_tag(&finish_tag_);
765 
766  if (!ctx_->sent_initial_metadata_) {
767  finish_ops_.SendInitialMetadata(&ctx_->initial_metadata_,
768  ctx_->initial_metadata_flags());
769  if (ctx_->compression_level_set()) {
770  finish_ops_.set_compression_level(ctx_->compression_level());
771  }
772  ctx_->MarkInitialMetadataSent();
773  }
774  finish_ops_.ServerSendStatus(&ctx_->trailing_metadata_, s);
775  finish_ops_.FillOps(&call_);
776  }
777 
778  void SendInitialMetadata() override {
779  ABSL_CHECK(!ctx_->sent_initial_metadata_);
780  this->Ref();
781  // The callback for this function should not be inlined because it invokes
782  // a user-controlled reaction, but any resulting OnDone can be inlined in
783  // the EventEngine thread to which this callback is dispatched.
784  meta_tag_.Set(
785  call_.call(),
786  [this](bool ok) {
787  ServerBidiReactor<RequestType, ResponseType>* reactor =
788  reactor_.load(std::memory_order_relaxed);
789  reactor->OnSendInitialMetadataDone(ok);
790  this->MaybeDone(/*inlineable_ondone=*/true);
791  },
792  &meta_ops_, /*can_inline=*/false);
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());
797  }
798  ctx_->MarkInitialMetadataSent();
799  meta_ops_.set_core_cq_tag(&meta_tag_);
800  meta_ops_.FillOps(&call_);
801  }
802 
803  void Write(const ResponseType* resp, grpc::WriteOptions options) override {
804  this->Ref();
805  if (options.is_last_message()) {
806  options.set_buffer_hint();
807  }
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());
813  }
814  ctx_->MarkInitialMetadataSent();
815  }
816  // TODO(vjpai): don't assert
817  ABSL_CHECK(
818  write_ops_.SendMessagePtr(resp, options, ctx_->memory_allocator())
819  .ok());
820  write_ops_.FillOps(&call_);
821  }
822 
823  void WriteAndFinish(const ResponseType* resp, grpc::WriteOptions options,
824  grpc::Status s) override {
825  // TODO(vjpai): don't assert
826  ABSL_CHECK(
827  finish_ops_.SendMessagePtr(resp, options, ctx_->memory_allocator())
828  .ok());
829  Finish(std::move(s));
830  }
831 
832  void Read(RequestType* req) override {
833  this->Ref();
834  read_ops_.RecvMessage(req);
835  read_ops_.FillOps(&call_);
836  }
837 
838  private:
839  friend class CallbackBidiHandler<RequestType, ResponseType>;
840 
841  ServerCallbackReaderWriterImpl(grpc::CallbackServerContext* ctx,
842  grpc::internal::Call* call,
843  std::function<void()> call_requester)
844  : ctx_(ctx), call_(*call), call_requester_(std::move(call_requester)) {}
845 
846  grpc_call* call() override { return call_.call(); }
847 
848  void SetupReactor(ServerBidiReactor<RequestType, ResponseType>* reactor) {
849  reactor_.store(reactor, std::memory_order_relaxed);
850  // The callbacks for these functions should not be inlined because they
851  // invoke user-controlled reactions, but any resulting OnDones can be
852  // inlined in the EventEngine thread to which a callback is dispatched.
853  write_tag_.Set(
854  call_.call(),
855  [this, reactor](bool ok) {
856  reactor->OnWriteDone(ok);
857  this->MaybeDone(/*inlineable_ondone=*/true);
858  },
859  &write_ops_, /*can_inline=*/false);
860  write_ops_.set_core_cq_tag(&write_tag_);
861  read_tag_.Set(
862  call_.call(),
863  [this, reactor](bool ok) {
864  if (GPR_UNLIKELY(!ok)) {
865  ctx_->MaybeMarkCancelledOnRead();
866  }
867  reactor->OnReadDone(ok);
868  this->MaybeDone(/*inlineable_ondone=*/true);
869  },
870  &read_ops_, /*can_inline=*/false);
871  read_ops_.set_core_cq_tag(&read_tag_);
872  this->BindReactor(reactor);
873  this->MaybeCallOnCancel(reactor);
874  // Inlineable OnDone can be false here because there is no bidi
875  // reactor that has an inlineable OnDone; this only applies to the
876  // DefaultReactor (which is unary).
877  this->MaybeDone(/*inlineable_ondone=*/false);
878  }
879 
880  void CallOnDone() override {
881  reactor_.load(std::memory_order_relaxed)->OnDone();
882  grpc_call* call = call_.call();
883  auto call_requester = std::move(call_requester_);
884  if (ctx_->context_allocator() != nullptr) {
885  ctx_->context_allocator()->Release(ctx_);
886  }
887  this->~ServerCallbackReaderWriterImpl(); // explicitly call destructor
888  grpc_call_unref(call);
889  call_requester();
890  }
891 
892  ServerReactor* reactor() override {
893  return reactor_.load(std::memory_order_relaxed);
894  }
895 
897  meta_ops_;
902  finish_ops_;
906  write_ops_;
909  read_ops_;
911 
912  grpc::CallbackServerContext* const ctx_;
913  grpc::internal::Call call_;
914  std::function<void()> call_requester_;
915  // The memory ordering of reactor_ follows ServerCallbackUnaryImpl.
916  std::atomic<ServerBidiReactor<RequestType, ResponseType>*> reactor_;
917  // callbacks_outstanding_ follows a refcount pattern
918  std::atomic<intptr_t> callbacks_outstanding_{
919  3}; // reserve for OnStarted, Finish, and CompletionOp
920  };
921 };
922 
923 } // namespace internal
924 
925 namespace experimental {
926 namespace internal {
927 
928 template <class RequestType>
930  public:
933  grpc::CallbackServerContext*, const RequestType*)>
934  get_reactor,
935  grpc::Service* service = nullptr)
936  : get_reactor_(std::move(get_reactor)), service_(service) {
937  ABSL_CHECK(service_ != nullptr && service_->is_virtual_service_);
938  }
939 
940  void RunHandler(const HandlerParameter& param) final {
941  // Arena allocate a controller structure (that includes request/response)
942  grpc_call_ref(param.call->call());
943  auto* allocator_state =
945  param.internal_data);
946 
947  grpc::Server* inner_server = nullptr;
948  inner_server = static_cast<grpc::Server*>(service_->server_);
949  ABSL_CHECK(inner_server != nullptr);
950 
951  auto* call = new (grpc_call_arena_alloc(param.call->call(),
952  sizeof(ServerCallbackSessionImpl)))
953  ServerCallbackSessionImpl(
954  static_cast<grpc::CallbackServerContext*>(param.server_context),
955  param.call, allocator_state, param.call_requester, inner_server);
956 
957  param.server_context->BeginCompletionOp(
958  param.call, [call](bool) { call->MaybeDone(); }, call);
959 
960  grpc::experimental::ServerSessionReactor* reactor = nullptr;
961  if (param.status.ok()) {
964  get_reactor_,
965  static_cast<grpc::CallbackServerContext*>(param.server_context),
966  call->request());
967  }
968 
969  if (reactor == nullptr) {
970  // if deserialization or reactor creator failed, we need to fail the call.
971  reactor = new (grpc_call_arena_alloc(param.call->call(),
975  }
976 
978  call->SetupReactor(reactor);
979  }
980 
982  grpc::Status* status, void** handler_data) final {
983  grpc::ByteBuffer buf;
984  buf.set_buffer(req);
985  RequestType* request = nullptr;
987  allocator_state = new (grpc_call_arena_alloc(
988  call, sizeof(grpc::internal::DefaultMessageHolder<RequestType,
989  grpc::ByteBuffer>)))
991  *handler_data = allocator_state;
992  request = allocator_state->request();
993  *status = grpc::Deserialize(&buf, request);
994  buf.Release();
995  if (status->ok()) {
996  return request;
997  }
998  return nullptr;
999  }
1000 
1001  private:
1003  grpc::CallbackServerContext*, const RequestType*)>
1004  get_reactor_;
1005  grpc::Service* service_;
1006 
1007  class ServerCallbackSessionImpl
1009  public:
1010  void Finish(grpc::Status s) override {
1011  if (ctx_->IsCancelled()) {
1012  MaybeDone(
1013  reactor_.load(std::memory_order_relaxed)->InternalInlineable());
1014  return;
1015  }
1016  // A callback that only contains a call to MaybeDone can be run as an
1017  // inline callback regardless of whether or not OnDone is inlineable
1018  // because if the actual OnDone callback needs to be scheduled, MaybeDone
1019  // is responsible for dispatching to an EventEngine thread if needed.
1020  // Thus, when setting up the finish_tag_, we can set its own callback to
1021  // inlineable.
1022  finish_tag_.Set(
1023  call_.call(),
1024  [this](bool) {
1025  this->MaybeDone(
1026  reactor_.load(std::memory_order_relaxed)->InternalInlineable());
1027  },
1028  &finish_ops_, /*can_inline=*/true);
1029  finish_ops_.set_core_cq_tag(&finish_tag_);
1030 
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());
1037  }
1038  ctx_->MarkInitialMetadataSent();
1039  }
1040  finish_ops_.ServerSendStatus(&ctx_->trailing_metadata_, s);
1041  finish_ops_.set_core_cq_tag(&finish_tag_);
1042  finish_ops_.FillOps(&call_);
1043  }
1044 
1045  void SendInitialMetadata() override {
1046  ABSL_CHECK(!ctx_->sent_initial_metadata_);
1047  this->Ref();
1048  // The callback for this function should not be marked inline because it
1049  // is directly invoking a user-controlled reaction
1050  // (OnSendInitialMetadataDone). Thus it must be dispatched to an
1051  // EventEngine thread. However, any OnDone needed after that can be
1052  // inlined because it is already running on an EventEngine thread.
1053  meta_tag_.Set(
1054  call_.call(),
1055  [this](bool ok) {
1056  grpc::experimental::ServerSessionReactor* reactor =
1057  reactor_.load(std::memory_order_relaxed);
1058  reactor->OnSendInitialMetadataDone(ok);
1059  this->MaybeDone(/*inlineable_ondone=*/true);
1060  },
1061  &meta_ops_, /*can_inline=*/false);
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());
1066  }
1067  ctx_->MarkInitialMetadataSent();
1068  meta_ops_.set_core_cq_tag(&meta_tag_);
1069  meta_ops_.FillOps(&call_);
1070  // We bind the inner server only when sending initial metadata because
1071  // this signals that the session handler has completed context
1072  // establishment and virtual RPCs can be started.
1073  BindInnerServer(inner_server_);
1074  }
1075 
1076  void BindInnerServer(grpc::Server* inner_server) override {
1078  call_.call(), inner_server, &transport_, &endpoint_);
1079  }
1080 
1081  void InitiateGracefulShutdown(
1082  absl::AnyInvocable<void(absl::Status)> on_shutdown) override {
1084  transport_, endpoint_, std::move(on_shutdown));
1085  }
1086 
1087  private:
1088  friend class CallbackSessionHandler<RequestType>;
1089 
1090  ServerCallbackSessionImpl(
1092  MessageHolder<RequestType, grpc::ByteBuffer>* allocator_state,
1093  std::function<void()> call_requester, grpc::Server* inner_server)
1094  : ctx_(ctx),
1095  call_(*call),
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);
1101  }
1102 
1103  grpc_call* call() override { return call_.call(); }
1104 
1109  void SetupReactor(grpc::experimental::ServerSessionReactor* reactor) {
1110  reactor_.store(reactor, std::memory_order_relaxed);
1111  this->BindReactor(reactor);
1112  this->MaybeCallOnCancel(reactor);
1113  this->MaybeDone(reactor->InternalInlineable());
1114  }
1115 
1116  const RequestType* request() { return allocator_state_->request(); }
1117 
1118  void CallOnDone() override {
1119  reactor_.load(std::memory_order_relaxed)->OnDone();
1120  grpc_call* call = call_.call();
1121  auto call_requester = std::move(call_requester_);
1122  allocator_state_->Release();
1123  if (ctx_->context_allocator() != nullptr) {
1124  ctx_->context_allocator()->Release(ctx_);
1125  }
1126  this->~ServerCallbackSessionImpl(); // explicitly call destructor
1127  grpc_call_unref(call);
1128  call_requester();
1129  }
1130 
1131  grpc::internal::ServerReactor* reactor() override {
1132  return reactor_.load(std::memory_order_relaxed);
1133  }
1134 
1136  meta_ops_;
1140  finish_ops_;
1142 
1143  grpc::CallbackServerContext* const ctx_;
1144  grpc::internal::Call call_;
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_;
1149  grpc::Server* inner_server_;
1150  // The memory ordering of reactor_ follows ServerCallbackUnaryImpl.
1151  std::atomic<grpc::experimental::ServerSessionReactor*> reactor_;
1152  // callbacks_outstanding_ follows a refcount pattern
1153  std::atomic<intptr_t> callbacks_outstanding_{
1154  3}; // reserve for start, Finish, and CompletionOp
1155  };
1156 };
1157 
1158 } // namespace internal
1159 } // namespace experimental
1160 } // namespace grpc
1161 
1162 #endif // GRPCPP_IMPL_SERVER_CALLBACK_HANDLERS_H
grpc::ServerReadReactor
ServerReadReactor is the interface for a client-streaming RPC.
Definition: server_callback.h:212
grpc::internal::CallbackWithSuccessTag
CallbackWithSuccessTag can be reused multiple times, and will be used in this fashion for streaming o...
Definition: callback_common.h:153
grpc::ServerCallbackReaderWriter::SendInitialMetadata
virtual void SendInitialMetadata()=0
message_allocator.h
grpc_call_arena_alloc
GRPCAPI void * grpc_call_arena_alloc(grpc_call *call, size_t size)
Allocate memory in the grpc_call arena: this memory is automatically discarded at call completion.
grpc::internal::CallOpServerSendStatus
Definition: call_op_set.h:656
grpc::Server
Represents a gRPC server.
Definition: server.h:58
grpc
An Alarm posts the user-provided tag to its associated completion queue or invokes the user-provided ...
Definition: alarm.h:33
grpc::internal::CallOpSet< grpc::internal::CallOpSendInitialMetadata >
grpc::CallbackServerContext
Definition: server_context.h:741
grpc::Deserialize
auto Deserialize(BufferPtr buffer, Message *msg)
Definition: serialization_traits.h:120
grpc::internal::CallOpSendMessage
Definition: call_op_set.h:287
grpc::internal::MethodHandler::HandlerParameter
Definition: rpc_service_method.h:43
grpc::MessageHolder::response
ResponseT * response()
Definition: message_allocator.h:46
grpc::internal::Call
Straightforward wrapping of the C call object.
Definition: call.h:34
rpc_service_method.h
grpc::Service
Descriptor of an RPC service and its various RPC methods.
Definition: service_type.h:66
grpc::internal::CallbackClientStreamingHandler::RunHandler
void RunHandler(const HandlerParameter &param) final
Definition: server_callback_handlers.h:269
status.h
grpc::ServerCallbackUnary::BindReactor
void BindReactor(Reactor *reactor)
Definition: server_callback.h:236
grpc::internal::CallOpSendInitialMetadata
Definition: call_op_set.h:217
grpc::internal::ServerCallbackCall::MaybeCallOnCancel
void MaybeCallOnCancel(ServerReactor *reactor)
Definition: server_callback.h:131
grpc::internal::CallbackClientStreamingHandler::CallbackClientStreamingHandler
CallbackClientStreamingHandler(std::function< ServerReadReactor< RequestType > *(grpc::CallbackServerContext *, ResponseType *)> get_reactor)
Definition: server_callback_handlers.h:264
grpc::Status::ok
bool ok() const
Is the status OK?
Definition: status.h:124
grpc_call_ref
GRPCAPI void grpc_call_ref(grpc_call *call)
Ref a call.
grpc::ServerCallbackReaderWriter
Definition: server_callback.h:292
grpc::internal::CallbackUnaryHandler::CallbackUnaryHandler
CallbackUnaryHandler(std::function< ServerUnaryReactor *(grpc::CallbackServerContext *, const RequestType *, ResponseType *)> get_reactor)
Definition: server_callback_handlers.h:40
grpc::Status
Did it work? If it didn't, why?
Definition: status.h:34
grpc::experimental::ServerSessionReactor
Definition: server_callback.h:820
grpc_call_unref
GRPCAPI void grpc_call_unref(grpc_call *call)
Unref a call.
grpc::internal::CallbackServerStreamingHandler::Deserialize
void * Deserialize(grpc_call *call, grpc_byte_buffer *req, grpc::Status *status, void **) final
Definition: server_callback_handlers.h:508
grpc::experimental::internal::CallbackSessionHandler::Deserialize
void * Deserialize(grpc_call *call, grpc_byte_buffer *req, grpc::Status *status, void **handler_data) final
Definition: server_callback_handlers.h:981
grpc::MessageAllocator< RequestType, ResponseType >
grpc.h
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_byte_buffer
Definition: grpc_types.h:41
grpc::ByteBuffer
A sequence of bytes.
Definition: byte_buffer.h:67
grpc::internal::CallbackUnaryHandler
Definition: server_callback_handlers.h:38
grpc::ServerCallbackWriter::SendInitialMetadata
virtual void SendInitialMetadata()=0
grpc::internal::ReturnPreexistingErrors
bool ReturnPreexistingErrors()
grpc::internal::DefaultMessageHolder
Definition: server_callback.h:191
grpc::internal::CallbackServerStreamingHandler
Definition: server_callback_handlers.h:460
grpc::experimental::ServerCallbackSession
Definition: server_callback.h:242
grpc::internal::CallbackServerStreamingHandler::CallbackServerStreamingHandler
CallbackServerStreamingHandler(std::function< ServerWriteReactor< ResponseType > *(grpc::CallbackServerContext *, const RequestType *)> get_reactor)
Definition: server_callback_handlers.h:462
grpc::internal::CallbackBidiHandler::RunHandler
void RunHandler(const HandlerParameter &param) final
Definition: server_callback_handlers.h:702
grpc::experimental::internal::InitiateSessionGracefulShutdown
void InitiateSessionGracefulShutdown(grpc_core::Transport *transport, grpc_endpoint *endpoint, absl::AnyInvocable< void(absl::Status)> on_shutdown)
grpc::internal::CallbackUnaryHandler::RunHandler
void RunHandler(const HandlerParameter &param) final
Definition: server_callback_handlers.h:51
grpc::ServerCallbackReader< RequestType >::BindReactor
void BindReactor(ServerReadReactor< RequestType > *reactor)
Definition: server_callback.h:269
grpc::MessageAllocator::AllocateMessages
virtual MessageHolder< RequestT, ResponseT > * AllocateMessages()=0
grpc::experimental::internal::CallbackSessionHandler::CallbackSessionHandler
CallbackSessionHandler(std::function< grpc::experimental::ServerSessionReactor *(grpc::CallbackServerContext *, const RequestType *)> get_reactor, grpc::Service *service=nullptr)
Definition: server_callback_handlers.h:931
grpc::internal::FinishOnlyReactor
Definition: server_context.h:90
grpc::internal::MethodHandler
Base class for running an RPC handler.
Definition: rpc_service_method.h:40
server_callback.h
grpc::experimental::internal::CallbackSessionHandler
Definition: server_callback_handlers.h:929
grpc::MessageHolder::request
RequestT * request()
Definition: message_allocator.h:45
grpc::WriteOptions
Per-message write options.
Definition: call_op_set.h:79
grpc::ServerWriteReactor
ServerWriteReactor is the interface for a server-streaming RPC.
Definition: server_callback.h:214
grpc::internal::UnimplementedUnaryReactor
FinishOnlyReactor< ServerUnaryReactor > UnimplementedUnaryReactor
Definition: server_callback.h:931
grpc::internal::CallbackBidiHandler
Definition: server_callback_handlers.h:695
grpc::UNIMPLEMENTED
@ UNIMPLEMENTED
Operation is not implemented or not supported/enabled in this service.
Definition: status_code_enum.h:117
grpc::ServerCallbackWriter
Definition: server_callback.h:275
call.h
grpc::internal::Call::call
grpc_call * call() const
Definition: call.h:55
server_context.h
std
Definition: async_unary_call.h:410
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::WriteOptions::set_buffer_hint
WriteOptions & set_buffer_hint()
Sets flag indicating that the write may be buffered and need not go out on the wire immediately.
Definition: call_op_set.h:117
grpc::internal::CallbackUnaryHandler::SetMessageAllocator
void SetMessageAllocator(MessageAllocator< RequestType, ResponseType > *allocator)
Definition: server_callback_handlers.h:46
grpc::ServerCallbackUnary::SendInitialMetadata
virtual void SendInitialMetadata()=0
service_type.h
grpc::internal::ServerReactor
Definition: server_callback.h:73
grpc::experimental::internal::CallbackSessionHandler::RunHandler
void RunHandler(const HandlerParameter &param) final
Definition: server_callback_handlers.h:940
grpc::internal::CallbackClientStreamingHandler
Definition: server_callback_handlers.h:262
grpc::internal::ServerCallbackCall::Ref
void Ref()
Increases the reference count.
Definition: server_callback.h:149
grpc::internal::CallbackBidiHandler::CallbackBidiHandler
CallbackBidiHandler(std::function< ServerBidiReactor< RequestType, ResponseType > *(grpc::CallbackServerContext *)> get_reactor)
Definition: server_callback_handlers.h:697
grpc::ServerCallbackReader
Definition: server_callback.h:261
server.h
grpc::WriteOptions::is_last_message
bool is_last_message() const
Get value for the flag indicating that this is the last message, and should be coalesced with trailin...
Definition: call_op_set.h:172
grpc::MessageHolder< RequestType, ResponseType >
grpc::ByteBuffer::Release
void Release()
Forget underlying byte buffer without destroying Use this only for un-owned byte buffers.
Definition: byte_buffer.h:151
grpc::experimental::internal::BindSessionToInnerServer
void BindSessionToInnerServer(grpc_call *call, grpc::Server *inner_server, grpc_core::Transport **out_transport, grpc_endpoint **out_endpoint)
grpc::internal::CatchingReactorGetter
Reactor * CatchingReactorGetter(Func &&func, Args &&... args)
Definition: callback_common.h:57
grpc::ServerCallbackUnary
Definition: server_callback.h:226
grpc::ServerUnaryReactor
Definition: server_callback.h:751
grpc::protobuf::util::Status
::absl::Status Status
Definition: config_protobuf.h:107
grpc::internal::CallbackServerStreamingHandler::RunHandler
void RunHandler(const HandlerParameter &param) final
Definition: server_callback_handlers.h:467
grpc::ServerCallbackReader::SendInitialMetadata
virtual void SendInitialMetadata()=0
grpc::internal::ServerCallbackCall::MaybeDone
void MaybeDone()
Definition: server_callback.h:117
grpc::internal::CallbackUnaryHandler::Deserialize
void * Deserialize(grpc_call *call, grpc_byte_buffer *req, grpc::Status *status, void **handler_data) final
Definition: server_callback_handlers.h:90