From 22f8ae0e3f5d68b92aecccdf96c1d841a0334311 Mon Sep 17 00:00:00 2001 From: qrort Date: Wed, 30 Nov 2022 23:47:12 +0300 Subject: validate canons without yatest_common --- .../threading/blocking_queue/blocking_queue.cpp | 3 + .../cpp/threading/blocking_queue/blocking_queue.h | 158 +++++++++++++++++++++ 2 files changed, 161 insertions(+) create mode 100644 library/cpp/threading/blocking_queue/blocking_queue.cpp create mode 100644 library/cpp/threading/blocking_queue/blocking_queue.h (limited to 'library/cpp/threading/blocking_queue') diff --git a/library/cpp/threading/blocking_queue/blocking_queue.cpp b/library/cpp/threading/blocking_queue/blocking_queue.cpp new file mode 100644 index 00000000000..db199c80be5 --- /dev/null +++ b/library/cpp/threading/blocking_queue/blocking_queue.cpp @@ -0,0 +1,3 @@ +#include "blocking_queue.h" + +// just check compilability diff --git a/library/cpp/threading/blocking_queue/blocking_queue.h b/library/cpp/threading/blocking_queue/blocking_queue.h new file mode 100644 index 00000000000..48d3762f68a --- /dev/null +++ b/library/cpp/threading/blocking_queue/blocking_queue.h @@ -0,0 +1,158 @@ +#pragma once + +#include +#include +#include +#include +#include +#include + +#include + +namespace NThreading { + /// + /// TBlockingQueue is a queue of elements of limited or unlimited size. + /// Queue provides Push and Pop operations that block if operation can't be executed + /// (queue is empty or maximum size is reached). + /// + /// Queue can be stopped, in that case all blocked operation will return `Nothing` / false. + /// + /// All operations are thread safe. + /// + /// + /// Example of usage: + /// TBlockingQueue queue; + /// + /// ... + /// + /// // thread 1 + /// queue.Push(42); + /// queue.Push(100500); + /// + /// ... + /// + /// // thread 2 + /// while (TMaybe number = queue.Pop()) { + /// ProcessNumber(number.GetRef()); + /// } + template + class TBlockingQueue { + public: + /// + /// Creates blocking queue with given maxSize + /// if maxSize == 0 then queue is unlimited + TBlockingQueue(size_t maxSize) + : MaxSize(maxSize == 0 ? Max() : maxSize) + , Stopped(false) + { + } + + /// + /// Blocks until queue has some elements or queue is stopped or deadline is reached. + /// Returns `Nothing` if queue is stopped or deadline is reached. + /// Returns element otherwise. + TMaybe Pop(TInstant deadline = TInstant::Max()) { + TGuard g(Lock); + + const auto canPop = [this]() { return CanPop(); }; + if (!CanPopCV.WaitD(Lock, deadline, canPop)) { + return Nothing(); + } + + if (Stopped && Queue.empty()) { + return Nothing(); + } + TElement e = std::move(Queue.front()); + Queue.pop_front(); + CanPushCV.Signal(); + return std::move(e); + } + + TMaybe Pop(TDuration duration) { + return Pop(TInstant::Now() + duration); + } + + /// + /// Blocks until queue has space for new elements or queue is stopped or deadline is reached. + /// Returns false exception if queue is stopped and push failed or deadline is reached. + /// Pushes element to queue and returns true otherwise. + bool Push(const TElement& e, TInstant deadline = TInstant::Max()) { + return PushRef(e, deadline); + } + + bool Push(TElement&& e, TInstant deadline = TInstant::Max()) { + return PushRef(std::move(e), deadline); + } + + bool Push(const TElement& e, TDuration duration) { + return Push(e, TInstant::Now() + duration); + } + + bool Push(TElement&& e, TDuration duration) { + return Push(std::move(e), TInstant::Now() + duration); + } + + /// + /// Stops the queue, all blocked operations will be aborted. + void Stop() { + TGuard g(Lock); + Stopped = true; + CanPopCV.BroadCast(); + CanPushCV.BroadCast(); + } + + /// + /// Checks whether queue is empty. + bool Empty() const { + TGuard g(Lock); + return Queue.empty(); + } + + /// + /// Returns size of the queue. + size_t Size() const { + TGuard g(Lock); + return Queue.size(); + } + + /// + /// Checks whether queue is stopped. + bool IsStopped() const { + TGuard g(Lock); + return Stopped; + } + + private: + bool CanPush() const { + return Queue.size() < MaxSize || Stopped; + } + + bool CanPop() const { + return !Queue.empty() || Stopped; + } + + template + bool PushRef(Ref e, TInstant deadline) { + TGuard g(Lock); + const auto canPush = [this]() { return CanPush(); }; + if (!CanPushCV.WaitD(Lock, deadline, canPush)) { + return false; + } + if (Stopped) { + return false; + } + Queue.push_back(std::forward(e)); + CanPopCV.Signal(); + return true; + } + + private: + TMutex Lock; + TCondVar CanPopCV; + TCondVar CanPushCV; + TDeque Queue; + size_t MaxSize; + bool Stopped; + }; + +} -- cgit v1.3