|
6 | 6 | #include <spdlog/spdlog.h> |
7 | 7 |
|
8 | 8 | SendThreadMain::SendThreadMain(CIOCPort* iocPort) |
9 | | - : _iocPort(iocPort), _running(false) |
| 9 | + : _iocPort(iocPort), _aiSocketCount(0) |
10 | 10 | { |
11 | 11 | } |
12 | 12 |
|
13 | | -void SendThreadMain::start() |
| 13 | +bool SendThreadMain::shutdown() |
14 | 14 | { |
15 | | - if (_running) |
16 | | - return; |
| 15 | + if (!Thread::shutdown()) |
| 16 | + return false; |
17 | 17 |
|
18 | | - _running = true; |
19 | | - _workerThread = std::thread(&SendThreadMain::run, this); |
20 | | -} |
21 | | - |
22 | | -void SendThreadMain::shutdown() |
23 | | -{ |
24 | | - if (!_running) |
25 | | - return; |
26 | | - |
27 | | - { |
28 | | - std::lock_guard<std::mutex> lock(_queueMutex); |
29 | | - _running = false; |
30 | | - _cv.notify_one(); |
31 | | - } |
32 | | - |
33 | | - _workerThread.join(); |
| 18 | + clear(); |
| 19 | + return true; |
34 | 20 | } |
35 | 21 |
|
36 | 22 | void SendThreadMain::queue(_SEND_DATA* sendData) |
37 | 23 | { |
38 | | - std::lock_guard<std::mutex> lock(_queueMutex); |
| 24 | + std::lock_guard<std::mutex> lock(_mutex); |
39 | 25 | _sendDataQueue.push(sendData); |
40 | 26 | _cv.notify_one(); |
41 | 27 | } |
42 | 28 |
|
43 | | -SendThreadMain::~SendThreadMain() |
44 | | -{ |
45 | | - shutdown(); |
46 | | - clear(); |
47 | | -} |
48 | | - |
49 | | -void SendThreadMain::run() |
| 29 | +void SendThreadMain::thread_loop() |
50 | 30 | { |
51 | 31 | while (_running) |
52 | 32 | { |
53 | | - std::unique_lock<std::mutex> lock(_queueMutex); |
| 33 | + std::unique_lock<std::mutex> lock(_mutex); |
54 | 34 | _cv.wait(lock); |
55 | 35 |
|
56 | 36 | if (!_running) |
57 | 37 | break; |
58 | 38 |
|
59 | | - int nRet = 0; |
60 | | - CGameSocket* pSocket = nullptr; |
61 | | - int size = 0, index = 0; |
| 39 | + tick(); |
| 40 | + } |
| 41 | +} |
62 | 42 |
|
63 | | - while (!_sendDataQueue.empty()) |
| 43 | +void SendThreadMain::tick() |
| 44 | +{ |
| 45 | + while (!_sendDataQueue.empty()) |
| 46 | + { |
| 47 | + _SEND_DATA* pSendData = _sendDataQueue.front(); |
| 48 | + |
| 49 | + int count = -1; |
| 50 | + for (int i = 0; i < MAX_SOCKET; i++) |
64 | 51 | { |
65 | | - _SEND_DATA* pSendData = _sendDataQueue.front(); |
| 52 | + CGameSocket* pSocket = (CGameSocket*) _iocPort->m_SockArray[i]; |
| 53 | + if (pSocket == nullptr) |
| 54 | + continue; |
66 | 55 |
|
67 | | - int count = -1; |
68 | | - for (int i = 0; i < MAX_SOCKET; i++) |
| 56 | + count++; |
| 57 | + |
| 58 | + if (_aiSocketCount == count) |
69 | 59 | { |
70 | | - CGameSocket* pSocket = (CGameSocket*) _iocPort->m_SockArray[i]; |
71 | | - if (pSocket == nullptr) |
| 60 | + int size = pSocket->Send(pSendData->pBuf, pSendData->sLength); |
| 61 | + if (size <= 0) |
| 62 | + { |
| 63 | + spdlog::error("SendThreadMain::tick: send failed: size={} socket_num={}", |
| 64 | + size, count); |
| 65 | + count--; |
72 | 66 | continue; |
| 67 | + } |
73 | 68 |
|
74 | | - count++; |
| 69 | + if (++_aiSocketCount >= MAX_AI_SOCKET) |
| 70 | + _aiSocketCount = 0; |
75 | 71 |
|
76 | | - if (_aiSocketCount == count) |
77 | | - { |
78 | | - size = pSocket->Send(pSendData->pBuf, pSendData->sLength); |
79 | | - if (size <= 0) |
80 | | - { |
81 | | - spdlog::error("SendThreadMain::run: send failed: size={} socket_num={}", |
82 | | - size, count); |
83 | | - count--; |
84 | | - continue; |
85 | | - } |
86 | | - |
87 | | - if (++_aiSocketCount >= MAX_AI_SOCKET) |
88 | | - _aiSocketCount = 0; |
89 | | - |
90 | | - //TRACE(_T("SendThreadMain - Send : size=%d, socket_num=%d\n"), size, count); |
91 | | - break; |
92 | | - } |
| 72 | + //TRACE(_T("SendThreadMain - Send : size=%d, socket_num=%d\n"), size, count); |
| 73 | + break; |
93 | 74 | } |
94 | | - |
95 | | - delete pSendData; |
96 | | - _sendDataQueue.pop(); |
97 | 75 | } |
| 76 | + |
| 77 | + delete pSendData; |
| 78 | + _sendDataQueue.pop(); |
98 | 79 | } |
99 | 80 | } |
100 | 81 |
|
101 | 82 | void SendThreadMain::clear() |
102 | 83 | { |
103 | | - std::lock_guard<std::mutex> lock(_queueMutex); |
| 84 | + std::lock_guard<std::mutex> lock(_mutex); |
104 | 85 | while (!_sendDataQueue.empty()); |
105 | 86 | { |
106 | 87 | delete _sendDataQueue.front(); |
107 | 88 | _sendDataQueue.pop(); |
108 | 89 | } |
109 | 90 | } |
| 91 | + |
| 92 | +SendThreadMain::~SendThreadMain() |
| 93 | +{ |
| 94 | + shutdown(); |
| 95 | + clear(); |
| 96 | +} |
0 commit comments