mirror of
https://github.com/cjcliffe/CubicSDR.git
synced 2026-07-30 22:14:13 -04:00
Flush queues on terminate() calls to unblock push()s and so ease threads termination
This commit is contained in:
@@ -68,10 +68,10 @@ DemodulatorInstance::DemodulatorInstance() {
|
||||
demodulatorPreThread->setInputQueue("IQDataInput",pipeIQInputData);
|
||||
demodulatorPreThread->setOutputQueue("IQDataOutput",pipeIQDemodData);
|
||||
|
||||
pipeAudioData = std::make_shared< AudioThreadInputQueue>();
|
||||
pipeAudioData = std::make_shared<AudioThreadInputQueue>();
|
||||
pipeAudioData->set_max_num_items(10);
|
||||
|
||||
threadQueueControl = std::make_shared< DemodulatorThreadControlCommandQueue>();
|
||||
threadQueueControl = std::make_shared<DemodulatorThreadControlCommandQueue>();
|
||||
threadQueueControl->set_max_num_items(2);
|
||||
|
||||
demodulatorThread = new DemodulatorThread(this);
|
||||
@@ -102,10 +102,10 @@ DemodulatorInstance::~DemodulatorInstance() {
|
||||
|
||||
#if ENABLE_DIGITAL_LAB
|
||||
delete activeOutput;
|
||||
#endif
|
||||
delete audioThread;
|
||||
delete demodulatorThread;
|
||||
#endif
|
||||
delete demodulatorPreThread;
|
||||
delete demodulatorThread;
|
||||
delete audioThread;
|
||||
|
||||
break;
|
||||
}
|
||||
@@ -174,10 +174,18 @@ void DemodulatorInstance::terminate() {
|
||||
|
||||
// std::cout << "Terminating demodulator audio thread.." << std::endl;
|
||||
audioThread->terminate();
|
||||
|
||||
// std::cout << "Terminating demodulator thread.." << std::endl;
|
||||
demodulatorThread->terminate();
|
||||
|
||||
// std::cout << "Terminating demodulator preprocessor thread.." << std::endl;
|
||||
demodulatorPreThread->terminate();
|
||||
|
||||
//that will actually unblock the currently blocked push().
|
||||
pipeIQInputData->flush();
|
||||
pipeAudioData->flush();
|
||||
pipeIQDemodData->flush();
|
||||
threadQueueControl->flush();
|
||||
}
|
||||
|
||||
std::string DemodulatorInstance::getLabel() {
|
||||
@@ -197,15 +205,23 @@ bool DemodulatorInstance::isTerminated() {
|
||||
bool demodTerminated = demodulatorThread->isTerminated();
|
||||
bool preDemodTerminated = demodulatorPreThread->isTerminated();
|
||||
|
||||
//Cleanup the worker threads, if the threads are indeed terminated
|
||||
if (audioTerminated) {
|
||||
//Cleanup the worker threads, if the threads are indeed terminated.
|
||||
// threads are linked as t_PreDemod ==> t_Demod ==> t_Audio
|
||||
//so terminate in the same order to starve the following threads in succession.
|
||||
//i.e waiting on timed-pop so able to se their stopping flag.
|
||||
|
||||
if (t_Audio) {
|
||||
t_Audio->join();
|
||||
if (preDemodTerminated) {
|
||||
|
||||
if (t_PreDemod) {
|
||||
|
||||
delete t_Audio;
|
||||
t_Audio = nullptr;
|
||||
}
|
||||
#ifdef __APPLE__
|
||||
pthread_join(t_PreDemod, NULL);
|
||||
#else
|
||||
t_PreDemod->join();
|
||||
delete t_PreDemod;
|
||||
#endif
|
||||
t_PreDemod = nullptr;
|
||||
}
|
||||
}
|
||||
|
||||
if (demodTerminated) {
|
||||
@@ -221,18 +237,14 @@ bool DemodulatorInstance::isTerminated() {
|
||||
}
|
||||
}
|
||||
|
||||
if (preDemodTerminated) {
|
||||
|
||||
if (t_PreDemod) {
|
||||
if (audioTerminated) {
|
||||
|
||||
#ifdef __APPLE__
|
||||
pthread_join(t_PreDemod, NULL);
|
||||
#else
|
||||
t_PreDemod->join();
|
||||
delete t_PreDemod;
|
||||
#endif
|
||||
t_PreDemod = nullptr;
|
||||
}
|
||||
if (t_Audio) {
|
||||
t_Audio->join();
|
||||
|
||||
delete t_Audio;
|
||||
t_Audio = nullptr;
|
||||
}
|
||||
}
|
||||
|
||||
bool terminated = audioTerminated && demodTerminated && preDemodTerminated;
|
||||
|
||||
@@ -352,6 +352,10 @@ void DemodulatorPreThread::terminate() {
|
||||
IOThread::terminate();
|
||||
workerThread->terminate();
|
||||
|
||||
//unblock the push()
|
||||
iqOutputQueue->flush();
|
||||
iqInputQueue->flush();
|
||||
|
||||
//wait blocking for termination here, it could be long with lots of modems and we MUST terminate properly,
|
||||
//else better kill the whole application...
|
||||
workerThread->isTerminated(5000);
|
||||
|
||||
@@ -241,7 +241,8 @@ void DemodulatorThread::run() {
|
||||
}
|
||||
|
||||
if ((ati || modemDigital) && localAudioVisOutputQueue != nullptr && localAudioVisOutputQueue->empty()) {
|
||||
AudioThreadInputPtr ati_vis(new AudioThreadInput);
|
||||
|
||||
AudioThreadInputPtr ati_vis = std::make_shared<AudioThreadInput>();
|
||||
|
||||
ati_vis->sampleRate = inp->sampleRate;
|
||||
ati_vis->inputRate = inp->sampleRate;
|
||||
@@ -348,6 +349,11 @@ void DemodulatorThread::run() {
|
||||
|
||||
void DemodulatorThread::terminate() {
|
||||
IOThread::terminate();
|
||||
|
||||
//unblock the curretly blocked push()
|
||||
iqInputQueue->flush();
|
||||
audioOutputQueue->flush();
|
||||
threadQueueControl->flush();
|
||||
}
|
||||
|
||||
bool DemodulatorThread::isMuted() {
|
||||
|
||||
@@ -116,4 +116,7 @@ void DemodulatorWorkerThread::run() {
|
||||
|
||||
void DemodulatorWorkerThread::terminate() {
|
||||
IOThread::terminate();
|
||||
//unblock the push()
|
||||
resultQueue->flush();
|
||||
commandQueue->flush();
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user