blocking_queue.h 3.67 KB
Newer Older
gab's avatar
gab committed
1 2
#pragma once

gabi's avatar
gabi committed
3 4 5 6
// blocking_queue:
// A blocking multi-consumer/multi-producer thread safe queue.
// Has max capacity and supports timeout on push or pop operations.

gab's avatar
gab committed
7 8 9 10 11 12
#include <chrono>
#include <memory>
#include <queue>
#include <mutex>
#include <condition_variable>

gabime's avatar
gabime committed
13 14
namespace c11log {
namespace details {
gabime's avatar
gabime committed
15

gab's avatar
gab committed
16
template<typename T>
gabime's avatar
gabime committed
17
class blocking_queue {
gab's avatar
gab committed
18
public:
gabime's avatar
gabime committed
19 20 21
    using queue_t = std::queue<T>;
    using size_type = typename queue_t::size_type;
    using clock = std::chrono::system_clock;
gabi's avatar
gabi committed
22

gabime's avatar
gabime committed
23
    explicit blocking_queue(size_type max_size) :
gabime's avatar
gabime committed
24 25 26
        _max_size(max_size),
        _q(),
        _mutex() {
gabime's avatar
gabime committed
27 28 29 30
    }
    blocking_queue(const blocking_queue&) = delete;
    blocking_queue& operator=(const blocking_queue&) = delete;
    ~blocking_queue() = default;
gab's avatar
gab committed
31

gabime's avatar
gabime committed
32
    size_type size() {
gabime's avatar
gabime committed
33 34
        std::lock_guard<std::mutex> lock(_mutex);
        return _q.size();
gabime's avatar
gabime committed
35
    }
gabi's avatar
gabi committed
36

gabime's avatar
gabime committed
37 38 39 40 41
    // Push copy of item into the back of the queue.
    // If the queue is full, block the calling thread util there is room or timeout have passed.
    // Return: false on timeout, true on successful push.
    template<typename Duration_Rep, typename Duration_Period, typename TT>
    bool push(TT&& item, const std::chrono::duration<Duration_Rep, Duration_Period>& timeout) {
gabime's avatar
gabime committed
42 43 44 45
        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;
gabime's avatar
gabime committed
46 47 48
            }))
            return false;
        }
gabime's avatar
gabime committed
49 50
        _q.push(std::forward<TT>(item));
        if (_q.size() <= 1) {
gabime's avatar
gabime committed
51
            ul.unlock(); //So the notified thread will have better chance to accuire the lock immediatly..
gabime's avatar
gabime committed
52
            _item_pushed_cond.notify_one();
gabime's avatar
gabime committed
53 54 55
        }
        return true;
    }
gabi's avatar
gabi committed
56

gabime's avatar
gabime committed
57 58 59 60
    // Push copy of item into the back of the queue.
    // If the queue is full, block the calling thread until there is room.
    template<typename TT>
    void push(TT&& item) {
gabime's avatar
gabime committed
61 62
		static constexpr std::chrono::hours one_hour(1);
        while (!push(std::forward<TT>(item), one_hour));
gabime's avatar
gabime committed
63
    }
gabi's avatar
gabi committed
64

gabime's avatar
gabime committed
65 66 67 68 69
    // Pop a copy of the front item in the queue into the given item ref.
    // If the queue is empty, block the calling thread util there is item to pop or timeout have passed.
    // Return: false on timeout , true on successful pop/
    template<class Duration_Rep, class Duration_Period>
    bool pop(T& item, const std::chrono::duration<Duration_Rep, Duration_Period>& timeout) {
gabime's avatar
gabime committed
70 71 72 73
        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();
gabime's avatar
gabime committed
74 75 76
            }))
            return false;
        }
gabime's avatar
gabime committed
77 78 79
        item = std::move(_q.front());
        _q.pop();
        if (_q.size() >= _max_size - 1) {
gabime's avatar
gabime committed
80
            ul.unlock(); //So the notified thread will have better chance to accuire the lock immediatly..
gabime's avatar
gabime committed
81
            _item_popped_cond.notify_one();
gabime's avatar
gabime committed
82 83 84
        }
        return true;
    }
gabi's avatar
gabi committed
85

gabime's avatar
gabime committed
86 87 88
    // Pop a copy of the front item in the queue into the given item ref.
    // If the queue is empty, block the calling thread util there is item to pop.
    void pop(T& item) {
gabime's avatar
gabime committed
89 90
		static constexpr std::chrono::hours one_hour(1);
        while (!pop(item, one_hour));
gabime's avatar
gabime committed
91
    }
gabi's avatar
gabi committed
92

gabime's avatar
gabime committed
93 94 95
    // Clear the queue
    void clear() {
        {
gabime's avatar
gabime committed
96 97
            std::unique_lock<std::mutex> ul(_mutex);
            queue_t().swap(_q);
gabime's avatar
gabime committed
98
        }
gabime's avatar
gabime committed
99
        _item_popped_cond.notify_all();
gabime's avatar
gabime committed
100
    }
gabi's avatar
gabi committed
101

gab's avatar
gab committed
102
private:
gabime's avatar
gabime committed
103 104 105 106 107
    size_type _max_size;
    std::queue<T> _q;
    std::mutex _mutex;
    std::condition_variable _item_pushed_cond;
    std::condition_variable _item_popped_cond;
gab's avatar
gab committed
108
};
gabime's avatar
gabime committed
109 110

}
gab's avatar
gab committed
111
}