ollie/kde/gui/lib9pstreamer.cpp

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();
}