pipequeue.h 4.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207
  1. #pragma once
  2. #include "lfqueue.h"
  3. #include <library/cpp/coroutine/engine/impl.h>
  4. #include <library/cpp/coroutine/engine/network.h>
  5. #include <library/cpp/deprecated/atomic/atomic.h>
  6. #include <util/system/pipe.h>
  7. #ifdef _linux_
  8. #include <sys/eventfd.h>
  9. #endif
  10. #if defined(_bionic_) && !defined(EFD_SEMAPHORE)
  11. #define EFD_SEMAPHORE 1
  12. #endif
  13. namespace NNeh {
  14. #ifdef _linux_
  15. class TSemaphoreEventFd {
  16. public:
  17. inline TSemaphoreEventFd() {
  18. F_ = eventfd(0, EFD_NONBLOCK | EFD_SEMAPHORE);
  19. if (F_ < 0) {
  20. ythrow TFileError() << "failed to create a eventfd";
  21. }
  22. }
  23. inline ~TSemaphoreEventFd() {
  24. close(F_);
  25. }
  26. inline size_t Acquire(TCont* c) {
  27. ui64 ev;
  28. return NCoro::ReadI(c, F_, &ev, sizeof ev).Processed();
  29. }
  30. inline void Release() {
  31. const static ui64 ev(1);
  32. (void)write(F_, &ev, sizeof ev);
  33. }
  34. private:
  35. int F_;
  36. };
  37. #endif
  38. class TSemaphorePipe {
  39. public:
  40. inline TSemaphorePipe() {
  41. TPipeHandle::Pipe(S_[0], S_[1]);
  42. SetNonBlock(S_[0]);
  43. SetNonBlock(S_[1]);
  44. }
  45. inline size_t Acquire(TCont* c) {
  46. char ch;
  47. return NCoro::ReadI(c, S_[0], &ch, 1).Processed();
  48. }
  49. inline size_t Acquire(TCont* c, char* buff, size_t buflen) {
  50. return NCoro::ReadI(c, S_[0], buff, buflen).Processed();
  51. }
  52. inline void Release() {
  53. char ch = 13;
  54. S_[1].Write(&ch, 1);
  55. }
  56. private:
  57. TPipeHandle S_[2];
  58. };
  59. class TPipeQueueBase {
  60. public:
  61. inline void Enqueue(void* job) {
  62. Q_.Enqueue(job);
  63. S_.Release();
  64. }
  65. inline void* Dequeue(TCont* c, char* ch, size_t buflen) {
  66. void* ret = nullptr;
  67. while (!Q_.Dequeue(&ret) && S_.Acquire(c, ch, buflen)) {
  68. }
  69. return ret;
  70. }
  71. inline void* Dequeue() noexcept {
  72. void* ret = nullptr;
  73. Q_.Dequeue(&ret);
  74. return ret;
  75. }
  76. private:
  77. TLockFreeQueue<void*> Q_;
  78. TSemaphorePipe S_;
  79. };
  80. template <class T, size_t buflen = 1>
  81. class TPipeQueue {
  82. public:
  83. template <class TPtr>
  84. inline void EnqueueSafe(TPtr req) {
  85. Enqueue(req.Get());
  86. req.Release();
  87. }
  88. inline void Enqueue(T* req) {
  89. Q_.Enqueue(req);
  90. }
  91. template <class TPtr>
  92. inline void DequeueSafe(TCont* c, TPtr& ret) {
  93. ret.Reset(Dequeue(c));
  94. }
  95. inline T* Dequeue(TCont* c) {
  96. char ch[buflen];
  97. return (T*)Q_.Dequeue(c, ch, sizeof(ch));
  98. }
  99. protected:
  100. TPipeQueueBase Q_;
  101. };
  102. //optimized for avoiding unnecessary usage semaphore + use eventfd on linux
  103. template <class T>
  104. struct TOneConsumerPipeQueue {
  105. inline TOneConsumerPipeQueue()
  106. : Signaled_(0)
  107. , SkipWait_(0)
  108. {
  109. }
  110. inline void Enqueue(T* job) {
  111. Q_.Enqueue(job);
  112. AtomicSet(SkipWait_, 1);
  113. if (AtomicCas(&Signaled_, 1, 0)) {
  114. S_.Release();
  115. }
  116. }
  117. inline T* Dequeue(TCont* c) {
  118. T* ret = nullptr;
  119. while (!Q_.Dequeue(&ret)) {
  120. AtomicSet(Signaled_, 0);
  121. if (!AtomicCas(&SkipWait_, 0, 1)) {
  122. if (!S_.Acquire(c)) {
  123. break;
  124. }
  125. }
  126. AtomicSet(Signaled_, 1);
  127. }
  128. return ret;
  129. }
  130. template <class TPtr>
  131. inline void EnqueueSafe(TPtr req) {
  132. Enqueue(req.Get());
  133. Y_UNUSED(req.Release());
  134. }
  135. template <class TPtr>
  136. inline void DequeueSafe(TCont* c, TPtr& ret) {
  137. ret.Reset(Dequeue(c));
  138. }
  139. protected:
  140. TLockFreeQueue<T*> Q_;
  141. #ifdef _linux_
  142. TSemaphoreEventFd S_;
  143. #else
  144. TSemaphorePipe S_;
  145. #endif
  146. TAtomic Signaled_;
  147. TAtomic SkipWait_;
  148. };
  149. template <class T, size_t buflen = 1>
  150. struct TAutoPipeQueue: public TPipeQueue<T, buflen> {
  151. ~TAutoPipeQueue() {
  152. while (T* t = (T*)TPipeQueue<T, buflen>::Q_.Dequeue()) {
  153. delete t;
  154. }
  155. }
  156. };
  157. template <class T>
  158. struct TAutoOneConsumerPipeQueue: public TOneConsumerPipeQueue<T> {
  159. ~TAutoOneConsumerPipeQueue() {
  160. T* ret = nullptr;
  161. while (TOneConsumerPipeQueue<T>::Q_.Dequeue(&ret)) {
  162. delete ret;
  163. }
  164. }
  165. };
  166. }