-
Notifications
You must be signed in to change notification settings - Fork 22
Expand file tree
/
Copy pathCThreadPool_Ret.hpp
More file actions
104 lines (102 loc) · 3.27 KB
/
Copy pathCThreadPool_Ret.hpp
File metadata and controls
104 lines (102 loc) · 3.27 KB
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
#ifndef CTHREADPOOL_RET
#define CTHREADPOOL_RET
#include<thread> //thread::hardware_concurrency
#include<type_traits> //invoke_result_t
#include<unordered_map>
#include<utility> //forward, move
#include"../../lib/header/thread/CWait_bounded_queue.hpp"
#include"CThreadPoolItem_Ret.hpp"
namespace nThread
{
//1. a fixed-sized threadpool
//2. can return value
template<class Ret>
class CThreadPool_Ret
{
public:
using size_type=std::invoke_result_t<decltype(std::thread::hardware_concurrency)>;
using thread_id=IThreadPoolItemBase::id;
private:
CWait_bounded_queue<CThreadPoolItem_Ret<Ret>*> waiting_queue_;
std::unordered_map<thread_id,CThreadPoolItem_Ret<Ret>> thr_;
public:
CThreadPool_Ret()
:CThreadPool_Ret{std::thread::hardware_concurrency()}{}
//1. determine the total of usable threads
//2. the value you pass will always equal to CThreadPool_Ret::size
explicit CThreadPool_Ret(size_type size)
:waiting_queue_{size},thr_{size}
{
while(size--)
{
CThreadPoolItem_Ret<Ret> item{waiting_queue_};
const auto id{item.get_id()};
waiting_queue_.emplace_not_ts(&thr_.emplace(id,std::move(item)).first->second);
}
}
//of course, why do you need to copy or move CThreadPool_Ret?
CThreadPool_Ret(const CThreadPool_Ret &)=delete;
//1. block until CThreadPool::empty is false and execute the func
//2. after returning from add, usable threads will reduce 1
//3. after returning from add, CThreadPool_Ret::valid(thread_id) will return true
//4. you must call CThreadPool_Ret::get after returning from add
template<class Func,class ... Args>
thread_id add(Func &&func,Args &&...args)
{
const auto temp{waiting_queue_.wait_and_pop()};
try
{
temp->assign(forward_as_lambda(std::forward<decltype(func)>(func),std::forward<decltype(args)>(args)...));
}catch(...)
{
waiting_queue_.emplace_and_notify(temp);
throw ;
}
return temp->get_id();
}
//1. return true if there does not have usable threads at that moment
//2. return false if there has usable threads at that moment
//3. non-block
bool empty() const noexcept
{
return waiting_queue_.empty();
}
//1. block until the thread_id completes the func
//2. after returning from get, CThreadPool_Ret::valid(thread_id) will return false
//3. if the thread_id is not valid, do not get the thread_id
inline Ret get(const thread_id id)
{
return thr_.at(id).get();
}
//1. return the total of usable threads
//2. is fixed after constructing
//3. non-block
inline size_type size() const noexcept
{
return static_cast<size_type>(thr_.size());
}
//1. return whether the thread_id has been get yet
//2. return true for the thread_id which was returned by CThreadPool_Ret::add
//3. return false for the thread_id which was used by CThreadPool_Ret::get
//4. non-block
inline bool valid(const thread_id id) const
{
return thr_.at(id).is_running();
}
//block until the thread_id completes the func
inline void wait(const thread_id id) const
{
thr_.at(id).wait();
}
void wait_all() const
{
for(const auto &val:thr_)
if(valid(val.first))
wait(val.first);
}
//of course, why do you need to copy or move CThreadPool_Ret?
CThreadPool_Ret& operator=(const CThreadPool_Ret &)=delete;
//get all the threads in destructor
};
}
#endif