Coverage Report

Created: 2026-07-16 06:32

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/work/openh264/codec/common/src/WelsThreadPool.cpp
Line
Count
Source
1
/*!
2
 * \copy
3
 *     Copyright (c)  2009-2015, Cisco Systems
4
 *     All rights reserved.
5
 *
6
 *     Redistribution and use in source and binary forms, with or without
7
 *     modification, are permitted provided that the following conditions
8
 *     are met:
9
 *
10
 *        * Redistributions of source code must retain the above copyright
11
 *          notice, this list of conditions and the following disclaimer.
12
 *
13
 *        * Redistributions in binary form must reproduce the above copyright
14
 *          notice, this list of conditions and the following disclaimer in
15
 *          the documentation and/or other materials provided with the
16
 *          distribution.
17
 *
18
 *     THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS
19
 *     "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT
20
 *     LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS
21
 *     FOR A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE
22
 *     COPYRIGHT HOLDER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT,
23
 *     INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING,
24
 *     BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES;
25
 *     LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER
26
 *     CAUSED AND ON ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT
27
 *     LIABILITY, OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN
28
 *     ANY WAY OUT OF THE USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE
29
 *     POSSIBILITY OF SUCH DAMAGE.
30
 *
31
 *
32
 * \file    WelsThreadPool.cpp
33
 *
34
 * \brief   functions for Thread Pool
35
 *
36
 * \date    5/09/2012 Created
37
 *
38
 *************************************************************************************
39
 */
