1. Added test
2. Fixes
This commit is contained in:
		@@ -21,17 +21,16 @@ public:
 | 
			
		||||
    using size_type = typename queue_t::size_type;
 | 
			
		||||
    using clock = std::chrono::system_clock;
 | 
			
		||||
 | 
			
		||||
    explicit blocking_queue(size_type max_size) :_max_size(max_size), _q()
 | 
			
		||||
    explicit blocking_queue(size_type max_size) :max_size_(max_size), q_()
 | 
			
		||||
    {}
 | 
			
		||||
    blocking_queue(const blocking_queue&) = delete;
 | 
			
		||||
    blocking_queue& operator=(const blocking_queue&) = delete;
 | 
			
		||||
    blocking_queue& operator=(const blocking_queue&) volatile = delete;
 | 
			
		||||
    blocking_queue& operator=(const blocking_queue&) = delete;    
 | 
			
		||||
    ~blocking_queue() = default;
 | 
			
		||||
 | 
			
		||||
    size_type size()
 | 
			
		||||
    {
 | 
			
		||||
        std::lock_guard<std::mutex> lock(_mutex);
 | 
			
		||||
        return _q.size();
 | 
			
		||||
        std::lock_guard<std::mutex> lock(mutex_);
 | 
			
		||||
        return q_.size();
 | 
			
		||||
    }
 | 
			
		||||
 | 
			
		||||
    // Push copy of item into the back of the queue. 
 | 
			
		||||
@@ -40,17 +39,17 @@ public:
 | 
			
		||||
    template<class Duration_Rep, class Duration_Period>
 | 
			
		||||
    bool push(const T& item, const std::chrono::duration<Duration_Rep, Duration_Period>& timeout)
 | 
			
		||||
    {
 | 
			
		||||
        std::unique_lock<std::mutex> ul(_mutex);
 | 
			
		||||
        if (_q.size() >= _max_size)
 | 
			
		||||
        std::unique_lock<std::mutex> ul(mutex_);
 | 
			
		||||
        if (q_.size() >= max_size_)
 | 
			
		||||
        {
 | 
			
		||||
            if (!_item_popped_cond.wait_until(ul, clock::now() + timeout, [this]() { return this->_q.size() < this->_max_size; }))
 | 
			
		||||
            if (!item_popped_cond_.wait_until(ul, clock::now() + timeout, [this]() { return this->q_.size() < this->max_size_; }))
 | 
			
		||||
                return false;
 | 
			
		||||
        }
 | 
			
		||||
        _q.push(item);
 | 
			
		||||
        if (_q.size() <= 1)
 | 
			
		||||
        q_.push(item);
 | 
			
		||||
        if (q_.size() <= 1)
 | 
			
		||||
        {
 | 
			
		||||
            ul.unlock(); //So the notified thread will have better chance to accuire the lock immediatly..
 | 
			
		||||
            _item_pushed_cond.notify_one();
 | 
			
		||||
            item_pushed_cond_.notify_one();
 | 
			
		||||
        }
 | 
			
		||||
        return true;
 | 
			
		||||
    }
 | 
			
		||||
@@ -68,18 +67,18 @@ public:
 | 
			
		||||
    template<class Duration_Rep, class Duration_Period>
 | 
			
		||||
    bool pop(T& item, const std::chrono::duration<Duration_Rep, Duration_Period>& timeout)
 | 
			
		||||
    {
 | 
			
		||||
        std::unique_lock<std::mutex> ul(_mutex);
 | 
			
		||||
        if (_q.empty())
 | 
			
		||||
        std::unique_lock<std::mutex> ul(mutex_);
 | 
			
		||||
        if (q_.empty())
 | 
			
		||||
        {
 | 
			
		||||
            if (!_item_pushed_cond.wait_until(ul, clock::now() + timeout, [this]() { return !this->_q.empty(); }))
 | 
			
		||||
            if (!item_pushed_cond_.wait_until(ul, clock::now() + timeout, [this]() { return !this->q_.empty(); }))
 | 
			
		||||
                return false;
 | 
			
		||||
        }
 | 
			
		||||
        item = _q.front();
 | 
			
		||||
        _q.pop();
 | 
			
		||||
        if (_q.size() >= _max_size - 1)
 | 
			
		||||
        item = q_.front();
 | 
			
		||||
        q_.pop();
 | 
			
		||||
        if (q_.size() >= max_size_ - 1)
 | 
			
		||||
        {
 | 
			
		||||
            ul.unlock(); //So the notified thread will have better chance to accuire the lock immediatly..
 | 
			
		||||
            _item_popped_cond.notify_one();
 | 
			
		||||
            item_popped_cond_.notify_one();
 | 
			
		||||
        }
 | 
			
		||||
        return true;
 | 
			
		||||
    }
 | 
			
		||||
@@ -91,12 +90,19 @@ public:
 | 
			
		||||
        while (!pop(item, std::chrono::hours::max()));
 | 
			
		||||
    }
 | 
			
		||||
 | 
			
		||||
    // Clear the queue
 | 
			
		||||
    void clear()
 | 
			
		||||
    {
 | 
			
		||||
        T item;
 | 
			
		||||
        while (pop(item, std::chrono::milliseconds(0)));        
 | 
			
		||||
    }
 | 
			
		||||
 | 
			
		||||
private:
 | 
			
		||||
    size_type _max_size;
 | 
			
		||||
    std::queue<T> _q;
 | 
			
		||||
    std::mutex _mutex;
 | 
			
		||||
    std::condition_variable _item_pushed_cond;
 | 
			
		||||
    std::condition_variable _item_popped_cond;
 | 
			
		||||
    size_type max_size_;
 | 
			
		||||
    std::queue<T> q_;
 | 
			
		||||
    std::mutex mutex_;
 | 
			
		||||
    std::condition_variable item_pushed_cond_;
 | 
			
		||||
    std::condition_variable item_popped_cond_;
 | 
			
		||||
};
 | 
			
		||||
}
 | 
			
		||||
}
 | 
			
		||||
		Reference in New Issue
	
	Block a user