Use Semaphore for pausing TX/RX threads.

flex-waveform-devel
Mooneer Salem 2025-09-06 01:30:39 -07:00
parent 150b3c9cbf
commit cf7d076ef1
4 changed files with 30 additions and 6 deletions

View File

@ -59,4 +59,19 @@ void FlexRealtimeHelper::setHelperRealTime()
dbus_connection_unref(bus);
}
#endif // defined(USE_RTKIT)
numRealtimeThreads_.fetch_add(1, std::memory_order_relaxed);
}
void FlexRealtimeHelper::clearHelperRealTime()
{
numRealtimeThreads_.fetch_sub(1, std::memory_order_relaxed);
}
void FlexRealtimeHelper::signalRealtimeThreads()
{
int numThreads = numRealtimeThreads_.load(std::memory_order_relaxed);
for (; numThreads > 0; numThreads--)
{
sem_.signal();
}
}

View File

@ -25,7 +25,9 @@
#include <thread>
#include <chrono>
#include <atomic>
#include "../util/IRealtimeHelper.h"
#include "../util/Semaphore.h"
using namespace std::chrono_literals;
@ -45,14 +47,20 @@ public:
// Lets audio system know that we're done with the work on the received
// audio.
virtual void stopRealTimeWork(bool fastMode = false) override { std::this_thread::sleep_for(10ms); }
virtual void stopRealTimeWork(bool fastMode = false) override { sem_.wait(); }
// Reverts real-time priority for current thread.
virtual void clearHelperRealTime() override { /* empty */ }
virtual void clearHelperRealTime() override;
// Returns true if real-time thread MUST sleep ASAP. Failure to do so
// may result in SIGKILL being sent to the process by the kernel.
virtual bool mustStopWork() override { return false; }
void signalRealtimeThreads();
private:
std::atomic<int> numRealtimeThreads_;
Semaphore sem_;
};
#endif // FLEX_REALTIME_HELPER_H

View File

@ -32,7 +32,7 @@ constexpr short FLOAT_TO_SHORT_MULTIPLIER = 32767;
using namespace std::placeholders;
FlexVitaTask::FlexVitaTask(std::shared_ptr<IRealtimeHelper> helper)
FlexVitaTask::FlexVitaTask(std::shared_ptr<FlexRealtimeHelper> helper)
: socket_(-1)
, rxStreamId_(0)
, txStreamId_(0)
@ -462,6 +462,7 @@ void FlexVitaTask::onReceiveVitaMessage_(vita_packet* packet, int length)
i++;
}
inFifo->write(audioInput, half_num_samples); // audio pipeline will resample
helper_->signalRealtimeThreads();
break;
}
default:

View File

@ -28,7 +28,7 @@
#include <arpa/inet.h>
#include "../util/ThreadedObject.h"
#include "../util/IRealtimeHelper.h"
#include "FlexRealtimeHelper.h"
#include "vita.h"
#include "../pipeline/paCallbackData.h"
@ -41,7 +41,7 @@ public:
enum { VITA_PORT = 4992 }; // Hardcoding VITA port because we can only handle one slice at a time.
FlexVitaTask(std::shared_ptr<IRealtimeHelper> helper);
FlexVitaTask(std::shared_ptr<FlexRealtimeHelper> helper);
virtual ~FlexVitaTask();
// Indicates to VitaTask that we've connected to the radio's TCP port.
@ -79,7 +79,7 @@ private:
int64_t lastVitaGenerationTime_;
int minPacketsRequired_;
int64_t timeBeyondExpectedUs_;
std::shared_ptr<IRealtimeHelper> helper_;
std::shared_ptr<FlexRealtimeHelper> helper_;
std::thread rxTxThread_;
bool rxTxThreadRunning_;