40
#include "typedefs.h"
41
#include "memory_align.h"
42
#include "WelsThreadPool.h"
43
44
namespace WelsCommon {
45
46
namespace {
47
48
0
CWelsLock& GetInitLock() {
49
0
  static CWelsLock *initLock = new CWelsLock;
50
0
  return *initLock;
51
0
}
52
53
}
54
55
int32_t CWelsThreadPool::m_iRefCount = 0;
56
int32_t CWelsThreadPool::m_iMaxThreadNum = DEFAULT_THREAD_NUM;
57
CWelsThreadPool* CWelsThreadPool::m_pThreadPoolSelf = NULL;
58
59
CWelsThreadPool::CWelsThreadPool() :
60
0
  m_cWaitedTasks (NULL), m_cIdleThreads (NULL), m_cBusyThreads (NULL) {
61
0
}
62
63
64
0
CWelsThreadPool::~CWelsThreadPool() {
65
  //fprintf(stdout, "CWelsThreadPool::~CWelsThreadPool: delete %x, %x, %x\n", m_cWaitedTasks, m_cIdleThreads, m_cBusyThreads);
66
0
  if (0 != m_iRefCount) {
67
0
    m_iRefCount = 0;
68
0
    Uninit();
69
0
  }
70
0
}
71
72
0
WELS_THREAD_ERROR_CODE CWelsThreadPool::SetThreadNum (int32_t iMaxThreadNum) {
73
0
  CWelsAutoLock  cLock (GetInitLock());
74
75
0
  if (m_iRefCount != 0) {
76
0
    return WELS_THREAD_ERROR_GENERAL;
77
0
  }
78
79
0
  if (iMaxThreadNum <= 0) {
80
0
    iMaxThreadNum = 1;
81
0
  }
82
0
  m_iMaxThreadNum = iMaxThreadNum;
83
0
  return WELS_THREAD_ERROR_OK;
84
0
}
85
86
87
0
CWelsThreadPool* CWelsThreadPool::AddReference() {
88
0
  CWelsAutoLock  cLock (GetInitLock());
89
0
  if (m_pThreadPoolSelf == NULL) {
90
0
    m_pThreadPoolSelf = new CWelsThreadPool();
91
0
    if (!m_pThreadPoolSelf) {
92
0
      return NULL;
93
0
    }
94
0
  }
95
96
0
  if (m_iRefCount == 0) {
97
0
    if (WELS_THREAD_ERROR_OK != m_pThreadPoolSelf->Init()) {
98
0
      m_pThreadPoolSelf->Uninit();
99
0
      delete m_pThreadPoolSelf;
100
0
      m_pThreadPoolSelf = NULL;
101
0
      return NULL;
102
0
    }
103
0
  }
104
105
  ////fprintf(stdout, "m_iRefCount=%d, iMaxThreadNum=%d\n", m_iRefCount, m_iMaxThreadNum);
106
107
0
  ++ m_iRefCount;
108
  //fprintf(stdout, "m_iRefCount2=%d\n", m_iRefCount);
109
0
  return m_pThreadPoolSelf;
110
0
}
111
112
0
void CWelsThreadPool::RemoveInstance() {
113
0
  CWelsAutoLock  cLock (GetInitLock());
114
  //fprintf(stdout, "m_iRefCount=%d\n", m_iRefCount);
115
0
  -- m_iRefCount;
116
0
  if (0 == m_iRefCount) {
117
0
    StopAllRunning();
118
0
    Uninit();
119
0
    if (m_pThreadPoolSelf) {
120
0
      delete m_pThreadPoolSelf;
121
0
      m_pThreadPoolSelf = NULL;
122
0
    }
123
    //fprintf(stdout, "m_iRefCount=%d, IdleThreadNum=%d, BusyThreadNum=%d, WaitedTask=%d\n", m_iRefCount, GetIdleThreadNum(), GetBusyThreadNum(), GetWaitedTaskNum());
124
0
  }
125
0
}
126
127
128
0
bool CWelsThreadPool::IsReferenced() {
129
0
  CWelsAutoLock  cLock (GetInitLock());
130
0
  return (m_iRefCount > 0);
131
0
}
132
133
134
0
WELS_THREAD_ERROR_CODE CWelsThreadPool::OnTaskStart (CWelsTaskThread* pThread, IWelsTask* pTask) {
135
0
  AddThreadToBusyList (pThread);
136
  //fprintf(stdout, "CWelsThreadPool::AddThreadToBusyList: Task %x at Thread %x\n", pTask, pThread);
137
0
  return WELS_THREAD_ERROR_OK;
138
0
}
139
140
0
WELS_THREAD_ERROR_CODE CWelsThreadPool::OnTaskStop (CWelsTaskThread* pThread, IWelsTask* pTask) {
141
  //fprintf(stdout, "CWelsThreadPool::OnTaskStop 0: Task %x at Thread %x Finished\n", pTask, pThread);
142
143
0
  RemoveThreadFromBusyList (pThread);
144
0
  AddThreadToIdleQueue (pThread);
145
146
0
  if (pTask && pTask->GetSink()) {
147
    //fprintf(stdout, "CWelsThreadPool::OnTaskStop 1: Task %x at Thread %x Finished, m_pSink=%x\n", pTask, pThread, pTask->GetSink());
148
0
    pTask->GetSink()->OnTaskExecuted();
149
    ////fprintf(stdout, "CWelsThreadPool::OnTaskStop 1: Task %x at Thread %x Finished, m_pSink=%x\n", pTask, pThread, pTask->GetSink());
150
0
  }
151
  //if (m_pSink) {
152
  //  m_pSink->OnTaskExecuted (pTask);
153
  //}
154
  //fprintf(stdout, "CWelsThreadPool::OnTaskStop 2: Task %x at Thread %x Finished\n", pTask, pThread);
155
156
0
  SignalThread();
157
158
  //fprintf(stdout, "ThreadPool: Task %x at Thread %x Finished\n", pTask, pThread);
159
0
  return WELS_THREAD_ERROR_OK;
160
0
}
161
162
0
WELS_THREAD_ERROR_CODE CWelsThreadPool::Init() {
163
  //fprintf(stdout, "Enter WelsThreadPool Init\n");
164
165
0
  CWelsAutoLock  cLock (m_cLockPool);
166
167
0
  m_cWaitedTasks = new CWelsNonDuplicatedList<IWelsTask>();
168
0
  m_cIdleThreads = new CWelsNonDuplicatedList<CWelsTaskThread>();
169
0
  m_cBusyThreads = new CWelsList<CWelsTaskThread>();
170
0
  if (NULL == m_cWaitedTasks || NULL == m_cIdleThreads || NULL == m_cBusyThreads) {
171
0
    return WELS_THREAD_ERROR_GENERAL;
172
0
  }
173
174
0
  for (int32_t i = 0; i < m_iMaxThreadNum; i++) {
175
0
    if (WELS_THREAD_ERROR_OK != CreateIdleThread()) {
176
0
      return WELS_THREAD_ERROR_GENERAL;
177
0
    }
178
0
  }
179
180
0
  if (WELS_THREAD_ERROR_OK != Start()) {
181
0
    return WELS_THREAD_ERROR_GENERAL;
182
0
  }
183
184
0
  return WELS_THREAD_ERROR_OK;
185
0
}
186
187
0
WELS_THREAD_ERROR_CODE CWelsThreadPool::StopAllRunning() {
188
0
  WELS_THREAD_ERROR_CODE iReturn = WELS_THREAD_ERROR_OK;
189
190
0
  ClearWaitedTasks();
191
192
0
  while (GetBusyThreadNum() > 0) {
193
    //WELS_INFO_TRACE ("CWelsThreadPool::Uninit - Waiting all thread to exit");
194
0
    WelsSleep (10);
195
0
  }
196
197
0
  if (GetIdleThreadNum() != m_iMaxThreadNum) {
198
0
    iReturn = WELS_THREAD_ERROR_GENERAL;
199
0
  }
200
201
0
  return iReturn;
202
0
}
203
204
0
WELS_THREAD_ERROR_CODE CWelsThreadPool::Uninit() {
205
0
  WELS_THREAD_ERROR_CODE iReturn = WELS_THREAD_ERROR_OK;
206
0
  CWelsAutoLock  cLock (m_cLockPool);
207
208
0
  iReturn = StopAllRunning();
209
0
  assert (0 == GetBusyThreadNum());
210
211
0
  m_cLockIdleTasks.Lock();
212
0
  while (m_cIdleThreads && m_cIdleThreads->size() > 0) {
213
0
    DestroyThread (m_cIdleThreads->begin());
214
0
    m_cIdleThreads->pop_front();
215
0
  }
216
0
  m_cLockIdleTasks.Unlock();
217
218
0
  Kill();
219
220
0
  WELS_DELETE_OP (m_cWaitedTasks);
221
0
  WELS_DELETE_OP (m_cIdleThreads);
222
0
  WELS_DELETE_OP (m_cBusyThreads);
223
224
0
  return iReturn;
225
0
}
226
227
0
void CWelsThreadPool::ExecuteTask() {
228
  //fprintf(stdout, "ThreadPool: scheduled tasks: ExecuteTask\n");
229
0
  CWelsTaskThread* pThread = NULL;
230
0
  IWelsTask*    pTask = NULL;
231
0
  while (GetWaitedTaskNum() > 0) {
232
    //fprintf(stdout, "ThreadPool:  ExecuteTask: waiting task %d\n", GetWaitedTaskNum());
233
0
    pThread = GetIdleThread();
234
0
    if (pThread == NULL) {
235
      //fprintf(stdout, "ThreadPool:  ExecuteTask: no IdleThread\n");
236
237
0
      break;
238
0
    }
239
0
    pTask = GetWaitedTask();
240
    //fprintf(stdout, "ThreadPool:  ExecuteTask = %x at thread %x\n", pTask, pThread);
241
0
    if (pTask) {
242
0
      pThread->SetTask (pTask);
243
0
    } else {
244
0
      AddThreadToIdleQueue (pThread);
245
0
    }
246
0
  }
247
0
}
248
249
0
WELS_THREAD_ERROR_CODE CWelsThreadPool::QueueTask (IWelsTask* pTask) {
250
0
  CWelsAutoLock  cLock (m_cLockPool);
251
252
  //fprintf(stdout, "CWelsThreadPool::QueueTask: %d, pTask=%x\n", m_iRefCount, pTask);
253
0
  if (GetWaitedTaskNum() == 0) {
254
0
    CWelsTaskThread* pThread = GetIdleThread();
255
256
0
    if (pThread != NULL) {
257
      //fprintf(stdout, "ThreadPool:  ExecuteTask = %x at thread %x\n", pTask, pThread);
258
0
      pThread->SetTask (pTask);
259
260
0
      return WELS_THREAD_ERROR_OK;
261
0
    }
262
0
  }
263
  //fprintf(stdout, "ThreadPool:  AddTaskToWaitedList: %x\n", pTask);
264
0
  if (false == AddTaskToWaitedList (pTask)) {
265
0
    return WELS_THREAD_ERROR_GENERAL;
266
0
  }
267
268
  //fprintf(stdout, "ThreadPool:  SignalThread: %x\n", pTask);
269
0
  SignalThread();
270
0
  return WELS_THREAD_ERROR_OK;
271
0
}
272
273
0
WELS_THREAD_ERROR_CODE CWelsThreadPool::CreateIdleThread() {
274
0
  CWelsTaskThread* pThread = new CWelsTaskThread (this);
275
276
0
  if (NULL == pThread) {
277
0
    return WELS_THREAD_ERROR_GENERAL;
278
0
  }
279
280
0
  if (WELS_THREAD_ERROR_OK != pThread->Start()) {
281
0
    WELS_DELETE_OP (pThread);
282
0
    return WELS_THREAD_ERROR_GENERAL;
283
0
  }
284
  //fprintf(stdout, "ThreadPool:  AddThreadToIdleQueue: %x\n", pThread);
285
0
  AddThreadToIdleQueue (pThread);
286
287
0
  return WELS_THREAD_ERROR_OK;
288
0
}
289
290
0
void  CWelsThreadPool::DestroyThread (CWelsTaskThread* pThread) {
291
0
  pThread->Kill();
292
0
  WELS_DELETE_OP (pThread);
293
294
0
  return;
295
0
}
296
297
0
WELS_THREAD_ERROR_CODE CWelsThreadPool::AddThreadToIdleQueue (CWelsTaskThread* pThread) {
298
0
  CWelsAutoLock cLock (m_cLockIdleTasks);
299
0
  m_cIdleThreads->push_back (pThread);
300
0
  return WELS_THREAD_ERROR_OK;
301
0
}
302
303
0
WELS_THREAD_ERROR_CODE CWelsThreadPool::AddThreadToBusyList (CWelsTaskThread* pThread) {
304
0
  CWelsAutoLock cLock (m_cLockBusyTasks);
305
0
  m_cBusyThreads->push_back (pThread);
306
0
  return WELS_THREAD_ERROR_OK;
307
0
}
308
309
0
WELS_THREAD_ERROR_CODE CWelsThreadPool::RemoveThreadFromBusyList (CWelsTaskThread* pThread) {
310
0
  CWelsAutoLock cLock (m_cLockBusyTasks);
311
0
  if (m_cBusyThreads->erase (pThread)) {
312
0
    return WELS_THREAD_ERROR_OK;
313
0
  } else {
314
0
    return WELS_THREAD_ERROR_GENERAL;
315
0
  }
316
0
}
317
318
0
bool  CWelsThreadPool::AddTaskToWaitedList (IWelsTask* pTask) {
319
0
  CWelsAutoLock  cLock (m_cLockWaitedTasks);
320
321
0
  return m_cWaitedTasks->push_back (pTask);
322
0
}
323
324
0
CWelsTaskThread*   CWelsThreadPool::GetIdleThread() {
325
0
  CWelsAutoLock cLock (m_cLockIdleTasks);
326
327
0
  if (NULL == m_cIdleThreads || m_cIdleThreads->size() == 0) {
328
0
    return NULL;
329
0
  }
330
331
  //fprintf(stdout, "CWelsThreadPool::GetIdleThread=%d\n", m_cIdleThreads->size());
332
333
0
  CWelsTaskThread* pThread = m_cIdleThreads->begin();
334
0
  m_cIdleThreads->pop_front();
335
0
  return pThread;
336
0
}
337
338
0
int32_t  CWelsThreadPool::GetBusyThreadNum() {
339
0
  return (m_cBusyThreads?m_cBusyThreads->size():0);
340
0
}
341
342
0
int32_t  CWelsThreadPool::GetIdleThreadNum() {
343
0
  return (m_cIdleThreads?m_cIdleThreads->size():0);
344
0
}
345
346
0
int32_t  CWelsThreadPool::GetWaitedTaskNum() {
347
0
  return (m_cWaitedTasks?m_cWaitedTasks->size():0);
348
0
}
349
350
0
IWelsTask* CWelsThreadPool::GetWaitedTask() {
351
0
  CWelsAutoLock lock (m_cLockWaitedTasks);
352
353
0
  if (NULL==m_cWaitedTasks || m_cWaitedTasks->size() == 0) {
354
0
    return NULL;
355
0
  }
356
357
0
  IWelsTask* pTask = m_cWaitedTasks->begin();
358
359
0
  m_cWaitedTasks->pop_front();
360
361
0
  return pTask;
362
0
}
363
364
0
void  CWelsThreadPool::ClearWaitedTasks() {
365
0
  CWelsAutoLock cLock (m_cLockWaitedTasks);
366
0
  if (NULL == m_cWaitedTasks) {
367
0
    return;
368
0
  }
369
0
  IWelsTask* pTask = NULL;
370
0
  while (0 != m_cWaitedTasks->size()) {
371
0
    pTask = m_cWaitedTasks->begin();
372
0
    if (pTask->GetSink()) {
373
0
      pTask->GetSink()->OnTaskCancelled();
374
0
    }
375
0
    m_cWaitedTasks->pop_front();
376
0
  }
377
0
}
378
379
}