grpc_client_low.h 45 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415
  1. #pragma once
  2. #include "grpc_common.h"
  3. #include <library/cpp/deprecated/atomic/atomic.h>
  4. #include <util/thread/factory.h>
  5. #include <util/string/builder.h>
  6. #include <grpc++/grpc++.h>
  7. #include <grpc++/support/async_stream.h>
  8. #include <grpc++/support/async_unary_call.h>
  9. #include <deque>
  10. #include <typeindex>
  11. #include <typeinfo>
  12. #include <variant>
  13. #include <vector>
  14. #include <unordered_map>
  15. #include <unordered_set>
  16. #include <mutex>
  17. #include <shared_mutex>
  18. /*
  19. * This file contains low level logic for grpc
  20. * This file should not be used in high level code without special reason
  21. */
  22. namespace NGrpc {
  23. const size_t DEFAULT_NUM_THREADS = 2;
  24. ////////////////////////////////////////////////////////////////////////////////
  25. void EnableGRpcTracing();
  26. ////////////////////////////////////////////////////////////////////////////////
  27. struct TTcpKeepAliveSettings {
  28. bool Enabled;
  29. size_t Idle;
  30. size_t Count;
  31. size_t Interval;
  32. };
  33. ////////////////////////////////////////////////////////////////////////////////
  34. // Common interface used to execute action from grpc cq routine
  35. class IQueueClientEvent {
  36. public:
  37. virtual ~IQueueClientEvent() = default;
  38. //! Execute an action defined by implementation
  39. virtual bool Execute(bool ok) = 0;
  40. //! Finish and destroy event
  41. virtual void Destroy() = 0;
  42. };
  43. // Implementation of IQueueClientEvent that reduces allocations
  44. template<class TSelf>
  45. class TQueueClientFixedEvent : private IQueueClientEvent {
  46. using TCallback = void (TSelf::*)(bool);
  47. public:
  48. TQueueClientFixedEvent(TSelf* self, TCallback callback)
  49. : Self(self)
  50. , Callback(callback)
  51. { }
  52. IQueueClientEvent* Prepare() {
  53. Self->Ref();
  54. return this;
  55. }
  56. private:
  57. bool Execute(bool ok) override {
  58. ((*Self).*Callback)(ok);
  59. return false;
  60. }
  61. void Destroy() override {
  62. Self->UnRef();
  63. }
  64. private:
  65. TSelf* const Self;
  66. TCallback const Callback;
  67. };
  68. class IQueueClientContext;
  69. using IQueueClientContextPtr = std::shared_ptr<IQueueClientContext>;
  70. // Provider of IQueueClientContext instances
  71. class IQueueClientContextProvider {
  72. public:
  73. virtual ~IQueueClientContextProvider() = default;
  74. virtual IQueueClientContextPtr CreateContext() = 0;
  75. };
  76. // Activity context for a low-level client
  77. class IQueueClientContext : public IQueueClientContextProvider {
  78. public:
  79. virtual ~IQueueClientContext() = default;
  80. //! Returns CompletionQueue associated with the client
  81. virtual grpc::CompletionQueue* CompletionQueue() = 0;
  82. //! Returns true if context has been cancelled
  83. virtual bool IsCancelled() const = 0;
  84. //! Tries to cancel context, calling all registered callbacks
  85. virtual bool Cancel() = 0;
  86. //! Subscribes callback to cancellation
  87. //
  88. // Note there's no way to unsubscribe, if subscription is temporary
  89. // make sure you create a new context with CreateContext and release
  90. // it as soon as it's no longer needed.
  91. virtual void SubscribeCancel(std::function<void()> callback) = 0;
  92. //! Subscribes callback to cancellation
  93. //
  94. // This alias is for compatibility with older code.
  95. void SubscribeStop(std::function<void()> callback) {
  96. SubscribeCancel(std::move(callback));
  97. }
  98. };
  99. // Represents grpc status and error message string
  100. struct TGrpcStatus {
  101. TString Msg;
  102. TString Details;
  103. int GRpcStatusCode;
  104. bool InternalError;
  105. TGrpcStatus()
  106. : GRpcStatusCode(grpc::StatusCode::OK)
  107. , InternalError(false)
  108. { }
  109. TGrpcStatus(TString msg, int statusCode, bool internalError)
  110. : Msg(std::move(msg))
  111. , GRpcStatusCode(statusCode)
  112. , InternalError(internalError)
  113. { }
  114. TGrpcStatus(grpc::StatusCode status, TString msg, TString details = {})
  115. : Msg(std::move(msg))
  116. , Details(std::move(details))
  117. , GRpcStatusCode(status)
  118. , InternalError(false)
  119. { }
  120. TGrpcStatus(const grpc::Status& status)
  121. : TGrpcStatus(status.error_code(), TString(status.error_message()), TString(status.error_details()))
  122. { }
  123. TGrpcStatus& operator=(const grpc::Status& status) {
  124. Msg = TString(status.error_message());
  125. Details = TString(status.error_details());
  126. GRpcStatusCode = status.error_code();
  127. InternalError = false;
  128. return *this;
  129. }
  130. static TGrpcStatus Internal(TString msg) {
  131. return { std::move(msg), -1, true };
  132. }
  133. bool Ok() const {
  134. return !InternalError && GRpcStatusCode == grpc::StatusCode::OK;
  135. }
  136. TStringBuilder ToDebugString() const {
  137. TStringBuilder ret;
  138. ret << "gRpcStatusCode: " << GRpcStatusCode;
  139. if(!Ok())
  140. ret << ", Msg: " << Msg << ", Details: " << Details << ", InternalError: " << InternalError;
  141. return ret;
  142. }
  143. };
  144. bool inline IsGRpcStatusGood(const TGrpcStatus& status) {
  145. return status.Ok();
  146. }
  147. // Response callback type - this callback will be called when request is finished
  148. // (or after getting each chunk in case of streaming mode)
  149. template<typename TResponse>
  150. using TResponseCallback = std::function<void (TGrpcStatus&&, TResponse&&)>;
  151. template<typename TResponse>
  152. using TAdvancedResponseCallback = std::function<void (const grpc::ClientContext&, TGrpcStatus&&, TResponse&&)>;
  153. // Call associated metadata
  154. struct TCallMeta {
  155. std::shared_ptr<grpc::CallCredentials> CallCredentials;
  156. std::vector<std::pair<TString, TString>> Aux;
  157. std::variant<TDuration, TInstant> Timeout; // timeout as duration from now or time point in future
  158. };
  159. class TGRpcRequestProcessorCommon {
  160. protected:
  161. void ApplyMeta(const TCallMeta& meta) {
  162. for (const auto& rec : meta.Aux) {
  163. Context.AddMetadata(rec.first, rec.second);
  164. }
  165. if (meta.CallCredentials) {
  166. Context.set_credentials(meta.CallCredentials);
  167. }
  168. if (const TDuration* timeout = std::get_if<TDuration>(&meta.Timeout)) {
  169. if (*timeout) {
  170. auto deadline = gpr_time_add(
  171. gpr_now(GPR_CLOCK_MONOTONIC),
  172. gpr_time_from_micros(timeout->MicroSeconds(), GPR_TIMESPAN));
  173. Context.set_deadline(deadline);
  174. }
  175. } else if (const TInstant* deadline = std::get_if<TInstant>(&meta.Timeout)) {
  176. if (*deadline) {
  177. Context.set_deadline(gpr_time_from_micros(deadline->MicroSeconds(), GPR_CLOCK_MONOTONIC));
  178. }
  179. }
  180. }
  181. void GetInitialMetadata(std::unordered_multimap<TString, TString>* metadata) {
  182. for (const auto& [key, value] : Context.GetServerInitialMetadata()) {
  183. metadata->emplace(
  184. TString(key.begin(), key.end()),
  185. TString(value.begin(), value.end())
  186. );
  187. }
  188. }
  189. grpc::Status Status;
  190. grpc::ClientContext Context;
  191. std::shared_ptr<IQueueClientContext> LocalContext;
  192. };
  193. template<typename TStub, typename TRequest, typename TResponse>
  194. class TSimpleRequestProcessor
  195. : public TThrRefBase
  196. , public IQueueClientEvent
  197. , public TGRpcRequestProcessorCommon {
  198. using TAsyncReaderPtr = std::unique_ptr<grpc::ClientAsyncResponseReader<TResponse>>;
  199. template<typename> friend class TServiceConnection;
  200. public:
  201. using TPtr = TIntrusivePtr<TSimpleRequestProcessor>;
  202. using TAsyncRequest = TAsyncReaderPtr (TStub::*)(grpc::ClientContext*, const TRequest&, grpc::CompletionQueue*);
  203. explicit TSimpleRequestProcessor(TResponseCallback<TResponse>&& callback)
  204. : Callback_(std::move(callback))
  205. { }
  206. ~TSimpleRequestProcessor() {
  207. if (!Replied_ && Callback_) {
  208. Callback_(TGrpcStatus::Internal("request left unhandled"), std::move(Reply_));
  209. Callback_ = nullptr; // free resources as early as possible
  210. }
  211. }
  212. bool Execute(bool ok) override {
  213. {
  214. std::unique_lock<std::mutex> guard(Mutex_);
  215. LocalContext.reset();
  216. }
  217. TGrpcStatus status;
  218. if (ok) {
  219. status = Status;
  220. } else {
  221. status = TGrpcStatus::Internal("Unexpected error");
  222. }
  223. Replied_ = true;
  224. Callback_(std::move(status), std::move(Reply_));
  225. Callback_ = nullptr; // free resources as early as possible
  226. return false;
  227. }
  228. void Destroy() override {
  229. UnRef();
  230. }
  231. private:
  232. IQueueClientEvent* FinishedEvent() {
  233. Ref();
  234. return this;
  235. }
  236. void Start(TStub& stub, TAsyncRequest asyncRequest, const TRequest& request, IQueueClientContextProvider* provider) {
  237. auto context = provider->CreateContext();
  238. if (!context) {
  239. Replied_ = true;
  240. Callback_(TGrpcStatus(grpc::StatusCode::CANCELLED, "Client is shutting down"), std::move(Reply_));
  241. Callback_ = nullptr;
  242. return;
  243. }
  244. {
  245. std::unique_lock<std::mutex> guard(Mutex_);
  246. LocalContext = context;
  247. Reader_ = (stub.*asyncRequest)(&Context, request, context->CompletionQueue());
  248. Reader_->Finish(&Reply_, &Status, FinishedEvent());
  249. }
  250. context->SubscribeStop([self = TPtr(this)] {
  251. self->Stop();
  252. });
  253. }
  254. void Stop() {
  255. Context.TryCancel();
  256. }
  257. TResponseCallback<TResponse> Callback_;
  258. TResponse Reply_;
  259. std::mutex Mutex_;
  260. TAsyncReaderPtr Reader_;
  261. bool Replied_ = false;
  262. };
  263. template<typename TStub, typename TRequest, typename TResponse>
  264. class TAdvancedRequestProcessor
  265. : public TThrRefBase
  266. , public IQueueClientEvent
  267. , public TGRpcRequestProcessorCommon {
  268. using TAsyncReaderPtr = std::unique_ptr<grpc::ClientAsyncResponseReader<TResponse>>;
  269. template<typename> friend class TServiceConnection;
  270. public:
  271. using TPtr = TIntrusivePtr<TAdvancedRequestProcessor>;
  272. using TAsyncRequest = TAsyncReaderPtr (TStub::*)(grpc::ClientContext*, const TRequest&, grpc::CompletionQueue*);
  273. explicit TAdvancedRequestProcessor(TAdvancedResponseCallback<TResponse>&& callback)
  274. : Callback_(std::move(callback))
  275. { }
  276. ~TAdvancedRequestProcessor() {
  277. if (!Replied_ && Callback_) {
  278. Callback_(Context, TGrpcStatus::Internal("request left unhandled"), std::move(Reply_));
  279. Callback_ = nullptr; // free resources as early as possible
  280. }
  281. }
  282. bool Execute(bool ok) override {
  283. {
  284. std::unique_lock<std::mutex> guard(Mutex_);
  285. LocalContext.reset();
  286. }
  287. TGrpcStatus status;
  288. if (ok) {
  289. status = Status;
  290. } else {
  291. status = TGrpcStatus::Internal("Unexpected error");
  292. }
  293. Replied_ = true;
  294. Callback_(Context, std::move(status), std::move(Reply_));
  295. Callback_ = nullptr; // free resources as early as possible
  296. return false;
  297. }
  298. void Destroy() override {
  299. UnRef();
  300. }
  301. private:
  302. IQueueClientEvent* FinishedEvent() {
  303. Ref();
  304. return this;
  305. }
  306. void Start(TStub& stub, TAsyncRequest asyncRequest, const TRequest& request, IQueueClientContextProvider* provider) {
  307. auto context = provider->CreateContext();
  308. if (!context) {
  309. Replied_ = true;
  310. Callback_(Context, TGrpcStatus(grpc::StatusCode::CANCELLED, "Client is shutting down"), std::move(Reply_));
  311. Callback_ = nullptr;
  312. return;
  313. }
  314. {
  315. std::unique_lock<std::mutex> guard(Mutex_);
  316. LocalContext = context;
  317. Reader_ = (stub.*asyncRequest)(&Context, request, context->CompletionQueue());
  318. Reader_->Finish(&Reply_, &Status, FinishedEvent());
  319. }
  320. context->SubscribeStop([self = TPtr(this)] {
  321. self->Stop();
  322. });
  323. }
  324. void Stop() {
  325. Context.TryCancel();
  326. }
  327. TAdvancedResponseCallback<TResponse> Callback_;
  328. TResponse Reply_;
  329. std::mutex Mutex_;
  330. TAsyncReaderPtr Reader_;
  331. bool Replied_ = false;
  332. };
  333. class IStreamRequestCtrl : public TThrRefBase {
  334. public:
  335. using TPtr = TIntrusivePtr<IStreamRequestCtrl>;
  336. /**
  337. * Asynchronously cancel the request
  338. */
  339. virtual void Cancel() = 0;
  340. };
  341. template<class TResponse>
  342. class IStreamRequestReadProcessor : public IStreamRequestCtrl {
  343. public:
  344. using TPtr = TIntrusivePtr<IStreamRequestReadProcessor>;
  345. using TReadCallback = std::function<void(TGrpcStatus&&)>;
  346. /**
  347. * Scheduled initial server metadata read from the stream
  348. */
  349. virtual void ReadInitialMetadata(std::unordered_multimap<TString, TString>* metadata, TReadCallback callback) = 0;
  350. /**
  351. * Scheduled response read from the stream
  352. * Callback will be called with the status if it failed
  353. * Only one Read or Finish call may be active at a time
  354. */
  355. virtual void Read(TResponse* response, TReadCallback callback) = 0;
  356. /**
  357. * Stop reading and gracefully finish the stream
  358. * Only one Read or Finish call may be active at a time
  359. */
  360. virtual void Finish(TReadCallback callback) = 0;
  361. /**
  362. * Additional callback to be called when stream has finished
  363. */
  364. virtual void AddFinishedCallback(TReadCallback callback) = 0;
  365. };
  366. template<class TRequest, class TResponse>
  367. class IStreamRequestReadWriteProcessor : public IStreamRequestReadProcessor<TResponse> {
  368. public:
  369. using TPtr = TIntrusivePtr<IStreamRequestReadWriteProcessor>;
  370. using TWriteCallback = std::function<void(TGrpcStatus&&)>;
  371. /**
  372. * Scheduled request write to the stream
  373. */
  374. virtual void Write(TRequest&& request, TWriteCallback callback = { }) = 0;
  375. };
  376. class TGRpcKeepAliveSocketMutator;
  377. // Class to hold stubs allocated on channel.
  378. // It is poor documented part of grpc. See KIKIMR-6109 and comment to this commit
  379. // Stub holds shared_ptr<ChannelInterface>, so we can destroy this holder even if
  380. // request processor using stub
  381. class TStubsHolder : public TNonCopyable {
  382. using TypeInfoRef = std::reference_wrapper<const std::type_info>;
  383. struct THasher {
  384. std::size_t operator()(TypeInfoRef code) const {
  385. return code.get().hash_code();
  386. }
  387. };
  388. struct TEqualTo {
  389. bool operator()(TypeInfoRef lhs, TypeInfoRef rhs) const {
  390. return lhs.get() == rhs.get();
  391. }
  392. };
  393. public:
  394. TStubsHolder(std::shared_ptr<grpc::ChannelInterface> channel)
  395. : ChannelInterface_(channel)
  396. {}
  397. // Returns true if channel can't be used to perform request now
  398. bool IsChannelBroken() const {
  399. auto state = ChannelInterface_->GetState(false);
  400. return state == GRPC_CHANNEL_SHUTDOWN ||
  401. state == GRPC_CHANNEL_TRANSIENT_FAILURE;
  402. }
  403. template<typename TStub>
  404. std::shared_ptr<TStub> GetOrCreateStub() {
  405. const auto& stubId = typeid(TStub);
  406. {
  407. std::shared_lock readGuard(RWMutex_);
  408. const auto it = Stubs_.find(stubId);
  409. if (it != Stubs_.end()) {
  410. return std::static_pointer_cast<TStub>(it->second);
  411. }
  412. }
  413. {
  414. std::unique_lock writeGuard(RWMutex_);
  415. auto it = Stubs_.emplace(stubId, nullptr);
  416. if (!it.second) {
  417. return std::static_pointer_cast<TStub>(it.first->second);
  418. } else {
  419. it.first->second = std::make_shared<TStub>(ChannelInterface_);
  420. return std::static_pointer_cast<TStub>(it.first->second);
  421. }
  422. }
  423. }
  424. const TInstant& GetLastUseTime() const {
  425. return LastUsed_;
  426. }
  427. void SetLastUseTime(const TInstant& time) {
  428. LastUsed_ = time;
  429. }
  430. private:
  431. TInstant LastUsed_ = Now();
  432. std::shared_mutex RWMutex_;
  433. std::unordered_map<TypeInfoRef, std::shared_ptr<void>, THasher, TEqualTo> Stubs_;
  434. std::shared_ptr<grpc::ChannelInterface> ChannelInterface_;
  435. };
  436. class TChannelPool {
  437. public:
  438. TChannelPool(const TTcpKeepAliveSettings& tcpKeepAliveSettings, const TDuration& expireTime = TDuration::Minutes(6));
  439. //Allows to CreateStub from TStubsHolder under lock
  440. //The callback will be called just during GetStubsHolderLocked call
  441. void GetStubsHolderLocked(const TString& channelId, const TGRpcClientConfig& config, std::function<void(TStubsHolder&)> cb);
  442. void DeleteChannel(const TString& channelId);
  443. void DeleteExpiredStubsHolders();
  444. private:
  445. std::shared_mutex RWMutex_;
  446. std::unordered_map<TString, TStubsHolder> Pool_;
  447. std::multimap<TInstant, TString> LastUsedQueue_;
  448. TTcpKeepAliveSettings TcpKeepAliveSettings_;
  449. TDuration ExpireTime_;
  450. TDuration UpdateReUseTime_;
  451. void EraseFromQueueByTime(const TInstant& lastUseTime, const TString& channelId);
  452. };
  453. template<class TResponse>
  454. using TStreamReaderCallback = std::function<void(TGrpcStatus&&, typename IStreamRequestReadProcessor<TResponse>::TPtr)>;
  455. template<typename TStub, typename TRequest, typename TResponse>
  456. class TStreamRequestReadProcessor
  457. : public IStreamRequestReadProcessor<TResponse>
  458. , public TGRpcRequestProcessorCommon {
  459. template<typename> friend class TServiceConnection;
  460. public:
  461. using TSelf = TStreamRequestReadProcessor;
  462. using TAsyncReaderPtr = std::unique_ptr<grpc::ClientAsyncReader<TResponse>>;
  463. using TAsyncRequest = TAsyncReaderPtr (TStub::*)(grpc::ClientContext*, const TRequest&, grpc::CompletionQueue*, void*);
  464. using TReaderCallback = TStreamReaderCallback<TResponse>;
  465. using TPtr = TIntrusivePtr<TSelf>;
  466. using TBase = IStreamRequestReadProcessor<TResponse>;
  467. using TReadCallback = typename TBase::TReadCallback;
  468. explicit TStreamRequestReadProcessor(TReaderCallback&& callback)
  469. : Callback(std::move(callback))
  470. {
  471. Y_VERIFY(Callback, "Missing connected callback");
  472. }
  473. void Cancel() override {
  474. Context.TryCancel();
  475. {
  476. std::unique_lock<std::mutex> guard(Mutex);
  477. Cancelled = true;
  478. if (Started && !ReadFinished) {
  479. if (!ReadActive) {
  480. ReadFinished = true;
  481. }
  482. if (ReadFinished) {
  483. Stream->Finish(&Status, OnFinishedTag.Prepare());
  484. }
  485. }
  486. }
  487. }
  488. void ReadInitialMetadata(std::unordered_multimap<TString, TString>* metadata, TReadCallback callback) override {
  489. TGrpcStatus status;
  490. {
  491. std::unique_lock<std::mutex> guard(Mutex);
  492. Y_VERIFY(!ReadActive, "Multiple Read/Finish calls detected");
  493. if (!Finished && !HasInitialMetadata) {
  494. ReadActive = true;
  495. ReadCallback = std::move(callback);
  496. InitialMetadata = metadata;
  497. if (!ReadFinished) {
  498. Stream->ReadInitialMetadata(OnReadDoneTag.Prepare());
  499. }
  500. return;
  501. }
  502. if (!HasInitialMetadata) {
  503. if (FinishedOk) {
  504. status = Status;
  505. } else {
  506. status = TGrpcStatus::Internal("Unexpected error");
  507. }
  508. } else {
  509. GetInitialMetadata(metadata);
  510. }
  511. }
  512. callback(std::move(status));
  513. }
  514. void Read(TResponse* message, TReadCallback callback) override {
  515. TGrpcStatus status;
  516. {
  517. std::unique_lock<std::mutex> guard(Mutex);
  518. Y_VERIFY(!ReadActive, "Multiple Read/Finish calls detected");
  519. if (!Finished) {
  520. ReadActive = true;
  521. ReadCallback = std::move(callback);
  522. if (!ReadFinished) {
  523. Stream->Read(message, OnReadDoneTag.Prepare());
  524. }
  525. return;
  526. }
  527. if (FinishedOk) {
  528. status = Status;
  529. } else {
  530. status = TGrpcStatus::Internal("Unexpected error");
  531. }
  532. }
  533. if (status.Ok()) {
  534. status = TGrpcStatus(grpc::StatusCode::OUT_OF_RANGE, "Read EOF");
  535. }
  536. callback(std::move(status));
  537. }
  538. void Finish(TReadCallback callback) override {
  539. TGrpcStatus status;
  540. {
  541. std::unique_lock<std::mutex> guard(Mutex);
  542. Y_VERIFY(!ReadActive, "Multiple Read/Finish calls detected");
  543. if (!Finished) {
  544. ReadActive = true;
  545. FinishCallback = std::move(callback);
  546. if (!ReadFinished) {
  547. ReadFinished = true;
  548. }
  549. Stream->Finish(&Status, OnFinishedTag.Prepare());
  550. return;
  551. }
  552. if (FinishedOk) {
  553. status = Status;
  554. } else {
  555. status = TGrpcStatus::Internal("Unexpected error");
  556. }
  557. }
  558. callback(std::move(status));
  559. }
  560. void AddFinishedCallback(TReadCallback callback) override {
  561. Y_VERIFY(callback, "Unexpected empty callback");
  562. TGrpcStatus status;
  563. {
  564. std::unique_lock<std::mutex> guard(Mutex);
  565. if (!Finished) {
  566. FinishedCallbacks.emplace_back().swap(callback);
  567. return;
  568. }
  569. if (FinishedOk) {
  570. status = Status;
  571. } else if (Cancelled) {
  572. status = TGrpcStatus(grpc::StatusCode::CANCELLED, "Stream cancelled");
  573. } else {
  574. status = TGrpcStatus::Internal("Unexpected error");
  575. }
  576. }
  577. callback(std::move(status));
  578. }
  579. private:
  580. void Start(TStub& stub, const TRequest& request, TAsyncRequest asyncRequest, IQueueClientContextProvider* provider) {
  581. auto context = provider->CreateContext();
  582. if (!context) {
  583. auto callback = std::move(Callback);
  584. TGrpcStatus status(grpc::StatusCode::CANCELLED, "Client is shutting down");
  585. callback(std::move(status), nullptr);
  586. return;
  587. }
  588. {
  589. std::unique_lock<std::mutex> guard(Mutex);
  590. LocalContext = context;
  591. Stream = (stub.*asyncRequest)(&Context, request, context->CompletionQueue(), OnStartDoneTag.Prepare());
  592. }
  593. context->SubscribeStop([self = TPtr(this)] {
  594. self->Cancel();
  595. });
  596. }
  597. void OnReadDone(bool ok) {
  598. TGrpcStatus status;
  599. TReadCallback callback;
  600. std::unordered_multimap<TString, TString>* initialMetadata = nullptr;
  601. {
  602. std::unique_lock<std::mutex> guard(Mutex);
  603. Y_VERIFY(ReadActive, "Unexpected Read done callback");
  604. Y_VERIFY(!ReadFinished, "Unexpected ReadFinished flag");
  605. if (!ok || Cancelled) {
  606. ReadFinished = true;
  607. Stream->Finish(&Status, OnFinishedTag.Prepare());
  608. if (!ok) {
  609. // Keep ReadActive=true, so callback is called
  610. // after the call is finished with an error
  611. return;
  612. }
  613. }
  614. callback = std::move(ReadCallback);
  615. ReadCallback = nullptr;
  616. ReadActive = false;
  617. initialMetadata = InitialMetadata;
  618. InitialMetadata = nullptr;
  619. HasInitialMetadata = true;
  620. }
  621. if (initialMetadata) {
  622. GetInitialMetadata(initialMetadata);
  623. }
  624. callback(std::move(status));
  625. }
  626. void OnStartDone(bool ok) {
  627. TReaderCallback callback;
  628. {
  629. std::unique_lock<std::mutex> guard(Mutex);
  630. Started = true;
  631. if (!ok || Cancelled) {
  632. ReadFinished = true;
  633. Stream->Finish(&Status, OnFinishedTag.Prepare());
  634. return;
  635. }
  636. callback = std::move(Callback);
  637. Callback = nullptr;
  638. }
  639. callback({ }, typename TBase::TPtr(this));
  640. }
  641. void OnFinished(bool ok) {
  642. TGrpcStatus status;
  643. std::vector<TReadCallback> finishedCallbacks;
  644. TReaderCallback startCallback;
  645. TReadCallback readCallback;
  646. TReadCallback finishCallback;
  647. {
  648. std::unique_lock<std::mutex> guard(Mutex);
  649. Finished = true;
  650. FinishedOk = ok;
  651. LocalContext.reset();
  652. if (ok) {
  653. status = Status;
  654. } else if (Cancelled) {
  655. status = TGrpcStatus(grpc::StatusCode::CANCELLED, "Stream cancelled");
  656. } else {
  657. status = TGrpcStatus::Internal("Unexpected error");
  658. }
  659. finishedCallbacks.swap(FinishedCallbacks);
  660. if (Callback) {
  661. Y_VERIFY(!ReadActive);
  662. startCallback = std::move(Callback);
  663. Callback = nullptr;
  664. } else if (ReadActive) {
  665. if (ReadCallback) {
  666. readCallback = std::move(ReadCallback);
  667. ReadCallback = nullptr;
  668. } else {
  669. finishCallback = std::move(FinishCallback);
  670. FinishCallback = nullptr;
  671. }
  672. ReadActive = false;
  673. }
  674. }
  675. for (auto& finishedCallback : finishedCallbacks) {
  676. auto statusCopy = status;
  677. finishedCallback(std::move(statusCopy));
  678. }
  679. if (startCallback) {
  680. if (status.Ok()) {
  681. status = TGrpcStatus(grpc::StatusCode::UNKNOWN, "Unknown stream failure");
  682. }
  683. startCallback(std::move(status), nullptr);
  684. } else if (readCallback) {
  685. if (status.Ok()) {
  686. status = TGrpcStatus(grpc::StatusCode::OUT_OF_RANGE, "Read EOF");
  687. }
  688. readCallback(std::move(status));
  689. } else if (finishCallback) {
  690. finishCallback(std::move(status));
  691. }
  692. }
  693. TReaderCallback Callback;
  694. TAsyncReaderPtr Stream;
  695. using TFixedEvent = TQueueClientFixedEvent<TSelf>;
  696. std::mutex Mutex;
  697. TFixedEvent OnReadDoneTag = { this, &TSelf::OnReadDone };
  698. TFixedEvent OnStartDoneTag = { this, &TSelf::OnStartDone };
  699. TFixedEvent OnFinishedTag = { this, &TSelf::OnFinished };
  700. TReadCallback ReadCallback;
  701. TReadCallback FinishCallback;
  702. std::vector<TReadCallback> FinishedCallbacks;
  703. std::unordered_multimap<TString, TString>* InitialMetadata = nullptr;
  704. bool Started = false;
  705. bool HasInitialMetadata = false;
  706. bool ReadActive = false;
  707. bool ReadFinished = false;
  708. bool Finished = false;
  709. bool Cancelled = false;
  710. bool FinishedOk = false;
  711. };
  712. template<class TRequest, class TResponse>
  713. using TStreamConnectedCallback = std::function<void(TGrpcStatus&&, typename IStreamRequestReadWriteProcessor<TRequest, TResponse>::TPtr)>;
  714. template<class TStub, class TRequest, class TResponse>
  715. class TStreamRequestReadWriteProcessor
  716. : public IStreamRequestReadWriteProcessor<TRequest, TResponse>
  717. , public TGRpcRequestProcessorCommon {
  718. public:
  719. using TSelf = TStreamRequestReadWriteProcessor;
  720. using TBase = IStreamRequestReadWriteProcessor<TRequest, TResponse>;
  721. using TPtr = TIntrusivePtr<TSelf>;
  722. using TConnectedCallback = TStreamConnectedCallback<TRequest, TResponse>;
  723. using TReadCallback = typename TBase::TReadCallback;
  724. using TWriteCallback = typename TBase::TWriteCallback;
  725. using TAsyncReaderWriterPtr = std::unique_ptr<grpc::ClientAsyncReaderWriter<TRequest, TResponse>>;
  726. using TAsyncRequest = TAsyncReaderWriterPtr (TStub::*)(grpc::ClientContext*, grpc::CompletionQueue*, void*);
  727. explicit TStreamRequestReadWriteProcessor(TConnectedCallback&& callback)
  728. : ConnectedCallback(std::move(callback))
  729. {
  730. Y_VERIFY(ConnectedCallback, "Missing connected callback");
  731. }
  732. void Cancel() override {
  733. Context.TryCancel();
  734. {
  735. std::unique_lock<std::mutex> guard(Mutex);
  736. Cancelled = true;
  737. if (Started && !(ReadFinished && WriteFinished)) {
  738. if (!ReadActive) {
  739. ReadFinished = true;
  740. }
  741. if (!WriteActive) {
  742. WriteFinished = true;
  743. }
  744. if (ReadFinished && WriteFinished) {
  745. Stream->Finish(&Status, OnFinishedTag.Prepare());
  746. }
  747. }
  748. }
  749. }
  750. void Write(TRequest&& request, TWriteCallback callback) override {
  751. TGrpcStatus status;
  752. {
  753. std::unique_lock<std::mutex> guard(Mutex);
  754. if (Cancelled || ReadFinished || WriteFinished) {
  755. status = TGrpcStatus(grpc::StatusCode::CANCELLED, "Write request dropped");
  756. } else if (WriteActive) {
  757. auto& item = WriteQueue.emplace_back();
  758. item.Callback.swap(callback);
  759. item.Request.Swap(&request);
  760. } else {
  761. WriteActive = true;
  762. WriteCallback.swap(callback);
  763. Stream->Write(request, OnWriteDoneTag.Prepare());
  764. }
  765. }
  766. if (!status.Ok() && callback) {
  767. callback(std::move(status));
  768. }
  769. }
  770. void ReadInitialMetadata(std::unordered_multimap<TString, TString>* metadata, TReadCallback callback) override {
  771. TGrpcStatus status;
  772. {
  773. std::unique_lock<std::mutex> guard(Mutex);
  774. Y_VERIFY(!ReadActive, "Multiple Read/Finish calls detected");
  775. if (!Finished && !HasInitialMetadata) {
  776. ReadActive = true;
  777. ReadCallback = std::move(callback);
  778. InitialMetadata = metadata;
  779. if (!ReadFinished) {
  780. Stream->ReadInitialMetadata(OnReadDoneTag.Prepare());
  781. }
  782. return;
  783. }
  784. if (!HasInitialMetadata) {
  785. if (FinishedOk) {
  786. status = Status;
  787. } else {
  788. status = TGrpcStatus::Internal("Unexpected error");
  789. }
  790. } else {
  791. GetInitialMetadata(metadata);
  792. }
  793. }
  794. callback(std::move(status));
  795. }
  796. void Read(TResponse* message, TReadCallback callback) override {
  797. TGrpcStatus status;
  798. {
  799. std::unique_lock<std::mutex> guard(Mutex);
  800. Y_VERIFY(!ReadActive, "Multiple Read/Finish calls detected");
  801. if (!Finished) {
  802. ReadActive = true;
  803. ReadCallback = std::move(callback);
  804. if (!ReadFinished) {
  805. Stream->Read(message, OnReadDoneTag.Prepare());
  806. }
  807. return;
  808. }
  809. if (FinishedOk) {
  810. status = Status;
  811. } else {
  812. status = TGrpcStatus::Internal("Unexpected error");
  813. }
  814. }
  815. if (status.Ok()) {
  816. status = TGrpcStatus(grpc::StatusCode::OUT_OF_RANGE, "Read EOF");
  817. }
  818. callback(std::move(status));
  819. }
  820. void Finish(TReadCallback callback) override {
  821. TGrpcStatus status;
  822. {
  823. std::unique_lock<std::mutex> guard(Mutex);
  824. Y_VERIFY(!ReadActive, "Multiple Read/Finish calls detected");
  825. if (!Finished) {
  826. ReadActive = true;
  827. FinishCallback = std::move(callback);
  828. if (!ReadFinished) {
  829. ReadFinished = true;
  830. if (!WriteActive) {
  831. WriteFinished = true;
  832. }
  833. if (WriteFinished) {
  834. Stream->Finish(&Status, OnFinishedTag.Prepare());
  835. }
  836. }
  837. return;
  838. }
  839. if (FinishedOk) {
  840. status = Status;
  841. } else {
  842. status = TGrpcStatus::Internal("Unexpected error");
  843. }
  844. }
  845. callback(std::move(status));
  846. }
  847. void AddFinishedCallback(TReadCallback callback) override {
  848. Y_VERIFY(callback, "Unexpected empty callback");
  849. TGrpcStatus status;
  850. {
  851. std::unique_lock<std::mutex> guard(Mutex);
  852. if (!Finished) {
  853. FinishedCallbacks.emplace_back().swap(callback);
  854. return;
  855. }
  856. if (FinishedOk) {
  857. status = Status;
  858. } else if (Cancelled) {
  859. status = TGrpcStatus(grpc::StatusCode::CANCELLED, "Stream cancelled");
  860. } else {
  861. status = TGrpcStatus::Internal("Unexpected error");
  862. }
  863. }
  864. callback(std::move(status));
  865. }
  866. private:
  867. template<typename> friend class TServiceConnection;
  868. void Start(TStub& stub, TAsyncRequest asyncRequest, IQueueClientContextProvider* provider) {
  869. auto context = provider->CreateContext();
  870. if (!context) {
  871. auto callback = std::move(ConnectedCallback);
  872. TGrpcStatus status(grpc::StatusCode::CANCELLED, "Client is shutting down");
  873. callback(std::move(status), nullptr);
  874. return;
  875. }
  876. {
  877. std::unique_lock<std::mutex> guard(Mutex);
  878. LocalContext = context;
  879. Stream = (stub.*asyncRequest)(&Context, context->CompletionQueue(), OnConnectedTag.Prepare());
  880. }
  881. context->SubscribeStop([self = TPtr(this)] {
  882. self->Cancel();
  883. });
  884. }
  885. private:
  886. void OnConnected(bool ok) {
  887. TConnectedCallback callback;
  888. {
  889. std::unique_lock<std::mutex> guard(Mutex);
  890. Started = true;
  891. if (!ok || Cancelled) {
  892. ReadFinished = true;
  893. WriteFinished = true;
  894. Stream->Finish(&Status, OnFinishedTag.Prepare());
  895. return;
  896. }
  897. callback = std::move(ConnectedCallback);
  898. ConnectedCallback = nullptr;
  899. }
  900. callback({ }, typename TBase::TPtr(this));
  901. }
  902. void OnReadDone(bool ok) {
  903. TGrpcStatus status;
  904. TReadCallback callback;
  905. std::unordered_multimap<TString, TString>* initialMetadata = nullptr;
  906. {
  907. std::unique_lock<std::mutex> guard(Mutex);
  908. Y_VERIFY(ReadActive, "Unexpected Read done callback");
  909. Y_VERIFY(!ReadFinished, "Unexpected ReadFinished flag");
  910. if (!ok || Cancelled || WriteFinished) {
  911. ReadFinished = true;
  912. if (!WriteActive) {
  913. WriteFinished = true;
  914. }
  915. if (WriteFinished) {
  916. Stream->Finish(&Status, OnFinishedTag.Prepare());
  917. }
  918. if (!ok) {
  919. // Keep ReadActive=true, so callback is called
  920. // after the call is finished with an error
  921. return;
  922. }
  923. }
  924. callback = std::move(ReadCallback);
  925. ReadCallback = nullptr;
  926. ReadActive = false;
  927. initialMetadata = InitialMetadata;
  928. InitialMetadata = nullptr;
  929. HasInitialMetadata = true;
  930. }
  931. if (initialMetadata) {
  932. GetInitialMetadata(initialMetadata);
  933. }
  934. callback(std::move(status));
  935. }
  936. void OnWriteDone(bool ok) {
  937. TWriteCallback okCallback;
  938. {
  939. std::unique_lock<std::mutex> guard(Mutex);
  940. Y_VERIFY(WriteActive, "Unexpected Write done callback");
  941. Y_VERIFY(!WriteFinished, "Unexpected WriteFinished flag");
  942. if (ok) {
  943. okCallback.swap(WriteCallback);
  944. } else if (WriteCallback) {
  945. // Put callback back on the queue until OnFinished
  946. auto& item = WriteQueue.emplace_front();
  947. item.Callback.swap(WriteCallback);
  948. }
  949. if (!ok || Cancelled) {
  950. WriteActive = false;
  951. WriteFinished = true;
  952. if (!ReadActive) {
  953. ReadFinished = true;
  954. }
  955. if (ReadFinished) {
  956. Stream->Finish(&Status, OnFinishedTag.Prepare());
  957. }
  958. } else if (!WriteQueue.empty()) {
  959. WriteCallback.swap(WriteQueue.front().Callback);
  960. Stream->Write(WriteQueue.front().Request, OnWriteDoneTag.Prepare());
  961. WriteQueue.pop_front();
  962. } else {
  963. WriteActive = false;
  964. if (ReadFinished) {
  965. WriteFinished = true;
  966. Stream->Finish(&Status, OnFinishedTag.Prepare());
  967. }
  968. }
  969. }
  970. if (okCallback) {
  971. okCallback(TGrpcStatus());
  972. }
  973. }
  974. void OnFinished(bool ok) {
  975. TGrpcStatus status;
  976. std::deque<TWriteItem> writesDropped;
  977. std::vector<TReadCallback> finishedCallbacks;
  978. TConnectedCallback connectedCallback;
  979. TReadCallback readCallback;
  980. TReadCallback finishCallback;
  981. {
  982. std::unique_lock<std::mutex> guard(Mutex);
  983. Finished = true;
  984. FinishedOk = ok;
  985. LocalContext.reset();
  986. if (ok) {
  987. status = Status;
  988. } else if (Cancelled) {
  989. status = TGrpcStatus(grpc::StatusCode::CANCELLED, "Stream cancelled");
  990. } else {
  991. status = TGrpcStatus::Internal("Unexpected error");
  992. }
  993. writesDropped.swap(WriteQueue);
  994. finishedCallbacks.swap(FinishedCallbacks);
  995. if (ConnectedCallback) {
  996. Y_VERIFY(!ReadActive);
  997. connectedCallback = std::move(ConnectedCallback);
  998. ConnectedCallback = nullptr;
  999. } else if (ReadActive) {
  1000. if (ReadCallback) {
  1001. readCallback = std::move(ReadCallback);
  1002. ReadCallback = nullptr;
  1003. } else {
  1004. finishCallback = std::move(FinishCallback);
  1005. FinishCallback = nullptr;
  1006. }
  1007. ReadActive = false;
  1008. }
  1009. }
  1010. for (auto& item : writesDropped) {
  1011. if (item.Callback) {
  1012. TGrpcStatus writeStatus = status;
  1013. if (writeStatus.Ok()) {
  1014. writeStatus = TGrpcStatus(grpc::StatusCode::CANCELLED, "Write request dropped");
  1015. }
  1016. item.Callback(std::move(writeStatus));
  1017. }
  1018. }
  1019. for (auto& finishedCallback : finishedCallbacks) {
  1020. TGrpcStatus statusCopy = status;
  1021. finishedCallback(std::move(statusCopy));
  1022. }
  1023. if (connectedCallback) {
  1024. if (status.Ok()) {
  1025. status = TGrpcStatus(grpc::StatusCode::UNKNOWN, "Unknown stream failure");
  1026. }
  1027. connectedCallback(std::move(status), nullptr);
  1028. } else if (readCallback) {
  1029. if (status.Ok()) {
  1030. status = TGrpcStatus(grpc::StatusCode::OUT_OF_RANGE, "Read EOF");
  1031. }
  1032. readCallback(std::move(status));
  1033. } else if (finishCallback) {
  1034. finishCallback(std::move(status));
  1035. }
  1036. }
  1037. private:
  1038. struct TWriteItem {
  1039. TWriteCallback Callback;
  1040. TRequest Request;
  1041. };
  1042. private:
  1043. using TFixedEvent = TQueueClientFixedEvent<TSelf>;
  1044. TFixedEvent OnConnectedTag = { this, &TSelf::OnConnected };
  1045. TFixedEvent OnReadDoneTag = { this, &TSelf::OnReadDone };
  1046. TFixedEvent OnWriteDoneTag = { this, &TSelf::OnWriteDone };
  1047. TFixedEvent OnFinishedTag = { this, &TSelf::OnFinished };
  1048. private:
  1049. std::mutex Mutex;
  1050. TAsyncReaderWriterPtr Stream;
  1051. TConnectedCallback ConnectedCallback;
  1052. TReadCallback ReadCallback;
  1053. TReadCallback FinishCallback;
  1054. std::vector<TReadCallback> FinishedCallbacks;
  1055. std::deque<TWriteItem> WriteQueue;
  1056. TWriteCallback WriteCallback;
  1057. std::unordered_multimap<TString, TString>* InitialMetadata = nullptr;
  1058. bool Started = false;
  1059. bool HasInitialMetadata = false;
  1060. bool ReadActive = false;
  1061. bool ReadFinished = false;
  1062. bool WriteActive = false;
  1063. bool WriteFinished = false;
  1064. bool Finished = false;
  1065. bool Cancelled = false;
  1066. bool FinishedOk = false;
  1067. };
  1068. class TGRpcClientLow;
  1069. template<typename TGRpcService>
  1070. class TServiceConnection {
  1071. using TStub = typename TGRpcService::Stub;
  1072. friend class TGRpcClientLow;
  1073. public:
  1074. /*
  1075. * Start simple request
  1076. */
  1077. template<typename TRequest, typename TResponse>
  1078. void DoRequest(const TRequest& request,
  1079. TResponseCallback<TResponse> callback,
  1080. typename TSimpleRequestProcessor<TStub, TRequest, TResponse>::TAsyncRequest asyncRequest,
  1081. const TCallMeta& metas = { },
  1082. IQueueClientContextProvider* provider = nullptr)
  1083. {
  1084. auto processor = MakeIntrusive<TSimpleRequestProcessor<TStub, TRequest, TResponse>>(std::move(callback));
  1085. processor->ApplyMeta(metas);
  1086. processor->Start(*Stub_, asyncRequest, request, provider ? provider : Provider_);
  1087. }
  1088. /*
  1089. * Start simple request
  1090. */
  1091. template<typename TRequest, typename TResponse>
  1092. void DoAdvancedRequest(const TRequest& request,
  1093. TAdvancedResponseCallback<TResponse> callback,
  1094. typename TAdvancedRequestProcessor<TStub, TRequest, TResponse>::TAsyncRequest asyncRequest,
  1095. const TCallMeta& metas = { },
  1096. IQueueClientContextProvider* provider = nullptr)
  1097. {
  1098. auto processor = MakeIntrusive<TAdvancedRequestProcessor<TStub, TRequest, TResponse>>(std::move(callback));
  1099. processor->ApplyMeta(metas);
  1100. processor->Start(*Stub_, asyncRequest, request, provider ? provider : Provider_);
  1101. }
  1102. /*
  1103. * Start bidirectional streamming
  1104. */
  1105. template<typename TRequest, typename TResponse>
  1106. void DoStreamRequest(TStreamConnectedCallback<TRequest, TResponse> callback,
  1107. typename TStreamRequestReadWriteProcessor<TStub, TRequest, TResponse>::TAsyncRequest asyncRequest,
  1108. const TCallMeta& metas = { },
  1109. IQueueClientContextProvider* provider = nullptr)
  1110. {
  1111. auto processor = MakeIntrusive<TStreamRequestReadWriteProcessor<TStub, TRequest, TResponse>>(std::move(callback));
  1112. processor->ApplyMeta(metas);
  1113. processor->Start(*Stub_, std::move(asyncRequest), provider ? provider : Provider_);
  1114. }
  1115. /*
  1116. * Start streaming response reading (one request, many responses)
  1117. */
  1118. template<typename TRequest, typename TResponse>
  1119. void DoStreamRequest(const TRequest& request,
  1120. TStreamReaderCallback<TResponse> callback,
  1121. typename TStreamRequestReadProcessor<TStub, TRequest, TResponse>::TAsyncRequest asyncRequest,
  1122. const TCallMeta& metas = { },
  1123. IQueueClientContextProvider* provider = nullptr)
  1124. {
  1125. auto processor = MakeIntrusive<TStreamRequestReadProcessor<TStub, TRequest, TResponse>>(std::move(callback));
  1126. processor->ApplyMeta(metas);
  1127. processor->Start(*Stub_, request, std::move(asyncRequest), provider ? provider : Provider_);
  1128. }
  1129. private:
  1130. TServiceConnection(std::shared_ptr<grpc::ChannelInterface> ci,
  1131. IQueueClientContextProvider* provider)
  1132. : Stub_(TGRpcService::NewStub(ci))
  1133. , Provider_(provider)
  1134. {
  1135. Y_VERIFY(Provider_, "Connection does not have a queue provider");
  1136. }
  1137. TServiceConnection(TStubsHolder& holder,
  1138. IQueueClientContextProvider* provider)
  1139. : Stub_(holder.GetOrCreateStub<TStub>())
  1140. , Provider_(provider)
  1141. {
  1142. Y_VERIFY(Provider_, "Connection does not have a queue provider");
  1143. }
  1144. std::shared_ptr<TStub> Stub_;
  1145. IQueueClientContextProvider* Provider_;
  1146. };
  1147. class TGRpcClientLow
  1148. : public IQueueClientContextProvider
  1149. {
  1150. class TContextImpl;
  1151. friend class TContextImpl;
  1152. enum ECqState : TAtomicBase {
  1153. WORKING = 0,
  1154. STOP_SILENT = 1,
  1155. STOP_EXPLICIT = 2,
  1156. };
  1157. public:
  1158. explicit TGRpcClientLow(size_t numWorkerThread = DEFAULT_NUM_THREADS, bool useCompletionQueuePerThread = false);
  1159. ~TGRpcClientLow();
  1160. // Tries to stop all currently running requests (via their stop callbacks)
  1161. // Will shutdown CQ and drain events once all requests have finished
  1162. // No new requests may be started after this call
  1163. void Stop(bool wait = false);
  1164. // Waits until all currently running requests finish execution
  1165. void WaitIdle();
  1166. inline bool IsStopping() const {
  1167. switch (GetCqState()) {
  1168. case WORKING:
  1169. return false;
  1170. case STOP_SILENT:
  1171. case STOP_EXPLICIT:
  1172. return true;
  1173. }
  1174. Y_UNREACHABLE();
  1175. }
  1176. IQueueClientContextPtr CreateContext() override;
  1177. template<typename TGRpcService>
  1178. std::unique_ptr<TServiceConnection<TGRpcService>> CreateGRpcServiceConnection(const TGRpcClientConfig& config) {
  1179. return std::unique_ptr<TServiceConnection<TGRpcService>>(new TServiceConnection<TGRpcService>(CreateChannelInterface(config), this));
  1180. }
  1181. template<typename TGRpcService>
  1182. std::unique_ptr<TServiceConnection<TGRpcService>> CreateGRpcServiceConnection(TStubsHolder& holder) {
  1183. return std::unique_ptr<TServiceConnection<TGRpcService>>(new TServiceConnection<TGRpcService>(holder, this));
  1184. }
  1185. // Tests only, not thread-safe
  1186. void AddWorkerThreadForTest();
  1187. private:
  1188. using IThreadRef = std::unique_ptr<IThreadFactory::IThread>;
  1189. using CompletionQueueRef = std::unique_ptr<grpc::CompletionQueue>;
  1190. void Init(size_t numWorkerThread);
  1191. inline ECqState GetCqState() const { return (ECqState) AtomicGet(CqState_); }
  1192. inline void SetCqState(ECqState state) { AtomicSet(CqState_, state); }
  1193. void StopInternal(bool silent);
  1194. void WaitInternal();
  1195. void ForgetContext(TContextImpl* context);
  1196. private:
  1197. bool UseCompletionQueuePerThread_;
  1198. std::vector<CompletionQueueRef> CQS_;
  1199. std::vector<IThreadRef> WorkerThreads_;
  1200. TAtomic CqState_ = -1;
  1201. std::mutex Mtx_;
  1202. std::condition_variable ContextsEmpty_;
  1203. std::unordered_set<TContextImpl*> Contexts_;
  1204. std::mutex JoinMutex_;
  1205. };
  1206. } // namespace NGRpc