57 Status
Push(
const T& element, int64_t timeout_ms);
66 Status
Pop(T* element, int64_t timeout_ms);
80 Status
Peek(
size_t pos, T* element, int64_t timeout_ms);
86 absl::MutexLock lock(mutex_);
87 stop_requested_ =
true;
88 not_empty_cv_.SignalAll();
89 not_full_cv_.SignalAll();
90 new_element_cv_.SignalAll();
95 absl::MutexLock lock(mutex_);
101 absl::MutexLock lock(mutex_);
108 absl::MutexLock lock(mutex_);
115 absl::MutexLock lock(mutex_);
116 return head_pos_ + q_.size() - 1;
122 absl::MutexLock lock(mutex_);
123 return stop_requested_;
128 void SlideHeadOnCenter(
size_t pos);
130 const size_t capacity_;
132 mutable absl::Mutex mutex_;
133 size_t head_pos_ ABSL_GUARDED_BY(mutex_);
135 ABSL_GUARDED_BY(mutex_);
136 absl::CondVar not_empty_cv_ ABSL_GUARDED_BY(mutex_);
137 absl::CondVar not_full_cv_ ABSL_GUARDED_BY(mutex_);
138 absl::CondVar new_element_cv_ ABSL_GUARDED_BY(mutex_);
140 ABSL_GUARDED_BY(mutex_);
160 absl::MutexLock lock(mutex_);
165 return Status(error::STOPPED,
"");
167 auto start = std::chrono::steady_clock::now();
168 auto timeout_delta = std::chrono::milliseconds(timeout_ms);
171 while (q_.size() == capacity_) {
172 if (timeout_ms < 0) {
174 not_full_cv_.Wait(&mutex_);
176 auto elapsed = std::chrono::steady_clock::now() - start;
177 if (elapsed < timeout_delta) {
179 not_full_cv_.WaitWithTimeout(
180 &mutex_, absl::FromChrono(timeout_delta - elapsed));
183 return Status(error::TIME_OUT,
"Time out on pushing.");
188 return Status(error::STOPPED,
"");
192 DCHECK_LT(q_.size(), capacity_);
197 not_empty_cv_.Signal();
198 new_element_cv_.Signal();
200 q_.push_back(element);
203 if (woken && q_.size() != capacity_)
204 not_full_cv_.Signal();
210 absl::MutexLock lock(mutex_);
213 auto start = std::chrono::steady_clock::now();
214 auto timeout_delta = std::chrono::milliseconds(timeout_ms);
218 return Status(error::STOPPED,
"");
220 if (timeout_ms < 0) {
222 not_empty_cv_.Wait(&mutex_);
224 auto elapsed = std::chrono::steady_clock::now() - start;
225 if (elapsed < timeout_delta) {
227 not_empty_cv_.WaitWithTimeout(
228 &mutex_, absl::FromChrono(timeout_delta - elapsed));
231 return Status(error::TIME_OUT,
"Time out on popping.");
238 if (q_.size() == capacity_)
239 not_full_cv_.Signal();
241 *element = q_.front();
246 if (woken && !q_.empty())
247 not_empty_cv_.Signal();
254 int64_t timeout_ms) {
255 absl::MutexLock lock(mutex_);
256 if (pos < head_pos_) {
257 return Status(error::INVALID_ARGUMENT,
258 absl::StrFormat(
"pos (%zu) is too small; head is at %zu.",
264 auto start = std::chrono::steady_clock::now();
265 auto timeout_delta = std::chrono::milliseconds(timeout_ms);
268 SlideHeadOnCenter(pos);
270 while (pos >= head_pos_ + q_.size()) {
272 return Status(error::STOPPED,
"");
274 if (timeout_ms < 0) {
276 new_element_cv_.Wait(&mutex_);
278 auto elapsed = std::chrono::steady_clock::now() - start;
279 if (elapsed < timeout_delta) {
281 new_element_cv_.WaitWithTimeout(
282 &mutex_, absl::FromChrono(timeout_delta - elapsed));
285 return Status(error::TIME_OUT,
"Time out on peeking.");
289 SlideHeadOnCenter(pos);
293 *element = q_[pos - head_pos_];
296 if (woken && !q_.empty())
297 new_element_cv_.Signal();