GitHub Viewer
// This file is part of OpenCV project.
// It is subject to the license terms in the LICENSE file found in the top-level directory
// of this distribution and at http://opencv.org/license.html.
#include "precomp.hpp"
#include "parallel_impl.hpp"
#ifdef HAVE_PTHREADS_PF
#include
#include
#include
//#undef CV_LOG_STRIP_LEVEL
//#define CV_LOG_STRIP_LEVEL CV_LOG_LEVEL_VERBOSE + 1
#include
//#define CV_PROFILE_THREADS 64
//#define getTickCount getCPUTickCount // use this if getTickCount() calls are expensive (and getCPUTickCount() is accurate)
//#define CV_USE_GLOBAL_WORKERS_COND_VAR // not effective on many-core systems (10+)
#include
// Spin lock's OS-level yield
#ifdef DECLARE_CV_YIELD
DECLARE_CV_YIELD
#endif
#ifndef CV_YIELD
# include
# define CV_YIELD() std::this_thread::yield()
#endif // CV_YIELD
// Spin lock's CPU-level yield (required for Hyper-Threading)
#ifdef DECLARE_CV_PAUSE
DECLARE_CV_PAUSE
#endif
#ifndef CV_PAUSE
# if defined __GNUC__ && (defined __i386__ || defined __x86_64__)
# if !defined(__SSE__)
static inline void cv_non_sse_mm_pause() { __asm__ __volatile__ ("rep; nop"); }
# define _mm_pause cv_non_sse_mm_pause
# endif
# define CV_PAUSE(v) do { for (int __delay = (v); __delay > 0; --__delay) { _mm_pause(); } } while (0)
# elif defined __GNUC__ && defined __aarch64__
# define CV_PAUSE(v) do { for (int __delay = (v); __delay > 0; --__delay) { asm volatile("yield" ::: "memory"); } } while (0)
# elif defined __GNUC__ && defined __arm__
# define CV_PAUSE(v) do { for (int __delay = (v); __delay > 0; --__delay) { asm volatile("" ::: "memory"); } } while (0)
# elif defined __GNUC__ && defined __PPC64__
# define CV_PAUSE(v) do { for (int __delay = (v); __delay > 0; --__delay) { asm volatile("or 27,27,27" ::: "memory"); } } while (0)
# else
# warning "Can't detect 'pause' (CPU-yield) instruction on the target platform. Specify CV_PAUSE() definition via compiler flags."
# define CV_PAUSE(...) do { /* no-op: works, but not effective */ } while (0)
# endif
#endif // CV_PAUSE
namespace cv
{
static int CV_ACTIVE_WAIT_PAUSE_LIMIT = (int)utils::getConfigurationParameterSizeT("OPENCV_THREAD_POOL_ACTIVE_WAIT_PAUSE_LIMIT", 16); // iterations
static int CV_WORKER_ACTIVE_WAIT = (int)utils::getConfigurationParameterSizeT("OPENCV_THREAD_POOL_ACTIVE_WAIT_WORKER", 2000); // iterations
static int CV_MAIN_THREAD_ACTIVE_WAIT = (int)utils::getConfigurationParameterSizeT("OPENCV_THREAD_POOL_ACTIVE_WAIT_MAIN", 10000); // iterations
static int CV_WORKER_ACTIVE_WAIT_THREADS_LIMIT = (int)utils::getConfigurationParameterSizeT("OPENCV_THREAD_POOL_ACTIVE_WAIT_THREADS_LIMIT", 0); // number of real cores
class WorkerThread;
class ParallelJob;
class ThreadPool
{
public:
static ThreadPool& instance()
{
CV_SINGLETON_LAZY_INIT_REF(ThreadPool, new ThreadPool())
}
static void stop()
{
ThreadPool& manager = instance();
manager.reconfigure(0);
}
void reconfigure(unsigned new_threads_count)
{
if (new_threads_count == threads.size())
return;
pthread_mutex_lock(&mutex);
reconfigure_(new_threads_count);
pthread_mutex_unlock(&mutex);
}
bool reconfigure_(unsigned new_threads_count); // internal implementation
void run(const Range& range, const ParallelLoopBody& body, double nstripes);
size_t getNumOfThreads();
void setNumOfThreads(unsigned n);
ThreadPool();
~ThreadPool();
unsigned num_threads;
pthread_mutex_t mutex; // guards fields (job/threads) from non-worker threads (concurrent parallel_for calls)
#if defined(CV_USE_GLOBAL_WORKERS_COND_VAR)
pthread_cond_t cond_thread_wake;
#endif
pthread_mutex_t mutex_notify;
pthread_cond_t cond_thread_task_complete;
std::vector< Ptr > threads;
Ptr job;
#ifdef CV_PROFILE_THREADS
double tickFreq;
int64 jobSubmitTime;
struct ThreadStatistics
{
ThreadStatistics() : threadWait(0)
{
reset();
}
void reset()
{
threadWake = 0;
threadExecuteStart = 0;
threadExecuteStop = 0;
executedTasks = 0;
keepActive = false;
threadPing = getTickCount();
}
int64 threadWait; // don't reset by default
int64 threadPing; // don't reset by default
int64 threadWake;
int64 threadExecuteStart;
int64 threadExecuteStop;
int64 threadFree;
unsigned executedTasks;
bool keepActive;
int64 dummy_[8]; // separate cache lines
void dump(int id, int64 baseTime, double tickFreq)
{
if (id < 0)
std::cout 0 ? (threadWait - baseTime) / tickFreq * 1e6 : -0.0,
threadPing > 0 ? (threadPing - baseTime) / tickFreq * 1e6 : -0.0);
if (threadWake > 0)
printf(" wake=% 6.1f",
(threadWake > 0 ? (threadWake - baseTime) / tickFreq * 1e6 : -0.0));
if (threadExecuteStart > 0)
{
printf(" exec=% 6.1f - % 6.1f tasksDone=%5u free=% 6.1f",
(threadExecuteStart > 0 ? (threadExecuteStart - baseTime) / tickFreq * 1e6 : -0.0),
(threadExecuteStop > 0 ? (threadExecuteStop - baseTime) / tickFreq * 1e6 : -0.0),
executedTasks,
(threadFree > 0 ? (threadFree - baseTime) / tickFreq * 1e6 : -0.0));
if (id >= 0)
printf(" active=%s\n", keepActive ? "true" : "false");
else
printf("\n");
}
else
printf(" ------------------------------------------------------------------------------\n");
}
};
ThreadStatistics threads_stat[CV_PROFILE_THREADS]; // 0 - main thread, 1..N - worker threads
#endif
};
class WorkerThread
{
public:
ThreadPool& thread_pool;
const unsigned id;
pthread_t posix_thread;
bool is_created;
volatile bool stop_thread;
volatile bool has_wake_signal;
Ptr job;
pthread_mutex_t mutex;
#if !defined(CV_USE_GLOBAL_WORKERS_COND_VAR)
volatile bool isActive;
pthread_cond_t cond_thread_wake;
#endif
WorkerThread(ThreadPool& thread_pool_, unsigned id_) :
thread_pool(thread_pool_),
id(id_),
posix_thread(0),
is_created(false),
stop_thread(false),
has_wake_signal(false)
#if !defined(CV_USE_GLOBAL_WORKERS_COND_VAR)
, isActive(true)
#endif
{
CV_LOG_VERBOSE(NULL, 1, "MainThread: initializing new worker: "