2019-02-25 14:48:22 +01:00
|
|
|
/* Icinga 2 | (c) 2012 Icinga GmbH | GPLv2+ */
|
2012-06-24 02:56:48 +02:00
|
|
|
|
2013-03-25 18:36:15 +01:00
|
|
|
#ifndef THREADPOOL_H
|
|
|
|
#define THREADPOOL_H
|
2012-06-24 02:56:48 +02:00
|
|
|
|
2019-08-14 17:12:59 +02:00
|
|
|
#include "base/atomic.hpp"
|
2020-02-18 11:39:17 +01:00
|
|
|
#include "base/configuration.hpp"
|
2019-04-01 17:05:16 +02:00
|
|
|
#include "base/exception.hpp"
|
|
|
|
#include "base/logger.hpp"
|
|
|
|
#include <cstddef>
|
|
|
|
#include <exception>
|
|
|
|
#include <functional>
|
|
|
|
#include <memory>
|
2017-11-21 12:12:58 +01:00
|
|
|
#include <thread>
|
2019-04-01 17:05:16 +02:00
|
|
|
#include <boost/asio/post.hpp>
|
|
|
|
#include <boost/asio/thread_pool.hpp>
|
|
|
|
#include <boost/thread/locks.hpp>
|
|
|
|
#include <boost/thread/shared_mutex.hpp>
|
2019-08-14 17:12:59 +02:00
|
|
|
#include <cstdint>
|
2013-03-15 18:21:29 +01:00
|
|
|
|
2012-06-24 02:56:48 +02:00
|
|
|
namespace icinga
|
|
|
|
{
|
|
|
|
|
2014-09-11 11:45:21 +02:00
|
|
|
enum SchedulerPolicy
|
|
|
|
{
|
|
|
|
DefaultScheduler,
|
|
|
|
LowLatencyScheduler
|
|
|
|
};
|
|
|
|
|
2012-09-17 13:35:55 +02:00
|
|
|
/**
|
2013-03-25 18:36:15 +01:00
|
|
|
* A thread pool.
|
2012-09-17 13:35:55 +02:00
|
|
|
*
|
|
|
|
* @ingroup base
|
|
|
|
*/
|
2017-12-31 07:22:16 +01:00
|
|
|
class ThreadPool
|
2012-06-24 02:56:48 +02:00
|
|
|
{
|
|
|
|
public:
|
2017-12-24 06:35:12 +01:00
|
|
|
typedef std::function<void ()> WorkFunction;
|
2013-03-25 18:36:15 +01:00
|
|
|
|
2023-01-27 16:32:29 +01:00
|
|
|
ThreadPool();
|
2018-01-04 04:25:35 +01:00
|
|
|
~ThreadPool();
|
2012-06-24 02:56:48 +02:00
|
|
|
|
2018-01-04 04:25:35 +01:00
|
|
|
void Start();
|
|
|
|
void Stop();
|
2023-05-23 14:41:35 +02:00
|
|
|
void Restart();
|
2013-02-18 14:40:24 +01:00
|
|
|
|
2019-04-01 17:05:16 +02:00
|
|
|
/**
|
|
|
|
* Appends a work item to the work queue. Work items will be processed in FIFO order.
|
|
|
|
*
|
|
|
|
* @param callback The callback function for the work item.
|
|
|
|
* @returns true if the item was queued, false otherwise.
|
|
|
|
*/
|
|
|
|
template<class T>
|
|
|
|
bool Post(T callback, SchedulerPolicy)
|
2013-12-06 21:46:50 +01:00
|
|
|
{
|
2019-04-01 17:05:16 +02:00
|
|
|
boost::shared_lock<decltype(m_Mutex)> lock (m_Mutex);
|
|
|
|
|
|
|
|
if (m_Pool) {
|
2019-08-14 17:12:59 +02:00
|
|
|
m_Pending.fetch_add(1);
|
|
|
|
|
|
|
|
boost::asio::post(*m_Pool, [this, callback]() {
|
|
|
|
m_Pending.fetch_sub(1);
|
|
|
|
|
2019-04-01 17:05:16 +02:00
|
|
|
try {
|
|
|
|
callback();
|
|
|
|
} catch (const std::exception& ex) {
|
|
|
|
Log(LogCritical, "ThreadPool")
|
|
|
|
<< "Exception thrown in event handler:\n"
|
|
|
|
<< DiagnosticInformation(ex);
|
|
|
|
} catch (...) {
|
|
|
|
Log(LogCritical, "ThreadPool", "Exception of unknown type thrown in event handler.");
|
|
|
|
}
|
|
|
|
});
|
|
|
|
|
|
|
|
return true;
|
|
|
|
} else {
|
|
|
|
return false;
|
|
|
|
}
|
|
|
|
}
|
2013-03-25 18:36:15 +01:00
|
|
|
|
2019-08-14 17:12:59 +02:00
|
|
|
/**
|
|
|
|
* Returns the amount of queued tasks not started yet.
|
|
|
|
*
|
|
|
|
* @returns amount of queued tasks.
|
|
|
|
*/
|
|
|
|
inline uint_fast64_t GetPending()
|
|
|
|
{
|
|
|
|
return m_Pending.load();
|
|
|
|
}
|
|
|
|
|
2019-04-01 17:05:16 +02:00
|
|
|
private:
|
|
|
|
boost::shared_mutex m_Mutex;
|
|
|
|
std::unique_ptr<boost::asio::thread_pool> m_Pool;
|
2019-08-14 17:12:59 +02:00
|
|
|
Atomic<uint_fast64_t> m_Pending;
|
2023-05-23 14:41:35 +02:00
|
|
|
|
|
|
|
void InitializePool();
|
2012-06-24 02:56:48 +02:00
|
|
|
};
|
|
|
|
|
|
|
|
}
|
|
|
|
|
2019-04-23 16:59:49 +02:00
|
|
|
#endif /* THREADPOOL_H */
|