2017-08-30 05:30:27 -04:00
|
|
|
#ifndef OSMIUM_THREAD_QUEUE_HPP
|
|
|
|
#define OSMIUM_THREAD_QUEUE_HPP
|
|
|
|
|
|
|
|
/*
|
|
|
|
|
2020-11-17 16:59:06 -05:00
|
|
|
This file is part of Osmium (https://osmcode.org/libosmium).
|
2017-08-30 05:30:27 -04:00
|
|
|
|
2022-08-16 13:26:21 -04:00
|
|
|
Copyright 2013-2022 Jochen Topf <jochen@topf.org> and others (see README).
|
2017-08-30 05:30:27 -04:00
|
|
|
|
|
|
|
Boost Software License - Version 1.0 - August 17th, 2003
|
|
|
|
|
|
|
|
Permission is hereby granted, free of charge, to any person or organization
|
|
|
|
obtaining a copy of the software and accompanying documentation covered by
|
|
|
|
this license (the "Software") to use, reproduce, display, distribute,
|
|
|
|
execute, and transmit the Software, and to prepare derivative works of the
|
|
|
|
Software, and to permit third-parties to whom the Software is furnished to
|
|
|
|
do so, all subject to the following:
|
|
|
|
|
|
|
|
The copyright notices in the Software and this entire statement, including
|
|
|
|
the above license grant, this restriction and the following disclaimer,
|
|
|
|
must be included in all copies of the Software, in whole or in part, and
|
|
|
|
all derivative works of the Software, unless such copies or derivative
|
|
|
|
works are solely in the form of machine-executable object code generated by
|
|
|
|
a source language processor.
|
|
|
|
|
|
|
|
THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
|
|
|
|
IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
|
|
|
|
FITNESS FOR A PARTICULAR PURPOSE, TITLE AND NON-INFRINGEMENT. IN NO EVENT
|
|
|
|
SHALL THE COPYRIGHT HOLDERS OR ANYONE DISTRIBUTING THE SOFTWARE BE LIABLE
|
|
|
|
FOR ANY DAMAGES OR OTHER LIABILITY, WHETHER IN CONTRACT, TORT OR OTHERWISE,
|
|
|
|
ARISING FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER
|
|
|
|
DEALINGS IN THE SOFTWARE.
|
|
|
|
|
|
|
|
*/
|
|
|
|
|
2022-08-16 13:26:21 -04:00
|
|
|
#include <atomic>
|
2017-08-30 05:30:27 -04:00
|
|
|
#include <chrono>
|
|
|
|
#include <condition_variable>
|
|
|
|
#include <cstddef>
|
|
|
|
#include <mutex>
|
|
|
|
#include <queue>
|
|
|
|
#include <string>
|
|
|
|
#include <utility> // IWYU pragma: keep
|
|
|
|
|
|
|
|
#ifdef OSMIUM_DEBUG_QUEUE_SIZE
|
|
|
|
# include <iostream>
|
|
|
|
#endif
|
|
|
|
|
|
|
|
namespace osmium {
|
|
|
|
|
|
|
|
namespace thread {
|
|
|
|
|
|
|
|
/**
|
|
|
|
* A thread-safe queue.
|
|
|
|
*/
|
|
|
|
template <typename T>
|
|
|
|
class Queue {
|
|
|
|
|
|
|
|
/// Maximum size of this queue. If the queue is full pushing to
|
|
|
|
/// the queue will block.
|
|
|
|
const std::size_t m_max_size;
|
|
|
|
|
|
|
|
/// Name of this queue (for debugging only).
|
|
|
|
const std::string m_name;
|
|
|
|
|
|
|
|
mutable std::mutex m_mutex;
|
|
|
|
|
|
|
|
std::queue<T> m_queue;
|
|
|
|
|
|
|
|
/// Used to signal consumers when data is available in the queue.
|
|
|
|
std::condition_variable m_data_available;
|
|
|
|
|
|
|
|
/// Used to signal producers when queue is not full.
|
|
|
|
std::condition_variable m_space_available;
|
|
|
|
|
2022-08-16 13:26:21 -04:00
|
|
|
std::atomic<bool> m_in_use{true};
|
|
|
|
|
2017-08-30 05:30:27 -04:00
|
|
|
#ifdef OSMIUM_DEBUG_QUEUE_SIZE
|
|
|
|
/// The largest size the queue has been so far.
|
|
|
|
std::size_t m_largest_size;
|
|
|
|
|
|
|
|
/// The number of times push() was called on the queue.
|
|
|
|
std::atomic<int> m_push_counter;
|
|
|
|
|
|
|
|
/// The number of times the queue was full and a thread pushing
|
|
|
|
/// to the queue was blocked.
|
|
|
|
std::atomic<int> m_full_counter;
|
|
|
|
|
|
|
|
/**
|
|
|
|
* The number of times wait_and_pop(with_timeout)() was called
|
|
|
|
* on the queue.
|
|
|
|
*/
|
|
|
|
std::atomic<int> m_pop_counter;
|
|
|
|
|
|
|
|
/// The number of times the queue was full and a thread pushing
|
|
|
|
/// to the queue was blocked.
|
|
|
|
std::atomic<int> m_empty_counter;
|
|
|
|
#endif
|
|
|
|
|
|
|
|
public:
|
|
|
|
|
|
|
|
/**
|
|
|
|
* Construct a multithreaded queue.
|
|
|
|
*
|
|
|
|
* @param max_size Maximum number of elements in the queue. Set to
|
|
|
|
* 0 for an unlimited size.
|
|
|
|
* @param name Optional name for this queue. (Used for debugging.)
|
|
|
|
*/
|
2018-04-19 15:03:25 -04:00
|
|
|
explicit Queue(std::size_t max_size = 0, std::string name = "") :
|
2017-08-30 05:30:27 -04:00
|
|
|
m_max_size(max_size),
|
2018-04-19 15:03:25 -04:00
|
|
|
m_name(std::move(name)),
|
|
|
|
m_queue()
|
2017-08-30 05:30:27 -04:00
|
|
|
#ifdef OSMIUM_DEBUG_QUEUE_SIZE
|
|
|
|
,
|
|
|
|
m_largest_size(0),
|
|
|
|
m_push_counter(0),
|
|
|
|
m_full_counter(0),
|
|
|
|
m_pop_counter(0),
|
|
|
|
m_empty_counter(0)
|
|
|
|
#endif
|
|
|
|
{
|
|
|
|
}
|
|
|
|
|
2018-04-19 15:03:25 -04:00
|
|
|
Queue(const Queue&) = delete;
|
|
|
|
Queue& operator=(const Queue&) = delete;
|
|
|
|
|
|
|
|
Queue(Queue&&) = delete;
|
|
|
|
Queue& operator=(Queue&&) = delete;
|
|
|
|
|
2017-08-30 05:30:27 -04:00
|
|
|
#ifdef OSMIUM_DEBUG_QUEUE_SIZE
|
2018-04-19 15:03:25 -04:00
|
|
|
~Queue() {
|
2017-08-30 05:30:27 -04:00
|
|
|
std::cerr << "queue '" << m_name
|
|
|
|
<< "' with max_size=" << m_max_size
|
|
|
|
<< " had largest size " << m_largest_size
|
|
|
|
<< " and was full " << m_full_counter
|
|
|
|
<< " times in " << m_push_counter
|
|
|
|
<< " push() calls and was empty " << m_empty_counter
|
|
|
|
<< " times in " << m_pop_counter
|
|
|
|
<< " pop() calls\n";
|
|
|
|
}
|
2018-04-19 15:03:25 -04:00
|
|
|
#else
|
|
|
|
~Queue() = default;
|
|
|
|
#endif
|
2017-08-30 05:30:27 -04:00
|
|
|
|
|
|
|
/**
|
|
|
|
* Push an element onto the queue. If the queue has a max size,
|
|
|
|
* this call will block if the queue is full.
|
|
|
|
*/
|
|
|
|
void push(T value) {
|
2022-08-16 13:26:21 -04:00
|
|
|
if (!m_in_use) {
|
|
|
|
return;
|
|
|
|
}
|
2018-04-19 15:03:25 -04:00
|
|
|
constexpr const std::chrono::milliseconds max_wait{10};
|
2017-08-30 05:30:27 -04:00
|
|
|
#ifdef OSMIUM_DEBUG_QUEUE_SIZE
|
|
|
|
++m_push_counter;
|
|
|
|
#endif
|
|
|
|
if (m_max_size) {
|
|
|
|
while (size() >= m_max_size) {
|
|
|
|
std::unique_lock<std::mutex> lock{m_mutex};
|
|
|
|
m_space_available.wait_for(lock, max_wait, [this] {
|
|
|
|
return m_queue.size() < m_max_size;
|
|
|
|
});
|
|
|
|
#ifdef OSMIUM_DEBUG_QUEUE_SIZE
|
|
|
|
++m_full_counter;
|
|
|
|
#endif
|
|
|
|
}
|
|
|
|
}
|
|
|
|
std::lock_guard<std::mutex> lock{m_mutex};
|
|
|
|
m_queue.push(std::move(value));
|
|
|
|
#ifdef OSMIUM_DEBUG_QUEUE_SIZE
|
|
|
|
if (m_largest_size < m_queue.size()) {
|
|
|
|
m_largest_size = m_queue.size();
|
|
|
|
}
|
|
|
|
#endif
|
|
|
|
m_data_available.notify_one();
|
|
|
|
}
|
|
|
|
|
|
|
|
void wait_and_pop(T& value) {
|
|
|
|
#ifdef OSMIUM_DEBUG_QUEUE_SIZE
|
|
|
|
++m_pop_counter;
|
|
|
|
#endif
|
|
|
|
std::unique_lock<std::mutex> lock{m_mutex};
|
|
|
|
#ifdef OSMIUM_DEBUG_QUEUE_SIZE
|
|
|
|
if (m_queue.empty()) {
|
|
|
|
++m_empty_counter;
|
|
|
|
}
|
|
|
|
#endif
|
|
|
|
m_data_available.wait(lock, [this] {
|
2022-08-16 13:26:21 -04:00
|
|
|
return !m_in_use || !m_queue.empty();
|
2017-08-30 05:30:27 -04:00
|
|
|
});
|
|
|
|
if (!m_queue.empty()) {
|
|
|
|
value = std::move(m_queue.front());
|
|
|
|
m_queue.pop();
|
|
|
|
lock.unlock();
|
|
|
|
if (m_max_size) {
|
|
|
|
m_space_available.notify_one();
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
bool try_pop(T& value) {
|
|
|
|
#ifdef OSMIUM_DEBUG_QUEUE_SIZE
|
|
|
|
++m_pop_counter;
|
|
|
|
#endif
|
|
|
|
{
|
|
|
|
std::lock_guard<std::mutex> lock{m_mutex};
|
|
|
|
if (m_queue.empty()) {
|
|
|
|
#ifdef OSMIUM_DEBUG_QUEUE_SIZE
|
|
|
|
++m_empty_counter;
|
|
|
|
#endif
|
|
|
|
return false;
|
|
|
|
}
|
|
|
|
value = std::move(m_queue.front());
|
|
|
|
m_queue.pop();
|
|
|
|
}
|
|
|
|
if (m_max_size) {
|
|
|
|
m_space_available.notify_one();
|
|
|
|
}
|
|
|
|
return true;
|
|
|
|
}
|
|
|
|
|
|
|
|
bool empty() const {
|
|
|
|
std::lock_guard<std::mutex> lock{m_mutex};
|
|
|
|
return m_queue.empty();
|
|
|
|
}
|
|
|
|
|
|
|
|
std::size_t size() const {
|
|
|
|
std::lock_guard<std::mutex> lock{m_mutex};
|
|
|
|
return m_queue.size();
|
|
|
|
}
|
|
|
|
|
2022-08-16 13:26:21 -04:00
|
|
|
bool in_use() const noexcept {
|
|
|
|
return m_in_use;
|
|
|
|
}
|
|
|
|
|
|
|
|
void shutdown() {
|
|
|
|
m_in_use = false;
|
|
|
|
std::lock_guard<std::mutex> lock{m_mutex};
|
|
|
|
while (!m_queue.empty()) {
|
|
|
|
m_queue.pop();
|
|
|
|
}
|
|
|
|
m_data_available.notify_all();
|
|
|
|
}
|
|
|
|
|
2017-08-30 05:30:27 -04:00
|
|
|
}; // class Queue
|
|
|
|
|
|
|
|
} // namespace thread
|
|
|
|
|
|
|
|
} // namespace osmium
|
|
|
|
|
|
|
|
#endif // OSMIUM_THREAD_QUEUE_HPP
|