00001
00002
00003
00004
00005
00006
00007
00008
00009 #ifndef __libutilxx__sella__util__ThreadPool_H__
00010 #define __libutilxx__sella__util__ThreadPool_H__
00011
00012 #include "../../common.h"
00013 #include "Task.h"
00014 #include "TaskQueue.h"
00015 #include "Log.h"
00016
00017 #include <sched.h>
00018
00019 #include <set>
00020 #include <thread>
00021 #include <chrono>
00022 #include <stdatomic.h>
00023 #include <condition_variable>
00024
00025 namespace sella {
00026 namespace util {
00027 class ThreadPool;
00028 }
00029 }
00030
00031 class sella::util::ThreadPool : protected TaskQueue {
00032 public:
00033 typedef std::shared_ptr<ThreadPool> shared;
00034
00035 public:
00036 static const size_t DefaultMaximumWorkers = 0;
00037 static const size_t DefaultMinimumWorkers = 2;
00038 static const size_t DefaultSpareWorkers = 2;
00039
00040 public:
00041 ThreadPool(size_t maximumWorkers = DefaultMaximumWorkers);
00042 ThreadPool(size_t maximumWorkers, size_t minimumWorkers, size_t spareWorkers = DefaultSpareWorkers);
00043 ThreadPool(Log::shared log, size_t maximumWorkers = DefaultMaximumWorkers);
00044 ThreadPool(Log::shared log, size_t maximumWorkers, size_t minimumWorkers, size_t spareWorkers = DefaultSpareWorkers);
00045 ThreadPool(const ThreadPool &other) = delete;
00046 virtual ~ThreadPool();
00047
00048 void schedule(Task::shared &task, int priority = 0);
00049 void schedule(Task::shared &task, const std::chrono::system_clock::duration &delay, int priority = 0);
00050 void schedule(Task::shared &task, const std::chrono::system_clock::time_point &time, int priority = 0);
00051 void schedule(Task::shared &task, time_t sec, suseconds_t usec, int priority = 0);
00052 void schedule(Task::shared &task, timeval &tv, int priority = 0);
00053 bool reschedule(Task::shared &task, const std::chrono::system_clock::duration &delay, int priority = 0, bool strict = false);
00054 bool reschedule(Task::shared &task, const std::chrono::system_clock::time_point &time, int priority = 0, bool strict = false);
00055 bool reschedule(Task::shared &task, time_t sec, suseconds_t usec = 0, int priority = 0, bool strict = false);
00056 bool reschedule(Task::shared &task, timeval &tv, int priority = 0, bool strict = false);
00057 bool unschedule(Task::shared &task, bool firstonly = true);
00058
00059 void pause(bool on = true);
00060 void resume(void) { pause(false); }
00061
00062 size_t getCPUCount(void);
00063 int getCPUSetCount(void);
00064 bool addCPUAffinity(int cpu);
00065 bool addCPUAffinity(std::set<int> cpus);
00066 bool setCPUAffinity(cpu_set_t &cpuset);
00067 bool clearCPUAffinity(void);
00068 bool clearCPUAffinity(int cpu);
00069 bool clearCPUAffinity(std::set<int> cpus);
00070
00071 size_t size(void) const;
00072 size_t workers(void) const;
00073 size_t active(void) const;
00074 bool empty(void) const;
00075 void clear(void);
00076
00077 Task::shared next(void);
00078 std::chrono::microseconds duration(void) const;
00079 std::chrono::microseconds max(void) const;
00080 bool ready(void) const;
00081
00082 bool wait(void);
00083 bool wait(const std::chrono::system_clock::duration &delay);
00084 bool wait(const std::chrono::system_clock::time_point &time);
00085
00086 using TaskQueue::capacity;
00087 void reserve(size_t capacity);
00088 void print(void) const;
00089 using TaskQueue::setLog;
00090
00091 static ThreadPool::shared shared_ptr(size_t maximumWorkers = DefaultMaximumWorkers);
00092 static ThreadPool::shared shared_ptr(size_t maximumWorkers, size_t minimumWorkers, size_t spareWorkers = DefaultSpareWorkers);
00093 static ThreadPool::shared shared_ptr(Log::shared log, size_t maximumWorkers = DefaultMaximumWorkers);
00094 static ThreadPool::shared shared_ptr(Log::shared log, size_t maximumWorkers, size_t minimumWorkers, size_t spareWorkers = DefaultSpareWorkers);
00095
00096 protected:
00097 std::thread managerTID;
00098 std::set<pthread_t> tids;
00099 std::atomic<size_t> threads;
00100 std::atomic<size_t> creating;
00101 std::atomic<size_t> actives;
00102 std::atomic<size_t> overage;
00103 bool shutdown;
00104 bool paused;
00105 bool allcpus;
00106 cpu_set_t cpuset;
00107
00108 mutable std::mutex heapMutex;
00109 mutable std::mutex tidsMutex;
00110 mutable std::mutex managerMutex;
00111 mutable std::mutex workerMutex;
00112 mutable std::mutex overageMutex;
00113 mutable std::mutex emptyMutex;
00114 std::condition_variable managerCV;
00115 std::condition_variable workerCV;
00116 std::condition_variable emptyCV;
00117
00118 size_t maximum;
00119 size_t minimum;
00120 size_t spares;
00121
00122 void manager(void);
00123 void worker(void);
00124 void normalize(void);
00125 void createWorker(void);
00126 bool setThreadsAffinity(void);
00127 bool setThreadAffinity(pthread_t tid);
00128 };
00129
00130 #endif
00131
00132
00133
00134