File: ThreadPool.hpp

package info (click to toggle)
sortmerna 4.3.7-3
  • links: PTS, VCS
  • area: main
  • in suites: forky, sid
  • size: 134,048 kB
  • sloc: cpp: 24,424; ansic: 15,923; python: 1,453; sh: 224; makefile: 31
file content (152 lines) | stat: -rw-r--r-- 4,186 bytes parent folder | download | duplicates (3)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
/*
 @copyright 2016-2021  Clarity Genomics BVBA
 @copyright 2012-2016  Bonsai Bioinformatics Research Group
 @copyright 2014-2016  Knight Lab, Department of Pediatrics, UCSD, La Jolla

 @parblock
 SortMeRNA - next-generation reads filter for metatranscriptomic or total RNA
 This is a free software: you can redistribute it and/or modify
 it under the terms of the GNU Lesser General Public License as published by
 the Free Software Foundation, either version 3 of the License, or
 (at your option) any later version.

 SortMeRNA is distributed in the hope that it will be useful,
 but WITHOUT ANY WARRANTY; without even the implied warranty of
 MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE.  See the
 GNU Lesser General Public License for more details.

 You should have received a copy of the GNU Lesser General Public License
 along with SortMeRNA. If not, see <http://www.gnu.org/licenses/>.
 @endparblock

 @contributors Jenya Kopylova   jenya.kopylov@gmail.com
			   Laurent No      laurent.noe@lifl.fr
			   Pierre Pericard  pierre.pericard@lifl.fr
			   Daniel McDonald  wasade@gmail.com
			   Mikal Salson    mikael.salson@lifl.fr
			   Hlne Touzet    helene.touzet@lifl.fr
			   Rob Knight       robknight@ucsd.edu
*/

/**
* file: ThreadPool.hpp
* created: 20170810 Thu
*
* https://stackoverflow.com/questions/26516683/reusing-thread-in-loop-c
*/

#pragma once

#include <iostream>
#include <sstream>
#include <thread>
#include <mutex>
#include <condition_variable>
#include <queue>
#include <functional>
#include <chrono>
#include <atomic>

#include "common.hpp"

/**
*  all the pool threads are initially in a waiting state until jobs are available for execution.
*/
class ThreadPool
{
protected:
	std::mutex job_queue_mx; // lock for pop/push on jobs_
	std::mutex job_done_mx; // lock for checking jobs_.empty and shutdown
	std::condition_variable cv_jobs;
	std::condition_variable cv_done;
	std::atomic_bool shutdown_;
	std::atomic_uint running_threads; // counter of running threads
	std::vector <std::thread> threads_;
	std::queue <std::function <void(void)>> jobs_;

public:
	ThreadPool(int numThreads) : shutdown_(false), running_threads(0)
	{
		// Create the specified number of threads
		threads_.reserve(numThreads);
		for (int i = 0; i < numThreads; ++i) {
			threads_.emplace_back(std::bind(&ThreadPool::threadEntry, this, i));
		}

		INFO("initialized Pool with [", numThreads, "] threads\n");
	}

	~ThreadPool()
	{
		{
			// Unblock any threads and tell them to stop
			std::lock_guard <std::mutex> lk(job_queue_mx);
			shutdown_ = true;
			cv_jobs.notify_all();
		}
		joinAll();
		PRN_MEM("ThreadPool destructor done.");
	}

protected:
		void threadEntry(int i)
		{
			std::function<void(void)> job;

			for (;;)
			{
				{
					std::unique_lock<std::mutex> lk(job_queue_mx);

					// Sleep until there is a job to execute or a shutdown flag set
					while (!shutdown_.load() && jobs_.empty())
						cv_jobs.wait(lk); // works
					//cv_jobs.wait(lk, [this] { return !shutdown_.load() && jobs_.empty(); });

					if (jobs_.empty()) // No jobs to do and shutting down
					{
						INFO("Thread  ", std::this_thread::get_id(), " job done");
						return;
					}

					job = std::move(jobs_.front());
					jobs_.pop();
					++running_threads;
				}
				// mutex 'job_queue_mx' released here

				job(); // Do the job without holding any locks
				--running_threads;

				INFO("number of running_threads= ", running_threads.load(), " jobs queue is empty= ", jobs_.empty());

				cv_done.notify_one(); // wake up the main thread waiting in 'waitAll'. Keep it here for multi-reference cases.
			} // ~for
		} // ~threadEntry

public:
	/** 
	 * lock the jobs queue and add a job
	 */
	void addJob(std::function <void()> func)
	{
		std::lock_guard<std::mutex> lk(job_queue_mx);
		jobs_.emplace(std::move(func));
		cv_jobs.notify_one();
	}

	// wait till no jobs running
	void waitAll()
	{
		std::unique_lock<std::mutex> lk(job_done_mx);
		cv_done.wait(lk, [this] { return running_threads.load() == 0 && jobs_.empty(); });
	}

	// Wait for all threads to stop
	void joinAll()
	{
		for (auto& thread : threads_)
			thread.join();
	}
}; // ~class ThreadPool