| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496 |
- /*
- * Copyright: JessMA Open Source (ldcsaa@gmail.com)
- *
- * Author : Bruce Liang
- * Website : https://github.com/ldcsaa
- * Project : https://github.com/ldcsaa/HP-Socket
- * Blog : http://www.cnblogs.com/ldcsaa
- * Wiki : http://www.oschina.net/p/hp-socket
- * QQ Group : 44636872, 75375912
- *
- * Licensed under the Apache License, Version 2.0 (the "License");
- * you may not use this file except in compliance with the License.
- * You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
- #include "HPThreadPool.h"
- #include "common/FuncHelper.h"
- #include <pthread.h>
- LPTSocketTask CreateSocketTaskObj( Fn_SocketTaskProc fnTaskProc,
- PVOID pSender, CONNID dwConnID,
- LPCBYTE pBuffer, INT iBuffLen, EnTaskBufferType enBuffType,
- WPARAM wParam, LPARAM lParam)
- {
- ASSERT(fnTaskProc != nullptr);
- ASSERT(iBuffLen >= 0);
- LPTSocketTask pTask = new TSocketTask;
- pTask->fn = fnTaskProc;
- pTask->sender = pSender;
- pTask->connID = dwConnID;
- pTask->bufLen = iBuffLen;
- pTask->bufType = enBuffType;
- pTask->wparam = wParam;
- pTask->lparam = lParam;
- if(enBuffType != TBT_COPY || !pBuffer)
- pTask->buf = pBuffer;
- else
- {
- pTask->buf = MALLOC(BYTE, iBuffLen);
- ::CopyMemory((LPBYTE)pTask->buf, pBuffer, iBuffLen);
- }
- return pTask;
- }
- void DestroySocketTaskObj(LPTSocketTask pTask)
- {
- if(pTask)
- {
- if(pTask->bufType != TBT_REFER && pTask->buf)
- FREE(pTask->buf);
- delete pTask;
- }
- }
- volatile UINT CHPThreadPool::sm_uiNum = 0;
- LPCTSTR CHPThreadPool::POOLED_THREAD_PREFIX = _T("hp-pool-");
- BOOL CHPThreadPool::Start(DWORD dwThreadCount, DWORD dwMaxQueueSize, EnRejectedPolicy enRejectedPolicy, DWORD dwStackSize)
- {
- if(!CheckStarting())
- return FALSE;
- m_dwStackSize = dwStackSize;
- m_dwMaxQueueSize = dwMaxQueueSize;
- m_enRejectedPolicy = enRejectedPolicy;
- FireStartup();
- if(!InternalAdjustThreadCount(dwThreadCount))
- {
- EXECUTE_RESTORE_ERROR(Stop());
- return FALSE;
- }
- m_enState = SS_STARTED;
- return TRUE;
- }
- BOOL CHPThreadPool::Stop(DWORD dwMaxWait)
- {
- if(!CheckStoping())
- return FALSE;
- ::WaitFor(15);
- Shutdown(dwMaxWait);
- FireShutdown();
- Reset();
- return TRUE;
- }
- BOOL CHPThreadPool::Shutdown(DWORD dwMaxWait)
- {
- BOOL isOK = TRUE;
- BOOL bLimited = (m_dwMaxQueueSize != 0);
- BOOL bInfinite = (dwMaxWait == (DWORD)INFINITE || dwMaxWait == 0);
- auto prdShutdown = [this]() {return m_stThreads.empty();};
- if(m_enRejectedPolicy == TRP_WAIT_FOR && bLimited)
- m_evQueue.SyncNotifyAll();
- VERIFY(DoAdjustThreadCount(0));
- if(bInfinite)
- m_evShutdown.Wait(prdShutdown);
- else
- m_evShutdown.WaitFor(dwMaxWait, prdShutdown);
- ASSERT(m_lsTasks.Size() == 0);
- ASSERT(m_stThreads.size() == 0);
- if(!m_lsTasks.IsEmpty())
- {
- TTask* pTask = nullptr;
- while(m_lsTasks.PopFront(&pTask))
- {
- if(pTask->freeArg)
- ::DestroySocketTaskObj((LPTSocketTask)pTask->arg);
- TTask::Destruct(pTask);
- }
- ::SetLastError(ERROR_CANCELLED);
- isOK = FALSE;
- }
- if(!m_stThreads.empty())
- {
- CCriSecLock lock(m_csThread);
- if(!m_stThreads.empty())
- {
- #if !defined(__ANDROID__)
- for(auto it = m_stThreads.begin(), end = m_stThreads.end(); it != end; ++it)
- pthread_cancel(*it);
- #endif
- m_stThreads.clear();
- ::SetLastError(ERROR_CANCELLED);
- isOK = FALSE;
- }
- }
- return isOK;
- }
- BOOL CHPThreadPool::Submit(Fn_TaskProc fnTaskProc, PVOID pvArg, DWORD dwMaxWait)
- {
- return DoSubmit(fnTaskProc, pvArg, FALSE, dwMaxWait);
- }
- BOOL CHPThreadPool::Submit(LPTSocketTask pTask, DWORD dwMaxWait)
- {
- return DoSubmit((Fn_TaskProc)pTask->fn, (PVOID)pTask, TRUE, dwMaxWait);
- }
- BOOL CHPThreadPool::DoSubmit(Fn_TaskProc fnTaskProc, PVOID pvArg, BOOL bFreeArg, DWORD dwMaxWait)
- {
- EnSubmitResult sr = DirectSubmit(fnTaskProc, pvArg, bFreeArg);
- if(sr != SUBMIT_FULL)
- return (sr == SUBMIT_OK);
- if(m_enRejectedPolicy == TRP_CALL_FAIL)
- {
- ::SetLastError(ERROR_DESTINATION_ELEMENT_FULL);
- return FALSE;
- }
- else if(m_enRejectedPolicy == TRP_WAIT_FOR)
- {
- return CycleWaitSubmit(fnTaskProc, pvArg, dwMaxWait, bFreeArg);
- }
- else if(m_enRejectedPolicy == TRP_CALLER_RUN)
- {
- DoRunTaskProc(fnTaskProc, pvArg, bFreeArg);
- }
- else
- {
- ASSERT(FALSE);
- ::SetLastError(ERROR_INVALID_PARAMETER);
- return FALSE;
- }
- return TRUE;
- }
- void CHPThreadPool::DoRunTaskProc(Fn_TaskProc fnTaskProc, PVOID pvArg, BOOL bFreeArg)
- {
- ::InterlockedIncrement(&m_dwTaskCount);
- fnTaskProc(pvArg);
- ::InterlockedDecrement(&m_dwTaskCount);
- if(bFreeArg)
- ::DestroySocketTaskObj((LPTSocketTask)pvArg);
- }
- CHPThreadPool::EnSubmitResult CHPThreadPool::DirectSubmit(Fn_TaskProc fnTaskProc, PVOID pvArg, BOOL bFreeArg)
- {
- if(!CheckStarted())
- return SUBMIT_ERROR;
- BOOL bLimited = (m_dwMaxQueueSize != 0);
- if(bLimited && m_lsTasks.Size() >= m_dwMaxQueueSize)
- return SUBMIT_FULL;
- else
- {
- TTask* pTask = TTask::Construct(fnTaskProc, pvArg, bFreeArg);
- m_lsTasks.PushBack(pTask);
- m_evTask.SyncNotifyOne();
- }
- return SUBMIT_OK;
- }
- BOOL CHPThreadPool::CycleWaitSubmit(Fn_TaskProc fnTaskProc, PVOID pvArg, DWORD dwMaxWait, BOOL bFreeArg)
- {
- ASSERT(m_dwMaxQueueSize != 0);
- DWORD dwTime = ::TimeGetTime();
- BOOL bInfinite = (dwMaxWait == (DWORD)INFINITE || dwMaxWait == 0);
- auto prdQueue = [this]() {return (m_lsTasks.Size() < m_dwMaxQueueSize);};
- while(CheckStarted())
- {
- EnSubmitResult sr = DirectSubmit(fnTaskProc, pvArg, bFreeArg);
- if(sr == SUBMIT_OK)
- return TRUE;
- if(sr == SUBMIT_ERROR)
- return FALSE;
- if(bInfinite)
- m_evQueue.Wait(prdQueue);
- else
- {
- DWORD dwNow = ::GetTimeGap32(dwTime);
- if(dwNow > dwMaxWait || !m_evQueue.WaitFor(chrono::milliseconds(dwMaxWait - dwNow), prdQueue))
- {
- ::SetLastError(ERROR_TIMEOUT);
- break;
- }
- }
- }
- return FALSE;
- }
- BOOL CHPThreadPool::AdjustThreadCount(DWORD dwNewThreadCount)
- {
- if(!CheckStarted())
- return FALSE;
- return InternalAdjustThreadCount(dwNewThreadCount);
- }
- BOOL CHPThreadPool::InternalAdjustThreadCount(DWORD dwNewThreadCount)
- {
- int iNewThreadCount = (int)dwNewThreadCount;
- if(iNewThreadCount == 0)
- iNewThreadCount = ::GetDefaultWorkerThreadCount();
- else if(iNewThreadCount < 0)
- iNewThreadCount = PROCESSOR_COUNT * (-iNewThreadCount);
- return DoAdjustThreadCount((DWORD)iNewThreadCount);
- }
- BOOL CHPThreadPool::DoAdjustThreadCount(DWORD dwNewThreadCount)
- {
- ASSERT((int)dwNewThreadCount >= 0);
- BOOL bRemove = FALSE;
- DWORD dwThreadCount = 0;
- if(dwNewThreadCount > m_dwThreadCount)
- {
- dwThreadCount = dwNewThreadCount - m_dwThreadCount;
- return CreateWorkerThreads(dwThreadCount);
- }
- else if(dwNewThreadCount < m_dwThreadCount)
- {
- bRemove = TRUE;
- dwThreadCount = m_dwThreadCount - dwNewThreadCount;
- ::InterlockedSub(&m_dwThreadCount, dwThreadCount);
- }
- if(bRemove)
- {
- for(DWORD i = 0; i < dwThreadCount; i++)
- m_evTask.SyncNotifyOne();
- }
- return TRUE;
- }
- BOOL CHPThreadPool::CreateWorkerThreads(DWORD dwThreadCount)
- {
- unique_ptr<pthread_attr_t> pThreadAttr;
- if(m_dwStackSize != 0)
- {
- pThreadAttr = make_unique<pthread_attr_t>();
- VERIFY_IS_NO_ERROR(pthread_attr_init(pThreadAttr.get()));
- int rs = pthread_attr_setstacksize(pThreadAttr.get(), m_dwStackSize);
- if(!IS_NO_ERROR(rs))
- {
- pthread_attr_destroy(pThreadAttr.get());
- ::SetLastError(rs);
- return FALSE;
- }
- }
- BOOL isOK = TRUE;
- for(DWORD i = 0; i < dwThreadCount; i++)
- {
- THR_ID dwThreadID;
- int rs = pthread_create(&dwThreadID, pThreadAttr.get(), ThreadProc, (PVOID)this);
- if(!IS_NO_ERROR(rs))
- {
- ::SetLastError(rs);
- isOK = FALSE;
- break;
- }
- ::InterlockedIncrement(&m_dwThreadCount);
- CCriSecLock lock(m_csThread);
- m_stThreads.emplace(dwThreadID);
- }
- if(pThreadAttr != nullptr)
- pthread_attr_destroy(pThreadAttr.get());
- return isOK;
- }
- PVOID CHPThreadPool::ThreadProc(LPVOID pv)
- {
- CHPThreadPool* pThis = (CHPThreadPool*)pv;
- ::SetSequenceThreadName(SELF_THREAD_ID, pThis->m_strPrefix, pThis->m_uiSeq);
- pThis->FireWorkerThreadStart();
-
- PVOID rs = (PVOID)(UINT_PTR)(pThis->WorkerProc());
-
- pThis->FireWorkerThreadEnd();
- return rs;
- }
- int CHPThreadPool::WorkerProc()
- {
- BOOL bLimited = (m_dwMaxQueueSize != 0);
- TTask* pTask = nullptr;
- auto prdTask = [this]() {return (!m_lsTasks.IsEmpty()) || (m_dwThreadCount < m_stThreads.size());};
- while(TRUE)
- {
- pTask = nullptr;
- while(m_lsTasks.PopFront(&pTask))
- {
- if(m_enRejectedPolicy == TRP_WAIT_FOR && bLimited)
- m_evQueue.SyncNotifyOne();
- DoRunTaskProc(pTask->fn, pTask->arg, pTask->freeArg);
- TTask::Destruct(pTask);
- }
-
- if(CheckWorkerThreadExit())
- break;
- m_evTask.Wait(prdTask);
- }
- return 0;
- }
- BOOL CHPThreadPool::CheckWorkerThreadExit()
- {
- BOOL bExit = FALSE;
- BOOL bShutdown = FALSE;
- if(m_dwThreadCount < m_stThreads.size())
- {
- CCriSecLock lock(m_csThread);
- if(m_dwThreadCount < m_stThreads.size())
- {
- VERIFY(m_stThreads.erase(SELF_THREAD_ID) == 1);
- bExit = TRUE;
- bShutdown = m_stThreads.empty();
- }
- }
- if(bExit)
- {
- pthread_detach(SELF_THREAD_ID);
- if(bShutdown)
- m_evShutdown.SyncNotifyOne();
- }
- return bExit;
- }
- BOOL CHPThreadPool::CheckStarting()
- {
- if(::InterlockedCompareExchange(&m_enState, SS_STARTING, SS_STOPPED) != SS_STOPPED)
- {
- ::SetLastError(ERROR_INVALID_STATE);
- return FALSE;
- }
- return TRUE;
- }
- BOOL CHPThreadPool::CheckStarted()
- {
- if(m_enState != SS_STARTED)
- {
- ::SetLastError(ERROR_INVALID_STATE);
- return FALSE;
- }
- return TRUE;
- }
- BOOL CHPThreadPool::CheckStoping()
- {
- if( ::InterlockedCompareExchange(&m_enState, SS_STOPPING, SS_STARTED) != SS_STARTED &&
- ::InterlockedCompareExchange(&m_enState, SS_STOPPING, SS_STARTING) != SS_STARTING)
- {
- while(m_enState != SS_STOPPED)
- ::WaitFor(5);
- ::SetLastError(ERROR_INVALID_STATE);
- return FALSE;
- }
- return TRUE;
- }
- void CHPThreadPool::Reset(BOOL bSetWaitEvent)
- {
- m_uiSeq = 0;
- m_dwStackSize = 0;
- m_dwTaskCount = 0;
- m_dwThreadCount = 0;
- m_dwMaxQueueSize = 0;
- m_enRejectedPolicy = TRP_CALL_FAIL;
- m_enState = SS_STOPPED;
- if(bSetWaitEvent)
- m_evWait.SyncNotifyAll();
- }
- void CHPThreadPool::MakePrefix()
- {
- UINT uiNumber = ::InterlockedIncrement(&sm_uiNum);
- m_strPrefix.Format(_T("%s%u-"), POOLED_THREAD_PREFIX, uiNumber);
- }
|