-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy pathThreadPool.h
More file actions
135 lines (132 loc) · 3.33 KB
/
Copy pathThreadPool.h
File metadata and controls
135 lines (132 loc) · 3.33 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
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
#pragma once
#include "Epoll.h"
#include "Thread.h"
#include "Function.h"
#include "Socket.h"
class CThreadPool
{
public:
CThreadPool() {
m_server = NULL;
timespec tp = { 0,0 };
clock_gettime(CLOCK_REALTIME, &tp);
char* buf = NULL;
asprintf(&buf, "%d.%d.sock", tp.tv_sec % 100000, tp.tv_nsec % 1000000);
if (buf != NULL) {
m_path = buf;
free(buf);
}//有问题的话,在start接口里面判断m_path来解决问题。
usleep(1);
}
~CThreadPool() {
Close();
}
CThreadPool(const CThreadPool&) = delete;
CThreadPool& operator=(const CThreadPool&) = delete;
public:
int Start(unsigned count) {
int ret = 0;
if (m_server != NULL)return -1;//已经初始化了
if (m_path.size() == 0)return -2;//构造函数失败!!!
m_server = new CSocket();
if (m_server == NULL)return -3;
ret = m_server->Init(CSockParam(m_path, SOCK_ISSERVER));
if (ret != 0)return -4;
ret = m_epoll.Create(count);
if (ret != 0)return -5;
ret = m_epoll.Add(*m_server, EpollData((void*)m_server));
if (ret != 0)return -6;
m_threads.resize(count);
for (unsigned i = 0; i < count; i++) {
m_threads[i] = new CThread(&CThreadPool::TaskDispatch, this);
if (m_threads[i] == NULL)return -7;
ret = m_threads[i]->Start();
if (ret != 0)return -8;
}
return 0;
}
void Close() {
m_epoll.Close();
if (m_server) {
CSocketBase* p = m_server;
m_server = NULL;
delete p;
}
for (auto thread : m_threads)
{
if (thread)delete thread;
}
m_threads.clear();
unlink(m_path);
}
template<typename _FUNCTION_, typename... _ARGS_>
int AddTask(_FUNCTION_ func, _ARGS_... args) {
static thread_local CSocket client;
int ret = 0;
if (client == -1) {
ret = client.Init(CSockParam(m_path, 0));
if (ret != 0)return -1;
ret = client.Link();
if (ret != 0)return -2;
}
CFunctionBase* base = new CFunction< _FUNCTION_, _ARGS_...>(func, args...);
if (base == NULL)return -3;
Buffer data(sizeof(base));
memcpy(data, &base, sizeof(base));
ret = client.Send(data);
if (ret != 0) {
delete base;
return -4;
}
return 0;
}
size_t Size()const { return m_threads.size(); }
private:
int TaskDispatch() {
while (m_epoll != -1) {
EPEvents events;
int ret = 0;
ssize_t esize = m_epoll.WaitEvents(events);
if (esize > 0) {
for (ssize_t i = 0; i < esize; i++) {
if (events[i].events & EPOLLIN) {
CSocketBase* pClient = NULL;
if (events[i].data.ptr == m_server) {//客户端请求连接
ret = m_server->Link(&pClient);
if (ret != 0)continue;
ret = m_epoll.Add(*pClient, EpollData((void*)pClient));
if (ret != 0) {
delete pClient;
continue;
}
}
else {//客户端的数据来了
pClient = (CSocketBase*)events[i].data.ptr;
if (pClient) {
CFunctionBase* base = NULL;
Buffer data(sizeof(base));
ret = pClient->Recv(data);
if (ret <= 0) {
m_epoll.Del(*pClient);
delete pClient;
continue;
}
memcpy(&base, (char*)data, sizeof(base));
if (base != NULL) {
(*base)();
delete base;
}
}
}
}
}
}
}
return 0;
}
private:
CEpoll m_epoll;
std::vector<CThread*> m_threads;
CSocketBase* m_server;
Buffer m_path;
};