238 lines
5.7 KiB
C++
238 lines
5.7 KiB
C++
#include "lib9pstreamer.h"
|
|
#include <QDebug>
|
|
|
|
// --- Worker implementation ---
|
|
|
|
Lib9pStreamer::Worker::Worker(const QString &path)
|
|
: m_path(path)
|
|
{
|
|
}
|
|
|
|
Lib9pStreamer::Worker::~Worker()
|
|
{
|
|
QMutexLocker lock(&m_fidMutex);
|
|
if (m_fid >= 0 && m_client) {
|
|
m_client->closeFid(m_fid);
|
|
m_fid = -1;
|
|
}
|
|
delete m_client;
|
|
m_client = nullptr;
|
|
}
|
|
|
|
void Lib9pStreamer::Worker::requestStop()
|
|
{
|
|
m_stopRequested = true;
|
|
// Close the fid to unblock any blocking read
|
|
QMutexLocker lock(&m_fidMutex);
|
|
if (m_fid >= 0 && m_client) {
|
|
m_client->closeFid(m_fid);
|
|
m_fid = -1;
|
|
}
|
|
}
|
|
|
|
void Lib9pStreamer::Worker::run()
|
|
{
|
|
// Create client in worker thread (Go library is thread-safe)
|
|
m_client = new Lib9pClient();
|
|
if (!m_client->connectDefault()) {
|
|
emit finished(false, QStringLiteral("Failed to connect: ") + m_client->lastError());
|
|
return;
|
|
}
|
|
|
|
int fid = m_client->open(m_path);
|
|
if (fid < 0) {
|
|
emit finished(false, QStringLiteral("Failed to open ") + m_path + ": " + m_client->lastError());
|
|
return;
|
|
}
|
|
|
|
{
|
|
QMutexLocker lock(&m_fidMutex);
|
|
if (m_stopRequested) {
|
|
// Stop requested while opening
|
|
m_client->closeFid(fid);
|
|
emit finished(true, QString());
|
|
return;
|
|
}
|
|
m_fid = fid;
|
|
}
|
|
|
|
// Read loop
|
|
char buf[8192];
|
|
while (!m_stopRequested) {
|
|
int n = m_client->readFid(m_fid, buf, sizeof(buf));
|
|
if (n < 0) {
|
|
// Error — could be from requestStop closing the fid
|
|
if (m_stopRequested) {
|
|
emit finished(true, QString());
|
|
} else {
|
|
emit finished(false, m_client->lastError());
|
|
}
|
|
return;
|
|
}
|
|
if (n == 0) {
|
|
// EOF — normal for event stream after returning data
|
|
emit finished(true, QString());
|
|
return;
|
|
}
|
|
emit dataReady(QByteArray(buf, n));
|
|
}
|
|
|
|
// Stop requested
|
|
emit finished(true, QString());
|
|
}
|
|
|
|
// --- Lib9pStreamer implementation ---
|
|
|
|
Lib9pStreamer::Lib9pStreamer(RestartPolicy policy, QObject *parent)
|
|
: QObject(parent)
|
|
, m_policy(policy)
|
|
{
|
|
m_restartTimer = new QTimer(this);
|
|
m_restartTimer->setSingleShot(true);
|
|
m_restartTimer->setInterval(100);
|
|
connect(m_restartTimer, &QTimer::timeout, this, &Lib9pStreamer::scheduleRestart);
|
|
}
|
|
|
|
Lib9pStreamer::~Lib9pStreamer()
|
|
{
|
|
stop();
|
|
}
|
|
|
|
void Lib9pStreamer::start(const QString &path)
|
|
{
|
|
if (m_worker) {
|
|
stop();
|
|
}
|
|
|
|
m_path = path;
|
|
m_stopping = false;
|
|
startWorker();
|
|
}
|
|
|
|
void Lib9pStreamer::stop()
|
|
{
|
|
m_stopping = true;
|
|
m_restartTimer->stop();
|
|
stopWorker();
|
|
}
|
|
|
|
bool Lib9pStreamer::isRunning() const
|
|
{
|
|
return m_worker != nullptr;
|
|
}
|
|
|
|
void Lib9pStreamer::setGuard(std::function<bool()> guard)
|
|
{
|
|
m_guard = guard;
|
|
}
|
|
|
|
void Lib9pStreamer::startWorker()
|
|
{
|
|
m_worker = new Worker(m_path);
|
|
m_workerThread = new QThread();
|
|
m_worker->moveToThread(m_workerThread);
|
|
|
|
connect(m_workerThread, &QThread::started, m_worker, &Worker::run);
|
|
connect(m_worker, &Worker::dataReady, this, &Lib9pStreamer::onWorkerData);
|
|
connect(m_worker, &Worker::finished, this, &Lib9pStreamer::onWorkerFinished);
|
|
|
|
m_workerThread->start();
|
|
}
|
|
|
|
void Lib9pStreamer::stopWorker()
|
|
{
|
|
if (!m_worker)
|
|
return;
|
|
|
|
m_worker->requestStop(); // Also closes fid to unblock read
|
|
|
|
// Give a very short time for the thread to notice the close.
|
|
// If it doesn't exit quickly, abandon it (it will clean up when the
|
|
// read finally returns).
|
|
if (m_workerThread->wait(10)) {
|
|
// Thread exited cleanly
|
|
delete m_worker;
|
|
delete m_workerThread;
|
|
} else {
|
|
// Thread still blocked — abandon it
|
|
// Disconnect signals so we don't get callbacks after we're gone
|
|
disconnect(m_worker, nullptr, this, nullptr);
|
|
// Let the thread run to completion on its own
|
|
connect(m_workerThread, &QThread::finished, m_worker, &QObject::deleteLater);
|
|
connect(m_workerThread, &QThread::finished, m_workerThread, &QObject::deleteLater);
|
|
}
|
|
|
|
m_worker = nullptr;
|
|
m_workerThread = nullptr;
|
|
}
|
|
|
|
void Lib9pStreamer::onWorkerData(const QByteArray &data)
|
|
{
|
|
if (!m_stopping && !data.isEmpty()) {
|
|
emit dataReady(data);
|
|
}
|
|
}
|
|
|
|
void Lib9pStreamer::onWorkerFinished(bool success, const QString &error)
|
|
{
|
|
// Clean up the worker
|
|
if (m_workerThread) {
|
|
m_workerThread->quit();
|
|
if (m_workerThread->wait(500)) {
|
|
// Thread exited cleanly
|
|
delete m_worker;
|
|
delete m_workerThread;
|
|
} else {
|
|
// Thread still running — abandon safely
|
|
disconnect(m_worker, nullptr, this, nullptr);
|
|
connect(m_workerThread, &QThread::finished, m_worker, &QObject::deleteLater);
|
|
connect(m_workerThread, &QThread::finished, m_workerThread, &QObject::deleteLater);
|
|
}
|
|
m_worker = nullptr;
|
|
m_workerThread = nullptr;
|
|
}
|
|
|
|
if (m_stopping) {
|
|
emit finished();
|
|
return;
|
|
}
|
|
|
|
if (!success && !error.isEmpty()) {
|
|
emit errorOccurred(error);
|
|
}
|
|
|
|
// Decide whether to restart
|
|
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 Lib9pStreamer::scheduleRestart()
|
|
{
|
|
if (m_stopping)
|
|
return;
|
|
|
|
// Re-check guard before restarting
|
|
if (m_policy == Guarded && m_guard && !m_guard()) {
|
|
emit finished();
|
|
return;
|
|
}
|
|
|
|
startWorker();
|
|
}
|