FazBrowse GitHub Viewer | Trending |
URL:
| Home
Tools: [Download Repo ZIP]   [Original HTTPS Page]

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: "

Back | FazBrowse Home | New Git URL