Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
18 commits
Select commit Hold shift + click to select a range
5a9abf5
copy the wait predicate into a test before waking the condition, to r…
odysseus654 Sep 23, 2021
e839857
breaking the pool _mutex into two smaller _slaveMutex and _poolMutex …
odysseus654 Sep 23, 2021
0c35b43
pull out data shared between threads / remove "friend" declarations
odysseus654 Sep 23, 2021
e60db78
add shared_mutex in an attempt to reduce the number of mutexes used i…
odysseus654 Sep 24, 2021
536b1ed
stripping out the shared_mutex for a contract-based threading model (…
odysseus654 Sep 24, 2021
aa2c307
suggested changes:
odysseus654 Sep 24, 2021
569afb1
resolve potential race condition
odysseus654 Sep 24, 2021
bc5d9de
renamed _pool to _data, to bring it more in line with other users
odysseus654 Sep 25, 2021
45d3312
remove last reference to _slaves
odysseus654 Sep 25, 2021
3a4e8a7
better management of numThreads when launching new workers
odysseus654 Sep 26, 2021
146da44
ensure that newly-created threads are properly configured before leav…
odysseus654 Sep 26, 2021
83f2b52
added some missed field initializations
odysseus654 Sep 26, 2021
df59a42
more tweaks to improve worker thread startup
odysseus654 Sep 27, 2021
5389d3f
New threads now unconditionally enter a wait without checking if numS…
odysseus654 Sep 28, 2021
ecba866
condition.wait is still letting threads through before it's time, sta…
odysseus654 Sep 28, 2021
f7e61b8
results of linux debugging, trying to get the wait conditions working…
odysseus654 Nov 8, 2021
1f9347d
Merge remote-tracking branch 'remotes/vircadia/master' into pr/pool-c…
odysseus654 Dec 4, 2021
20df505
Merge remote-tracking branch 'remotes/vircadia/master' into pr/pool-c…
odysseus654 Jan 25, 2022
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
204 changes: 121 additions & 83 deletions assignment-client/src/audio/AudioMixerSlavePool.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -18,65 +18,78 @@

#include <ThreadHelpers.h>

void AudioMixerSlaveThread::run() {
while (true) {
wait();

// iterate over all available nodes
SharedNodePointer node;
while (try_pop(node)) {
(this->*_function)(node);
}
void AudioMixerWorkerThread::run() {
assert(_data.numStarted < _data.numThreads);
_data.numStarted++;

bool starting = true;
while (true) {
bool stopping = _stop;
notify(stopping);
if (stopping) {
return;
}

wait(starting);
starting = false;

if (_function) {
// iterate over all available nodes
SharedNodePointer node;
while (try_pop(node)) {
(this->*_function)(node);
}
}
}
}

void AudioMixerSlaveThread::wait() {
void AudioMixerWorkerThread::wait(bool starting) {
assert(_data.numStarted <= _data.numThreads);

{
Lock lock(_pool._mutex);
_pool._slaveCondition.wait(lock, [&] {
assert(_pool._numStarted <= _pool._numThreads);
return _pool._numStarted != _pool._numThreads;
});
++_pool._numStarted;
Lock workerLock(_data.workerMutex);
if (starting || _data.numStarted == _data.numThreads) {
do { // this is equivalent to the two-parameter "wait" call, except we're doing the test afterwards
_data.workerCondition.wait(workerLock);
} while (_data.numStarted == _data.numThreads);
}
_data.numStarted++;
}

if (_pool._configure) {
_pool._configure(*this);
if (_data.configure) {
_data.configure(*this);
}
_function = _pool._function;
_function = _data.function;
}

void AudioMixerSlaveThread::notify(bool stopping) {
{
Lock lock(_pool._mutex);
assert(_pool._numFinished < _pool._numThreads);
++_pool._numFinished;
if (stopping) {
++_pool._numStopped;
}
void AudioMixerWorkerThread::notify(bool stopping) {
assert(_data.numFinished < _data.numThreads);
assert(_data.numFinished <= _data.numStarted);
int numFinished = ++_data.numFinished;
if (stopping) {
_data.numStopped++;
assert(_data.numStopped <= _data.numFinished);
}

if (numFinished == _data.numThreads) {
Lock poolLock(_data.poolMutex);
_data.poolCondition.notify_one();
}
_pool._poolCondition.notify_one();
}

bool AudioMixerSlaveThread::try_pop(SharedNodePointer& node) {
return _pool._queue.try_pop(node);
bool AudioMixerWorkerThread::try_pop(SharedNodePointer& node) {
return _data.queue.try_pop(node);
}

void AudioMixerSlavePool::processPackets(ConstIter begin, ConstIter end) {
_function = &AudioMixerSlave::processPackets;
_configure = [](AudioMixerSlave& slave) {};
_data.function = &AudioMixerSlave::processPackets;
_data.configure = [](AudioMixerSlave& slave) {};
run(begin, end);
}

void AudioMixerSlavePool::mix(ConstIter begin, ConstIter end, unsigned int frame, int numToRetain) {
_function = &AudioMixerSlave::mix;
_configure = [=](AudioMixerSlave& slave) {
_data.function = &AudioMixerSlave::mix;
_data.configure = [=](AudioMixerSlave& slave) {
slave.configureMix(_begin, _end, frame, numToRetain);
};

Expand All @@ -89,46 +102,49 @@ void AudioMixerSlavePool::run(ConstIter begin, ConstIter end) {

// fill the queue
std::for_each(_begin, _end, [&](const SharedNodePointer& node) {
_queue.push(node);
_data.queue.push(node);
});

// run
{
Lock lock(_mutex);

// run
_numStarted = _numFinished = 0;
_slaveCondition.notify_all();

// wait
_poolCondition.wait(lock, [&] {
assert(_numFinished <= _numThreads);
return _numFinished == _numThreads;
});
Lock workerLock(_data.workerMutex);
_data.numStarted = _data.numFinished = 0;
_data.workerCondition.notify_all();
}

assert(_numStarted == _numThreads);
// wait
{
Lock poolLock(_data.poolMutex);
if (_data.numFinished < _data.numThreads) {
_data.poolCondition.wait(poolLock, [&] {
assert(_data.numFinished <= _data.numThreads);
return _data.numFinished == _data.numThreads;
});
}
}
assert(_data.numStarted == _data.numThreads);

assert(_queue.empty());
assert(_data.queue.empty());
}

void AudioMixerSlavePool::each(std::function<void(AudioMixerSlave& slave)> functor) {
for (auto& slave : _slaves) {
functor(*slave.get());
for (auto& worker : _workers) {
functor(*worker.get());
}
}

#ifdef DEBUG_EVENT_QUEUE
void AudioMixerSlavePool::queueStats(QJsonObject& stats) {
unsigned i = 0;
for (auto& slave : _slaves) {
int queueSize = ::hifi::qt::getEventQueueSize(slave.get());
for (auto& worker : _workers) {
int queueSize = ::hifi::qt::getEventQueueSize(worker.get());
QString queueName = QString("audio_thread_event_queue_%1").arg(i);
stats[queueName] = queueSize;

i++;
}
}
#endif // DEBUG_EVENT_QUEUE
#endif // DEBUG_EVENT_QUEUE

void AudioMixerSlavePool::setNumThreads(int numThreads) {
// clamp to allowed size
Expand All @@ -151,54 +167,76 @@ void AudioMixerSlavePool::setNumThreads(int numThreads) {
}

void AudioMixerSlavePool::resize(int numThreads) {
assert(_numThreads == (int)_slaves.size());
assert(_data.numThreads == (int)_workers.size());

qDebug("%s: set %d threads (was %d)", __FUNCTION__, numThreads, _numThreads);
qDebug("%s: set %d threads (was %d)", __FUNCTION__, numThreads, _data.numThreads);

Lock lock(_mutex);
if (numThreads > _data.numThreads) {
assert(_data.numFinished == _data.numThreads);

if (numThreads > _numThreads) {
// start new slaves
for (int i = 0; i < numThreads - _numThreads; ++i) {
auto slave = new AudioMixerSlaveThread(*this, _workerSharedData);
QObject::connect(slave, &QThread::started, [] { setThreadName("AudioMixerSlaveThread"); });
slave->start();
_slaves.emplace_back(slave);
{
Lock workerLock(_data.workerMutex);
while (numThreads > _data.numThreads) {
_data.numThreads++;
auto worker = new AudioMixerWorkerThread(_data, _workerSharedData);
QObject::connect(worker, &QThread::started, [] { setThreadName("AudioMixerSlaveThread"); });
worker->start();
_workers.emplace_back(worker);
}
}

// wait for the new workers to wake up and enter the wait
Lock poolLock(_data.poolMutex);
if (_data.numFinished < _data.numThreads) {
_data.poolCondition.wait(poolLock, [&] {
assert(_data.numFinished <= _data.numThreads);
return _data.numFinished == _data.numThreads;
});
}
} else if (numThreads < _numThreads) {
auto extraBegin = _slaves.begin() + numThreads;
} else if (numThreads < _data.numThreads) {
auto extraBegin = _workers.begin() + numThreads;

// mark slaves to stop...
auto slave = extraBegin;
while (slave != _slaves.end()) {
(*slave)->_stop = true;
++slave;
auto worker = extraBegin;
while (worker != _workers.end()) {
(*worker)->stop();
++worker;
}

// ...cycle them until they do stop...
_numStopped = 0;
while (_numStopped != (_numThreads - numThreads)) {
_numStarted = _numFinished = _numStopped;
_slaveCondition.notify_all();
_poolCondition.wait(lock, [&] {
assert(_numFinished <= _numThreads);
return _numFinished == _numThreads;
});
_data.numStopped = 0;
while (_data.numStopped != (_data.numThreads - numThreads)) {
{
Lock workerLock(_data.workerMutex);
_data.numStarted = _data.numFinished = 0;
_data.workerCondition.notify_all();
}

Lock poolLock(_data.poolMutex);
if (_data.numFinished < _data.numThreads) {
_data.poolCondition.wait(poolLock, [&] {
assert(_data.numFinished <= _data.numThreads);
return _data.numFinished == _data.numThreads;
});
}
assert(_data.numStopped == (_data.numThreads - numThreads));
}

// ...wait for threads to finish...
slave = extraBegin;
while (slave != _slaves.end()) {
QThread* thread = reinterpret_cast<QThread*>(slave->get());
worker = extraBegin;
while (worker != _workers.end()) {
QThread* thread = reinterpret_cast<QThread*>(worker->get());
static const int MAX_THREAD_WAIT_TIME = 10;
thread->wait(MAX_THREAD_WAIT_TIME);
++slave;
++worker;
}

// ...and erase them
_slaves.erase(extraBegin, _slaves.end());
_workers.erase(extraBegin, _workers.end());
}

_numThreads = _numStarted = _numFinished = numThreads;
assert(_numThreads == (int)_slaves.size());
_data.numThreads = _data.numStarted = _data.numFinished = numThreads;
assert(_data.numThreads == (int)_workers.size());
qDebug("%s: completed", __FUNCTION__);
}
Loading