2012-11-15 16:51:55 +00:00
|
|
|
/* FasTC
|
|
|
|
* Copyright (c) 2012 University of North Carolina at Chapel Hill. All rights reserved.
|
|
|
|
*
|
|
|
|
* Permission to use, copy, modify, and distribute this software and its documentation for educational,
|
|
|
|
* research, and non-profit purposes, without fee, and without a written agreement is hereby granted,
|
|
|
|
* provided that the above copyright notice, this paragraph, and the following four paragraphs appear
|
|
|
|
* in all copies.
|
|
|
|
*
|
|
|
|
* Permission to incorporate this software into commercial products may be obtained by contacting the
|
|
|
|
* authors or the Office of Technology Development at the University of North Carolina at Chapel Hill <otd@unc.edu>.
|
|
|
|
*
|
|
|
|
* This software program and documentation are copyrighted by the University of North Carolina at Chapel Hill.
|
|
|
|
* The software program and documentation are supplied "as is," without any accompanying services from the
|
|
|
|
* University of North Carolina at Chapel Hill or the authors. The University of North Carolina at Chapel Hill
|
|
|
|
* and the authors do not warrant that the operation of the program will be uninterrupted or error-free. The
|
|
|
|
* end-user understands that the program was developed for research purposes and is advised not to rely
|
|
|
|
* exclusively on the program for any reason.
|
|
|
|
*
|
|
|
|
* IN NO EVENT SHALL THE UNIVERSITY OF NORTH CAROLINA AT CHAPEL HILL OR THE AUTHORS BE LIABLE TO ANY PARTY FOR
|
|
|
|
* DIRECT, INDIRECT, SPECIAL, INCIDENTAL, OR CONSEQUENTIAL DAMAGES, INCLUDING LOST PROFITS, ARISING OUT OF THE
|
|
|
|
* USE OF THIS SOFTWARE AND ITS DOCUMENTATION, EVEN IF THE UNIVERSITY OF NORTH CAROLINA AT CHAPEL HILL OR THE
|
|
|
|
* AUTHORS HAVE BEEN ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
|
|
|
|
*
|
|
|
|
* THE UNIVERSITY OF NORTH CAROLINA AT CHAPEL HILL AND THE AUTHORS SPECIFICALLY DISCLAIM ANY WARRANTIES, INCLUDING,
|
|
|
|
* BUT NOT LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE AND ANY
|
|
|
|
* STATUTORY WARRANTY OF NON-INFRINGEMENT. THE SOFTWARE PROVIDED HEREUNDER IS ON AN "AS IS" BASIS, AND THE UNIVERSITY
|
|
|
|
* OF NORTH CAROLINA AT CHAPEL HILL AND THE AUTHORS HAVE NO OBLIGATIONS TO PROVIDE MAINTENANCE, SUPPORT, UPDATES,
|
|
|
|
* ENHANCEMENTS, OR MODIFICATIONS.
|
|
|
|
*
|
|
|
|
* Please send all BUG REPORTS to <pavel@cs.unc.edu>.
|
|
|
|
*
|
|
|
|
* The authors may be contacted via:
|
|
|
|
*
|
|
|
|
* Pavel Krajcevski
|
|
|
|
* Dept of Computer Science
|
|
|
|
* 201 S Columbia St
|
|
|
|
* Frederick P. Brooks, Jr. Computer Science Bldg
|
|
|
|
* Chapel Hill, NC 27599-3175
|
|
|
|
* USA
|
|
|
|
*
|
|
|
|
* <http://gamma.cs.unc.edu/FasTC/>
|
|
|
|
*/
|
|
|
|
|
2012-09-21 20:57:45 +00:00
|
|
|
#include "WorkerQueue.h"
|
|
|
|
|
2012-11-01 22:56:13 +00:00
|
|
|
#include <algorithm>
|
2013-09-13 23:36:37 +00:00
|
|
|
#include <cstdlib>
|
|
|
|
#include <cstdio>
|
|
|
|
#include <cassert>
|
|
|
|
#include <iostream>
|
2012-09-21 22:14:38 +00:00
|
|
|
|
2014-01-21 19:46:25 +00:00
|
|
|
#include "BPTCCompressor.h"
|
2012-09-21 22:14:38 +00:00
|
|
|
|
2013-11-08 21:21:01 +00:00
|
|
|
using FasTC::CompressionJob;
|
|
|
|
|
2012-09-21 22:14:38 +00:00
|
|
|
template <typename T>
|
|
|
|
static inline void clamp(T &x, const T &min, const T &max) {
|
|
|
|
if(x < min) x = min;
|
|
|
|
else if(x > max) x = max;
|
|
|
|
}
|
2012-09-21 20:57:45 +00:00
|
|
|
|
|
|
|
WorkerThread::WorkerThread(WorkerQueue * parent, uint32 idx)
|
2012-09-29 19:36:42 +00:00
|
|
|
: TCCallable()
|
|
|
|
, m_ThreadIdx(idx)
|
2012-09-21 20:57:45 +00:00
|
|
|
, m_Parent(parent)
|
|
|
|
{ }
|
|
|
|
|
|
|
|
void WorkerThread::operator()() {
|
|
|
|
|
|
|
|
if(!m_Parent) {
|
|
|
|
fprintf(stderr, "%s\n", "Illegal worker thread initialization -- parent is NULL.");
|
|
|
|
return;
|
|
|
|
}
|
|
|
|
|
|
|
|
CompressionFunc f = m_Parent->GetCompressionFunc();
|
2012-11-01 22:56:13 +00:00
|
|
|
CompressionFuncWithStats fStat = m_Parent->GetCompressionFuncWithStats();
|
2013-09-29 01:42:24 +00:00
|
|
|
std::ostream *logStream = m_Parent->GetLogStream();
|
2012-11-01 22:56:13 +00:00
|
|
|
|
2013-09-29 01:42:24 +00:00
|
|
|
if(!(f || (fStat && logStream))) {
|
2012-09-21 20:57:45 +00:00
|
|
|
fprintf(stderr, "%s\n", "Illegal worker queue initialization -- compression func is NULL.");
|
|
|
|
return;
|
|
|
|
}
|
|
|
|
|
2012-09-25 21:05:52 +00:00
|
|
|
bool quitFlag = false;
|
|
|
|
while(!quitFlag) {
|
2012-09-21 22:14:38 +00:00
|
|
|
|
2013-08-26 20:54:08 +00:00
|
|
|
switch(m_Parent->AcceptThreadData(m_ThreadIdx)) {
|
2012-09-25 21:05:52 +00:00
|
|
|
|
|
|
|
case eAction_Quit:
|
|
|
|
{
|
2013-08-26 20:54:08 +00:00
|
|
|
quitFlag = true;
|
|
|
|
break;
|
2012-09-25 21:05:52 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
case eAction_Wait:
|
|
|
|
{
|
2013-08-26 20:54:08 +00:00
|
|
|
TCThread::Yield();
|
|
|
|
break;
|
2012-09-25 21:05:52 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
case eAction_DoWork:
|
|
|
|
{
|
2013-11-08 21:21:01 +00:00
|
|
|
const CompressionJob &job = m_Parent->GetCompressionJob();
|
|
|
|
|
|
|
|
uint32 start[2];
|
|
|
|
m_Parent->GetStartForThread(m_ThreadIdx, start);
|
2013-03-09 18:36:39 +00:00
|
|
|
|
2013-11-08 21:21:01 +00:00
|
|
|
uint32 end[2];
|
|
|
|
m_Parent->GetEndForThread(m_ThreadIdx, end);
|
|
|
|
|
|
|
|
CompressionJob cj (job.Format(),
|
|
|
|
job.InBuf(), job.OutBuf(),
|
|
|
|
job.Width(), job.Height(),
|
|
|
|
start[0], start[1],
|
|
|
|
end[0], end[1]);
|
2013-08-26 20:54:08 +00:00
|
|
|
if(f)
|
|
|
|
(*f)(cj);
|
|
|
|
else
|
2013-09-29 01:42:24 +00:00
|
|
|
(*fStat)(cj, logStream);
|
2012-11-01 22:56:13 +00:00
|
|
|
|
2013-08-26 20:54:08 +00:00
|
|
|
break;
|
2012-09-25 21:05:52 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
default:
|
|
|
|
{
|
2013-08-26 20:54:08 +00:00
|
|
|
fprintf(stderr, "Unrecognized thread command!\n");
|
|
|
|
quitFlag = true;
|
|
|
|
break;
|
2012-09-25 21:05:52 +00:00
|
|
|
}
|
2012-09-21 22:14:38 +00:00
|
|
|
}
|
2012-09-21 20:57:45 +00:00
|
|
|
}
|
|
|
|
|
2012-09-21 22:14:38 +00:00
|
|
|
m_Parent->NotifyWorkerFinished();
|
2012-09-21 20:57:45 +00:00
|
|
|
|
2012-09-21 22:14:38 +00:00
|
|
|
return;
|
2012-09-21 20:57:45 +00:00
|
|
|
}
|
|
|
|
|
2012-09-21 22:14:38 +00:00
|
|
|
WorkerQueue::WorkerQueue(
|
2012-09-21 22:43:35 +00:00
|
|
|
uint32 numCompressions,
|
2012-09-21 22:14:38 +00:00
|
|
|
uint32 numThreads,
|
|
|
|
uint32 jobSize,
|
2013-11-08 21:21:01 +00:00
|
|
|
const CompressionJob &job,
|
|
|
|
CompressionFunc func
|
2012-09-21 22:14:38 +00:00
|
|
|
)
|
2012-09-21 22:43:35 +00:00
|
|
|
: m_NumCompressions(0)
|
2012-11-01 22:56:13 +00:00
|
|
|
, m_TotalNumCompressions(std::max(uint32(1), numCompressions))
|
2012-09-21 22:43:35 +00:00
|
|
|
, m_NumThreads(numThreads)
|
2012-09-25 21:05:52 +00:00
|
|
|
, m_WaitingThreads(0)
|
2012-09-21 22:14:38 +00:00
|
|
|
, m_ActiveThreads(0)
|
2012-11-01 22:56:13 +00:00
|
|
|
, m_JobSize(std::max(uint32(1), jobSize))
|
2013-11-08 21:21:01 +00:00
|
|
|
, m_Job(job)
|
2012-09-26 17:31:39 +00:00
|
|
|
, m_NextBlock(0)
|
2012-09-21 22:14:38 +00:00
|
|
|
, m_CompressionFunc(func)
|
2012-11-01 22:56:13 +00:00
|
|
|
, m_CompressionFuncWithStats(NULL)
|
2013-09-29 01:42:24 +00:00
|
|
|
, m_LogStream(NULL)
|
2012-11-01 22:56:13 +00:00
|
|
|
{
|
|
|
|
clamp(m_NumThreads, uint32(1), uint32(kMaxNumWorkerThreads));
|
|
|
|
}
|
|
|
|
|
|
|
|
WorkerQueue::WorkerQueue(
|
|
|
|
uint32 numCompressions,
|
|
|
|
uint32 numThreads,
|
|
|
|
uint32 jobSize,
|
2013-11-08 21:21:01 +00:00
|
|
|
const CompressionJob &job,
|
2012-11-01 22:56:13 +00:00
|
|
|
CompressionFuncWithStats func,
|
2013-11-08 21:21:01 +00:00
|
|
|
std::ostream *logStream
|
2012-11-01 22:56:13 +00:00
|
|
|
)
|
|
|
|
: m_NumCompressions(0)
|
|
|
|
, m_TotalNumCompressions(std::max(uint32(1), numCompressions))
|
|
|
|
, m_NumThreads(numThreads)
|
|
|
|
, m_WaitingThreads(0)
|
|
|
|
, m_ActiveThreads(0)
|
|
|
|
, m_JobSize(std::max(uint32(1), jobSize))
|
2013-11-08 21:21:01 +00:00
|
|
|
, m_Job(job)
|
2012-11-01 22:56:13 +00:00
|
|
|
, m_NextBlock(0)
|
|
|
|
, m_CompressionFunc(NULL)
|
|
|
|
, m_CompressionFuncWithStats(func)
|
2013-09-29 01:42:24 +00:00
|
|
|
, m_LogStream(logStream)
|
2012-09-21 22:14:38 +00:00
|
|
|
{
|
2012-09-21 22:43:35 +00:00
|
|
|
clamp(m_NumThreads, uint32(1), uint32(kMaxNumWorkerThreads));
|
2012-09-21 20:57:45 +00:00
|
|
|
}
|
|
|
|
|
2012-09-21 22:14:38 +00:00
|
|
|
void WorkerQueue::Run() {
|
|
|
|
|
|
|
|
// Spawn a bunch of threads...
|
2012-09-29 19:36:42 +00:00
|
|
|
TCLock lock(m_Mutex);
|
2012-11-07 22:10:26 +00:00
|
|
|
for(uint32 i = 0; i < m_NumThreads; i++) {
|
2012-09-29 19:36:42 +00:00
|
|
|
m_Workers[i] = new WorkerThread(this, i);
|
|
|
|
m_ThreadHandles[m_ActiveThreads] = new TCThread(*m_Workers[i]);
|
2012-09-21 22:14:38 +00:00
|
|
|
m_ActiveThreads++;
|
2012-09-21 20:57:45 +00:00
|
|
|
}
|
|
|
|
|
2012-09-21 22:43:35 +00:00
|
|
|
m_StopWatch.Reset();
|
|
|
|
m_StopWatch.Start();
|
|
|
|
|
2012-09-26 17:31:39 +00:00
|
|
|
m_NextBlock = 0;
|
2012-09-25 21:05:52 +00:00
|
|
|
m_WaitingThreads = 0;
|
|
|
|
|
2012-09-21 22:14:38 +00:00
|
|
|
// Wait for them to finish...
|
|
|
|
while(m_ActiveThreads > 0) {
|
2012-09-29 19:36:42 +00:00
|
|
|
m_CV.Wait(lock);
|
2012-09-21 20:57:45 +00:00
|
|
|
}
|
|
|
|
|
2012-09-21 22:43:35 +00:00
|
|
|
m_StopWatch.Stop();
|
|
|
|
|
2012-09-21 22:14:38 +00:00
|
|
|
// Join them all together..
|
2012-11-07 22:10:26 +00:00
|
|
|
for(uint32 i = 0; i < m_NumThreads; i++) {
|
2012-09-29 19:36:42 +00:00
|
|
|
m_ThreadHandles[i]->Join();
|
2012-09-21 22:14:38 +00:00
|
|
|
delete m_ThreadHandles[i];
|
2012-09-29 19:36:42 +00:00
|
|
|
delete m_Workers[i];
|
2012-09-21 20:57:45 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2012-09-21 22:14:38 +00:00
|
|
|
void WorkerQueue::NotifyWorkerFinished() {
|
|
|
|
{
|
2012-09-29 19:36:42 +00:00
|
|
|
TCLock lock(m_Mutex);
|
2012-09-21 22:14:38 +00:00
|
|
|
m_ActiveThreads--;
|
|
|
|
}
|
2012-09-29 19:36:42 +00:00
|
|
|
m_CV.NotifyOne();
|
2012-09-21 22:14:38 +00:00
|
|
|
}
|
2012-09-21 20:57:45 +00:00
|
|
|
|
2012-09-21 22:14:38 +00:00
|
|
|
WorkerThread::EAction WorkerQueue::AcceptThreadData(uint32 threadIdx) {
|
2013-01-29 01:20:52 +00:00
|
|
|
if(threadIdx >= m_ActiveThreads) {
|
2012-09-21 22:14:38 +00:00
|
|
|
return WorkerThread::eAction_Quit;
|
2012-09-21 20:57:45 +00:00
|
|
|
}
|
|
|
|
|
2012-09-21 22:14:38 +00:00
|
|
|
// How many blocks total do we have?
|
2013-11-08 21:21:01 +00:00
|
|
|
uint32 blockDim[2];
|
|
|
|
GetBlockDimensions(m_Job.Format(), blockDim);
|
|
|
|
const uint32 totalBlocks = (m_Job.Width() * m_Job.Height()) / (blockDim[0] * blockDim[1]);
|
2012-09-21 22:14:38 +00:00
|
|
|
|
|
|
|
// Make sure we have exclusive access...
|
2012-09-29 19:36:42 +00:00
|
|
|
TCLock lock(m_Mutex);
|
2012-09-21 22:14:38 +00:00
|
|
|
|
|
|
|
// If we've completed all blocks, then mark the thread for
|
|
|
|
// completion.
|
2012-09-25 21:05:52 +00:00
|
|
|
if(m_NextBlock == totalBlocks) {
|
|
|
|
if(m_NumCompressions < m_TotalNumCompressions) {
|
|
|
|
if(++m_WaitingThreads == m_ActiveThreads) {
|
2013-08-26 20:54:08 +00:00
|
|
|
m_NextBlock = 0;
|
|
|
|
m_WaitingThreads = 0;
|
2012-09-25 21:05:52 +00:00
|
|
|
} else {
|
2013-08-26 20:54:08 +00:00
|
|
|
return WorkerThread::eAction_Wait;
|
2012-09-25 21:05:52 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
else {
|
|
|
|
return WorkerThread::eAction_Quit;
|
|
|
|
}
|
2012-09-21 20:57:45 +00:00
|
|
|
}
|
|
|
|
|
2012-09-21 22:14:38 +00:00
|
|
|
// Otherwise, this thread's offset is the current block...
|
|
|
|
m_Offsets[threadIdx] = m_NextBlock;
|
|
|
|
|
|
|
|
// The number of blocks to process is either the job size
|
|
|
|
// or the number of blocks remaining.
|
2012-11-01 22:56:13 +00:00
|
|
|
int blocksProcessed = std::min(m_JobSize, totalBlocks - m_NextBlock);
|
2012-09-21 22:14:38 +00:00
|
|
|
m_NumBlocks[threadIdx] = blocksProcessed;
|
2012-09-21 20:57:45 +00:00
|
|
|
|
2012-09-21 22:14:38 +00:00
|
|
|
// Make sure the next block is updated.
|
|
|
|
m_NextBlock += blocksProcessed;
|
2012-09-21 20:57:45 +00:00
|
|
|
|
2012-09-25 21:05:52 +00:00
|
|
|
if(m_NextBlock == totalBlocks) {
|
|
|
|
++m_NumCompressions;
|
2012-09-21 22:43:35 +00:00
|
|
|
}
|
|
|
|
|
2012-09-21 22:14:38 +00:00
|
|
|
return WorkerThread::eAction_DoWork;
|
2012-09-21 20:57:45 +00:00
|
|
|
}
|
|
|
|
|
2013-11-08 21:21:01 +00:00
|
|
|
void WorkerQueue::GetStartForThread(const uint32 threadIdx, uint32 (&start)[2]) {
|
2012-09-21 22:14:38 +00:00
|
|
|
assert(threadIdx >= 0);
|
2013-11-11 23:45:09 +00:00
|
|
|
assert(threadIdx < m_NumThreads);
|
2012-09-21 22:14:38 +00:00
|
|
|
assert(m_Offsets[threadIdx] >= 0);
|
2012-09-21 20:57:45 +00:00
|
|
|
|
2013-11-08 21:21:01 +00:00
|
|
|
const uint32 blockIdx = m_Offsets[threadIdx];
|
|
|
|
m_Job.BlockIdxToCoords(blockIdx, start);
|
2012-09-21 20:57:45 +00:00
|
|
|
}
|
|
|
|
|
2013-11-08 21:21:01 +00:00
|
|
|
void WorkerQueue::GetEndForThread(const uint32 threadIdx, uint32 (&end)[2]) {
|
2012-09-21 22:14:38 +00:00
|
|
|
assert(threadIdx >= 0);
|
2013-11-11 23:45:09 +00:00
|
|
|
assert(threadIdx < m_NumThreads);
|
2013-11-08 21:21:01 +00:00
|
|
|
assert(m_Offsets[threadIdx] >= 0);
|
|
|
|
assert(m_NumBlocks[threadIdx] >= 0);
|
2012-09-21 20:57:45 +00:00
|
|
|
|
2013-11-08 21:21:01 +00:00
|
|
|
const uint32 blockIdx = m_Offsets[threadIdx] + m_NumBlocks[threadIdx];
|
|
|
|
m_Job.BlockIdxToCoords(blockIdx, end);
|
2012-09-21 20:57:45 +00:00
|
|
|
}
|