Coverage Report

Created: 2026-07-30 06:27

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/work/vvdec/source/Lib/Utilities/ThreadPool.cpp
Line
Count
Source
1
/* -----------------------------------------------------------------------------
2
The copyright in this software is being made available under the Clear BSD
3
License, included below. No patent rights, trademark rights and/or 
4
other Intellectual Property Rights other than the copyrights concerning 
5
the Software are granted under this license.
6
7
The Clear BSD License
8
9
Copyright (c) 2018-2026, Fraunhofer-Gesellschaft zur Förderung der angewandten Forschung e.V. & The VVdeC Authors.
10
All rights reserved.
11
12
Redistribution and use in source and binary forms, with or without modification,
13
are permitted (subject to the limitations in the disclaimer below) provided that
14
the following conditions are met:
15
16
     * Redistributions of source code must retain the above copyright notice,
17
     this list of conditions and the following disclaimer.
18
19
     * Redistributions in binary form must reproduce the above copyright
20
     notice, this list of conditions and the following disclaimer in the
21
     documentation and/or other materials provided with the distribution.
22
23
     * Neither the name of the copyright holder nor the names of its
24
     contributors may be used to endorse or promote products derived from this
25
     software without specific prior written permission.
26
27
NO EXPRESS OR IMPLIED LICENSES TO ANY PARTY'S PATENT RIGHTS ARE GRANTED BY
28
THIS LICENSE. THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND
29
CONTRIBUTORS "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT
30
LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A
31
PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT HOLDER OR
32
CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL,
33
EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED TO,
34
PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA, OR PROFITS; OR
35
BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY THEORY OF LIABILITY, WHETHER
36
IN CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE)
37
ARISING IN ANY WAY OUT OF THE USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE
38
POSSIBILITY OF SUCH DAMAGE.
39
40
41
------------------------------------------------------------------------------------------- */
42
43
#include "ThreadPool.h"
44
45
#include <chrono>
46
47
#if __linux
48
# include <pthread.h>
49
#endif
50
51
52
namespace vvdec
53
{
54
using namespace std::chrono_literals;
55
56
std::mutex Barrier::s_exceptionLock{};
57
58
// block threads after busy-waiting this long
59
14
const static auto VVDEC_BUSY_WAIT_TIME_MIN = [] {
60
14
  const char* env = getenv( "VVDEC_BUSY_WAIT_TIME_MIN" );
61
14
  if( env )
62
0
    return std::chrono::microseconds( int( atof( env ) * 1000 ) );
63
14
  return std::chrono::microseconds( 1ms );
64
14
}();
65
66
// block last waiting thread after this time
67
14
const static auto VVDEC_BUSY_WAIT_TIME_MAX = [] {
68
14
  const char* env = getenv( "VVDEC_BUSY_WAIT_TIME_MAX" );
69
14
  if( env )
70
0
    return atoi( env ) * 1ms;
71
14
  return 5ms;
72
14
}();
73
74
75
struct ThreadPool::TaskException : public std::exception
76
{
77
  explicit TaskException( std::exception_ptr e, ThreadPool::Slot& task )
78
0
    : m_originalException( e )
79
0
    , m_task( task )
80
0
  {}
81
  std::exception_ptr m_originalException;
82
  ThreadPool::Slot&  m_task;
83
};
84
85
class ScopeIncDecCounter
86
{
87
  std::atomic_uint& m_cntr;
88
89
public:
90
21.4k
  explicit ScopeIncDecCounter( std::atomic_uint& counter ): m_cntr( counter ) { m_cntr.fetch_add( 1, std::memory_order_relaxed ); }
91
21.5k
  ~ScopeIncDecCounter()                                                       { m_cntr.fetch_sub( 1, std::memory_order_relaxed ); }
92
  CLASS_COPY_MOVE_DELETE( ScopeIncDecCounter )
93
};
94
95
// ---------------------------------------------------------------------------
96
// Thread Pool
97
// ---------------------------------------------------------------------------
98
99
ThreadPool::ThreadPool( int numThreads, const char* threadPoolName )
100
676
  : m_poolName( threadPoolName )
101
676
  , m_threads( numThreads < 0 ? std::thread::hardware_concurrency() : numThreads )
102
676
  , m_poolPause( m_threads.size() )
103
676
{
104
676
  int tid = 0;
105
676
  for( auto& t: m_threads )
106
21.6k
  {
107
21.6k
    t = std::thread( &ThreadPool::threadProc, this, tid++ );
108
21.6k
  }
109
676
}
110
111
ThreadPool::~ThreadPool()
112
676
{
113
676
  m_exitThreads = true;
114
115
676
  waitForThreads();
116
676
}
117
118
bool ThreadPool::processTasksOnMainThread()
119
0
{
120
0
  CHECK_FATAL( m_threads.size() != 0, "should not be used with multiple threads" );
121
122
0
  TaskIterator taskIt;
123
0
  while( true )
124
0
  {
125
0
    try
126
0
    {
127
0
      if( !taskIt.isValid() )
128
0
      {
129
0
        taskIt = findNextTask( 0, m_tasks.begin() );
130
0
      }
131
0
      else
132
0
      {
133
0
        taskIt = findNextTask( 0, taskIt );
134
0
      }
135
136
0
      if( !taskIt.isValid() )
137
0
      {
138
0
        break;
139
0
      }
140
0
      processTask( 0, *taskIt );
141
0
    }
142
0
    catch( TaskException& e )
143
0
    {
144
0
      handleTaskException( e.m_originalException, e.m_task.done, e.m_task.counter, &e.m_task.state );
145
0
    }
146
0
  }
147
148
  // return true if all done (-> false if some tasks blocked due to barriers)
149
0
  return std::all_of( m_tasks.begin(), m_tasks.end(), []( Slot& t ) { return t.state == FREE; } );
150
0
}
151
152
void ThreadPool::shutdown( bool block )
153
676
{
154
676
  m_exitThreads = true;
155
676
  if( block )
156
676
  {
157
676
    waitForThreads();
158
676
  }
159
676
}
160
161
void ThreadPool::waitForThreads()
162
1.35k
{
163
1.35k
  m_poolPause.unpauseIfPaused( m_poolPause.acquireLock() );
164
165
1.35k
  for( auto& t: m_threads )
166
43.2k
  {
167
43.2k
    if( t.joinable() )
168
21.6k
      t.join();
169
43.2k
  }
170
1.35k
}
171
172
void ThreadPool::checkAndThrowThreadPoolException()
173
0
{
174
0
  if( !m_exceptionFlag.load() )
175
0
  {
176
0
    return;
177
0
  }
178
179
0
  msg( WARNING, "ThreadPool is in exception state." );
180
181
0
  std::exception_ptr tmp = m_threadPoolException;
182
0
  m_threadPoolException  = tmp;
183
0
  m_exceptionFlag.store( false );
184
185
0
  std::rethrow_exception( tmp );
186
0
}
187
188
void ThreadPool::threadProc( int threadId )
189
21.6k
{
190
21.6k
#if __linux
191
21.6k
  if( !m_poolName.empty() )
192
21.6k
  {
193
21.6k
    std::string threadName( m_poolName + std::to_string( threadId ) );
194
21.6k
    pthread_setname_np( pthread_self(), threadName.c_str() );
195
21.6k
  }
196
21.6k
#endif
197
198
21.6k
  auto nextTaskIt = m_tasks.begin();
199
21.6k
  while( !m_exitThreads )
200
21.6k
  {
201
21.6k
    try
202
21.6k
    {
203
21.6k
      auto taskIt = findNextTask( threadId, nextTaskIt );
204
21.6k
      if( !taskIt.isValid() )   // immediately try again without any delay
205
21.1k
      {
206
21.1k
        taskIt = findNextTask( threadId, nextTaskIt );
207
21.1k
      }
208
21.6k
      if( !taskIt.isValid() )   // still nothing found, go into idle loop searching for more tasks
209
21.0k
      {
210
21.0k
        ITT_TASKSTART( itt_domain_thrd, itt_handle_TPspinWait );
211
212
21.0k
        std::unique_lock<std::mutex> idleLock( m_idleMutex, std::defer_lock );
213
214
21.0k
        const auto startWait = std::chrono::steady_clock::now();
215
21.0k
        bool       didBlock  = false;   // if the previous iteration did block we don't want to yield in this iteration
216
285k
        while( !m_exitThreads )
217
264k
        {
218
264k
          if( !didBlock )
219
263k
          {
220
263k
            std::this_thread::yield();
221
263k
          }
222
264k
          didBlock = false;
223
224
264k
          taskIt = findNextTask( threadId, nextTaskIt );
225
264k
          if( taskIt.isValid() || m_exitThreads.load( std::memory_order_relaxed ) )
226
45
          {
227
45
            break;
228
45
          }
229
230
263k
          if( !idleLock.owns_lock() )
231
153k
          {
232
154k
            if( VVDEC_BUSY_WAIT_TIME_MIN.count() == 0 || std::chrono::steady_clock::now() - startWait > VVDEC_BUSY_WAIT_TIME_MIN )
233
21.4k
            {
234
21.4k
              ITT_TASKSTART( itt_domain_thrd, itt_handle_TPblocked );
235
21.4k
              ScopeIncDecCounter cntr( m_poolPause.m_waitingForLockThreads );
236
21.4k
              idleLock.lock();
237
21.4k
              didBlock = true;
238
21.4k
              ITT_TASKEND( itt_domain_thrd, itt_handle_TPblocked );
239
21.4k
            }
240
153k
          }
241
109k
          else if( std::chrono::steady_clock::now() - startWait > VVDEC_BUSY_WAIT_TIME_MAX )
242
28.5k
          {
243
#if THREAD_POOL_TASK_NAMES
244
            printWaitingTasks();
245
#endif
246
28.5k
            didBlock = m_poolPause.pauseIfAllOtherThreadsWaiting(
247
28.5k
              [&]
248
28.5k
              {
249
649
                taskIt = findNextTask( threadId, nextTaskIt );
250
649
                return taskIt.isValid() || m_exitThreads;
251
649
              } );
252
28.5k
            if( taskIt.isValid() )
253
0
            {
254
0
              break;
255
0
            }
256
28.5k
          }
257
263k
        }
258
259
21.0k
        ITT_TASKEND( itt_domain_thrd, itt_handle_TPspinWait );
260
21.0k
      }
261
21.6k
      if( m_exitThreads )
262
21.6k
      {
263
21.6k
        return;
264
21.6k
      }
265
266
18.4E
      processTask( threadId, *taskIt );
267
268
18.4E
      nextTaskIt = taskIt;
269
18.4E
      nextTaskIt.incWrap();
270
18.4E
    }
271
21.6k
    catch( TaskException& e )
272
21.6k
    {
273
0
      handleTaskException( e.m_originalException, e.m_task.done, e.m_task.counter, &e.m_task.state );
274
0
    }
275
21.6k
    catch( std::exception& e )
276
21.6k
    {
277
0
      msg( ERROR, "ERROR: Caught unexpected exception from within the thread pool: %s", e.what() );
278
279
0
      if( m_exceptionFlag.exchange( true ) )
280
0
      {
281
0
        msg( ERROR, "ERROR: Another exception has already happend in the thread pool, but we can only store one." );
282
0
        return;
283
0
      }
284
0
      m_threadPoolException = std::current_exception();
285
0
      return;
286
0
    }
287
21.6k
  }
288
21.6k
}
289
290
bool ThreadPool::checkTaskReady( int threadId, CBarrierVec& barriers, ThreadPool::TaskFunc readyCheck, void* taskParam )
291
0
{
292
0
  if( !barriers.empty() )
293
0
  {
294
    // don't break early, because isBlocked() also checks exception state
295
0
    if( std::count_if( barriers.cbegin(), barriers.cend(), []( const Barrier* b ) { return b && b->isBlocked(); } ) )
296
0
    {
297
0
      return false;
298
0
    }
299
0
    barriers.clear();
300
0
  }
301
302
0
  if( readyCheck && readyCheck( threadId, taskParam ) == false )
303
0
  {
304
0
    return false;
305
0
  }
306
307
0
  return true;
308
0
}
309
310
ThreadPool::TaskIterator ThreadPool::findNextTask( int threadId, TaskIterator startSearch )
311
307k
{
312
307k
  if( !startSearch.isValid() )
313
0
  {
314
0
    startSearch = m_tasks.begin();
315
0
  }
316
307k
  bool first = true;
317
16.7M
  for( auto it = startSearch; it != startSearch || first; it.incWrap() )
318
16.4M
  {
319
16.4M
    first = false;
320
16.4M
    try
321
16.4M
    {
322
16.4M
      Slot& task     = *it;
323
16.4M
      auto  expected = WAITING;
324
16.4M
      if( task.state == expected && task.state.compare_exchange_strong( expected, RUNNING ) )
325
0
      {
326
0
        if( checkTaskReady( threadId, task.barriers, task.readyCheck, task.param ) )
327
0
        {
328
0
          return it;
329
0
        }
330
331
        // reschedule
332
0
        task.state = WAITING;
333
0
      }
334
16.4M
    }
335
16.4M
    catch( ... )
336
16.4M
    {
337
0
      throw TaskException( std::current_exception(), *it );
338
0
    }
339
16.4M
  }
340
300k
  return {};
341
307k
}
342
343
bool ThreadPool::processTask( int threadId, ThreadPool::Slot& task )
344
0
{
345
0
  try
346
0
  {
347
0
    const bool success = task.func( threadId, task.param );
348
0
    if( !success )
349
0
    {
350
0
      task.state = WAITING;
351
0
      return false;
352
0
    }
353
354
0
    if( task.done != nullptr )
355
0
    {
356
0
      task.done->unlock();
357
0
    }
358
0
    if( task.counter != nullptr )
359
0
    {
360
0
      task.counter->decrement_nothrow();
361
0
    }
362
0
  }
363
0
  catch( ... )
364
0
  {
365
0
    throw TaskException( std::current_exception(), task );
366
0
  }
367
368
0
  task.state = FREE;
369
370
0
  return true;
371
0
}
372
373
bool ThreadPool::bypassTaskQueue( TaskFunc func, void* param, WaitCounter* counter, Barrier* done, CBarrierVec& barriers, TaskFunc readyCheck )
374
0
{
375
0
  CHECKD( numThreads() > 0, "the task queue should only be bypassed, when running single-threaded." );
376
0
  try
377
0
  {
378
    // if singlethreaded, execute all pending tasks
379
0
    bool waiting_tasks = m_nextFillSlot != m_tasks.begin();
380
0
    bool is_ready      = checkTaskReady( 0, barriers, (TaskFunc)readyCheck, param );
381
0
    if( !is_ready && waiting_tasks )
382
0
    {
383
0
      waiting_tasks = processTasksOnMainThread();
384
0
      is_ready      = checkTaskReady( 0, barriers, (TaskFunc)readyCheck, param );
385
0
    }
386
387
    // when no barriers block this task, execute it directly
388
0
    if( is_ready )
389
0
    {
390
0
      if( func( 0, param ) )
391
0
      {
392
0
        if( done != nullptr )
393
0
        {
394
0
          done->unlock();
395
0
        }
396
397
0
        if( waiting_tasks )
398
0
        {
399
0
          processTasksOnMainThread();
400
0
        }
401
0
        return true;
402
0
      }
403
0
    }
404
0
  }
405
0
  catch( ... )
406
0
  {
407
0
    handleTaskException( std::current_exception(), done, counter, nullptr );
408
0
  }
409
410
  // direct execution of the task failed
411
0
  return false;
412
0
}
413
414
void ThreadPool::handleTaskException( const std::exception_ptr e, Barrier* done, WaitCounter* counter, std::atomic<TaskState>* slot_state )
415
0
{
416
0
  if( done != nullptr )
417
0
  {
418
0
    done->setException( e );
419
0
  }
420
0
  if( counter != nullptr )
421
0
  {
422
0
    counter->setException( e );
423
0
    counter->decrement_nothrow();
424
0
  }
425
426
0
  if( slot_state != nullptr )
427
0
  {
428
0
    *slot_state = FREE;
429
0
  }
430
0
}
431
432
#if THREAD_POOL_TASK_NAMES
433
void ThreadPool::printWaitingTasks()
434
{
435
  std::cerr << "Waiting tasks:" << std::endl;
436
  int count = 0;
437
  for( auto& t: m_tasks )
438
  {
439
    if( t.state == WAITING )
440
    {
441
      ++count;
442
      std::cerr << t.taskName << std::endl;
443
    }
444
  }
445
  std::cerr << std::endl << count << " total tasks waiting" << std::endl;
446
}
447
#endif  // THREAD_POOL_TASK_NAMES
448
449
// ---------------------------------------------------------------------------
450
// Chunked Task Queue
451
// ---------------------------------------------------------------------------
452
453
ThreadPool::ChunkedTaskQueue::~ChunkedTaskQueue()
454
676
{
455
676
  Chunk* next = m_firstChunk.m_next;
456
676
  while( next )
457
0
  {
458
0
    Chunk* curr = next;
459
0
    next = curr->m_next;
460
0
    delete curr;
461
0
  }
462
676
}
463
464
ThreadPool::ChunkedTaskQueue::Iterator ThreadPool::ChunkedTaskQueue::grow()
465
0
{
466
0
  std::lock_guard<std::mutex> l( m_resizeMutex );   // prevent concurrent growth of the queue. Read access while growing is no problem
467
468
0
  m_lastChunk->m_next = new Chunk( &m_firstChunk );
469
0
  m_lastChunk         = m_lastChunk->m_next;
470
471
0
  return Iterator{ &m_lastChunk->m_slots.front(), m_lastChunk };
472
0
}
473
474
ThreadPool::ChunkedTaskQueue::Iterator& ThreadPool::ChunkedTaskQueue::Iterator::operator++()
475
0
{
476
0
  CHECKD( m_slot == nullptr, "incrementing invalid iterator" );
477
0
  CHECKD( m_chunk == nullptr, "incrementing invalid iterator" );
478
479
0
  if( m_slot != &m_chunk->m_slots.back() )
480
0
  {
481
0
    ++m_slot;
482
0
  }
483
0
  else
484
0
  {
485
0
    m_chunk = m_chunk->m_next;
486
0
    if( m_chunk )
487
0
    {
488
0
      m_slot = &m_chunk->m_slots.front();
489
0
    }
490
0
    else
491
0
    {
492
0
      m_slot = nullptr;
493
0
    }
494
0
  }
495
0
  return *this;
496
0
}
497
498
ThreadPool::ChunkedTaskQueue::Iterator& ThreadPool::ChunkedTaskQueue::Iterator::incWrap()
499
16.4M
{
500
16.4M
  CHECKD( m_slot == nullptr, "incrementing invalid iterator" );
501
16.4M
  CHECKD( m_chunk == nullptr, "incrementing invalid iterator" );
502
503
16.4M
  if( m_slot != &m_chunk->m_slots.back() )
504
16.3M
  {
505
16.3M
    ++m_slot;
506
16.3M
  }
507
96.3k
  else
508
96.3k
  {
509
96.3k
    if( m_chunk->m_next )
510
0
    {
511
0
      m_chunk = m_chunk->m_next;
512
0
    }
513
96.3k
    else
514
96.3k
    {
515
96.3k
      m_chunk = &m_chunk->m_firstChunk;
516
96.3k
    }
517
96.3k
    m_slot = &m_chunk->m_slots.front();
518
96.3k
  }
519
16.4M
  return *this;
520
16.4M
}
521
522
void ThreadPool::PoolPause::unpauseIfPaused( std::unique_lock<std::mutex> lockOwnership )
523
2.02k
{
524
2.02k
  CHECKD( lockOwnership.mutex() != &m_allThreadsWaitingMutex, "wrong mutex passed into ThreadPool::PoolPause::unpauseIfPaused()" );
525
2.02k
  CHECKD( !lockOwnership.owns_lock(), "lock passed into ThreadPool::PoolPause::unpauseIfPaused() does not own lock" );
526
  // All threads may be sleeping. If so, wake up.
527
2.02k
  m_allThreadsWaiting = false;
528
2.02k
  m_allThreadsWaitingCV.notify_all();
529
2.02k
}
530
531
template<typename Predicate>
532
bool ThreadPool::PoolPause::pauseIfAllOtherThreadsWaiting( Predicate predicate )
533
28.5k
{
534
28.5k
  if( m_nrThreads == 0 )
535
0
  {
536
0
    return false;
537
0
  }
538
28.5k
  const auto nrWaiting = m_waitingForLockThreads.load( std::memory_order_relaxed );
539
28.5k
  if( nrWaiting < m_nrThreads - 1 )
540
27.9k
  {
541
27.9k
    return false;
542
27.9k
  }
543
544
  // All threads are waiting. This (current) threads is the one which locked `idleLock`. All
545
  // other threads are waiting in the above condition for `idleLock.lock();`.
546
  // The only way how more work for the threads can come in is if addBarrierTask is called
547
  // or if the thread pool is closed or destroyed.
548
649
  std::unique_lock<std::mutex> lock( m_allThreadsWaitingMutex );
549
649
  m_allThreadsWaiting = true;
550
1.29k
  m_allThreadsWaitingCV.wait( lock, [this, &predicate] { return !m_allThreadsWaiting || predicate(); } );
551
649
  return true;
552
28.5k
}
553
554
}   // namespace vvdec