163 lines
3.7 KiB
C++
163 lines
3.7 KiB
C++
#include "nativestreamer.h"
|
|
#include "lib9pclient.h"
|
|
#include <QDebug>
|
|
|
|
NativeStreamer::NativeStreamer(Lib9pClient *client, RestartPolicy policy, QObject *parent)
|
|
: QObject(parent)
|
|
, m_client(client)
|
|
, m_policy(policy)
|
|
{
|
|
m_restartTimer = new QTimer(this);
|
|
m_restartTimer->setSingleShot(true);
|
|
m_restartTimer->setInterval(500); // 500ms delay before restart
|
|
connect(m_restartTimer, &QTimer::timeout, this, [this]() {
|
|
if (!m_stopping.load() && !m_path.isEmpty()) {
|
|
start(m_path);
|
|
}
|
|
});
|
|
}
|
|
|
|
NativeStreamer::~NativeStreamer()
|
|
{
|
|
stop();
|
|
}
|
|
|
|
void NativeStreamer::start(const QString &path)
|
|
{
|
|
if (m_running.load()) {
|
|
stop();
|
|
}
|
|
|
|
m_path = path;
|
|
m_running.store(true);
|
|
m_stopping.store(false);
|
|
|
|
// Open the file first (on main thread, before spawning worker)
|
|
int fid = m_client->open(path);
|
|
if (fid < 0) {
|
|
emit errorOccurred(m_client->lastError());
|
|
m_running.store(false);
|
|
scheduleRestart();
|
|
return;
|
|
}
|
|
m_fid.store(fid);
|
|
|
|
// Start worker thread
|
|
m_thread = QThread::create([this]() { run(); });
|
|
connect(m_thread, &QThread::finished, this, &NativeStreamer::onThreadFinished);
|
|
m_thread->start();
|
|
}
|
|
|
|
void NativeStreamer::stop()
|
|
{
|
|
m_stopping.store(true);
|
|
m_running.store(false);
|
|
m_restartTimer->stop();
|
|
|
|
// Close the fid to unblock any pending read
|
|
int fid = m_fid.exchange(-1);
|
|
if (fid >= 0) {
|
|
m_client->closeFid(fid);
|
|
}
|
|
|
|
if (m_thread) {
|
|
// Disconnect to prevent onThreadFinished from running after we delete
|
|
disconnect(m_thread, &QThread::finished, this, &NativeStreamer::onThreadFinished);
|
|
|
|
// Wait for thread to finish (it should exit after fid is closed)
|
|
if (!m_thread->wait(2000)) {
|
|
// Thread didn't exit in time, force terminate
|
|
qWarning() << "NativeStreamer: thread didn't exit cleanly, terminating";
|
|
m_thread->terminate();
|
|
m_thread->wait();
|
|
}
|
|
delete m_thread;
|
|
m_thread = nullptr;
|
|
}
|
|
}
|
|
|
|
bool NativeStreamer::isRunning() const
|
|
{
|
|
return m_running.load();
|
|
}
|
|
|
|
void NativeStreamer::setGuard(std::function<bool()> guard)
|
|
{
|
|
m_guard = guard;
|
|
}
|
|
|
|
void NativeStreamer::onThreadFinished()
|
|
{
|
|
// Thread finished naturally (not via stop()) - clean up and maybe restart
|
|
if (m_thread) {
|
|
m_thread->deleteLater();
|
|
m_thread = nullptr;
|
|
}
|
|
|
|
// Close any remaining fid
|
|
int fid = m_fid.exchange(-1);
|
|
if (fid >= 0) {
|
|
m_client->closeFid(fid);
|
|
}
|
|
|
|
scheduleRestart();
|
|
}
|
|
|
|
void NativeStreamer::scheduleRestart()
|
|
{
|
|
if (m_stopping.load()) {
|
|
return;
|
|
}
|
|
|
|
bool shouldRestart = false;
|
|
switch (m_policy) {
|
|
case Oneshot:
|
|
shouldRestart = false;
|
|
break;
|
|
case Looping:
|
|
shouldRestart = true;
|
|
break;
|
|
case Guarded:
|
|
shouldRestart = m_guard && m_guard();
|
|
break;
|
|
}
|
|
|
|
if (shouldRestart) {
|
|
m_restartTimer->start();
|
|
} else {
|
|
emit finished();
|
|
}
|
|
}
|
|
|
|
void NativeStreamer::run()
|
|
{
|
|
constexpr int BUF_SIZE = 8192;
|
|
char buf[BUF_SIZE];
|
|
|
|
while (m_running.load()) {
|
|
int fid = m_fid.load();
|
|
if (fid < 0) break;
|
|
|
|
int n = m_client->readFid(fid, buf, BUF_SIZE);
|
|
if (n < 0) {
|
|
// Error
|
|
if (m_running.load()) {
|
|
emit errorOccurred(m_client->lastError());
|
|
}
|
|
break;
|
|
}
|
|
if (n == 0) {
|
|
// EOF
|
|
break;
|
|
}
|
|
|
|
// Emit data on main thread via queued connection
|
|
QByteArray data(buf, n);
|
|
QMetaObject::invokeMethod(this, [this, data]() {
|
|
emit dataReady(data);
|
|
}, Qt::QueuedConnection);
|
|
}
|
|
|
|
m_running.store(false);
|
|
}
|