nativestreamer: use subprocess instead of blocking thread

Go blocking reads can't be interrupted by closing the fid from another
thread. Revert to subprocess-based streaming using ollie-9p, which can
be killed cleanly.

Native library is still used for all request/response operations.
Only streaming (chat, statewait, eventwait) uses subprocess.
This commit is contained in:
Levi Neely 2026-08-07 18:05:53 +02:00
parent 16883248f0
commit 6fda37ae96
3 changed files with 71 additions and 103 deletions

View File

@ -1,17 +1,33 @@
#include "nativestreamer.h" #include "nativestreamer.h"
#include "lib9pclient.h" #include <QStandardPaths>
#include <QDebug> #include <QDebug>
NativeStreamer::NativeStreamer(Lib9pClient *client, RestartPolicy policy, QObject *parent) QString NativeStreamer::ollie9pBin()
{
return QStandardPaths::findExecutable("ollie-9p");
}
QString NativeStreamer::serverAddr()
{
QString ns = qEnvironmentVariable("NAMESPACE");
if (ns.isEmpty()) {
// Plan 9 default: /tmp/ns.$USER.$DISPLAY
QString user = qEnvironmentVariable("USER");
QString display = qEnvironmentVariable("DISPLAY");
ns = "/tmp/ns." + user + "." + display;
}
return "unix!" + ns + "/ollie";
}
NativeStreamer::NativeStreamer(RestartPolicy policy, QObject *parent)
: QObject(parent) : QObject(parent)
, m_client(client)
, m_policy(policy) , m_policy(policy)
{ {
m_restartTimer = new QTimer(this); m_restartTimer = new QTimer(this);
m_restartTimer->setSingleShot(true); m_restartTimer->setSingleShot(true);
m_restartTimer->setInterval(500); // 500ms delay before restart m_restartTimer->setInterval(500); // 500ms delay before restart
connect(m_restartTimer, &QTimer::timeout, this, [this]() { connect(m_restartTimer, &QTimer::timeout, this, [this]() {
if (!m_stopping.load() && !m_path.isEmpty()) { if (!m_stopping && !m_path.isEmpty()) {
start(m_path); start(m_path);
} }
}); });
@ -24,61 +40,40 @@ NativeStreamer::~NativeStreamer()
void NativeStreamer::start(const QString &path) void NativeStreamer::start(const QString &path)
{ {
if (m_running.load()) { if (m_process) {
stop(); stop();
} }
m_path = path; m_path = path;
m_running.store(true); m_stopping = false;
m_stopping.store(false);
// Open the file first (on main thread, before spawning worker) m_process = new QProcess(this);
int fid = m_client->open(path); connect(m_process, &QProcess::readyReadStandardOutput, this, &NativeStreamer::onReadyRead);
if (fid < 0) { connect(m_process, QOverload<int, QProcess::ExitStatus>::of(&QProcess::finished),
emit errorOccurred(m_client->lastError()); this, &NativeStreamer::onProcessFinished);
m_running.store(false);
scheduleRestart();
return;
}
m_fid.store(fid);
// Start worker thread m_process->start(ollie9pBin(), {"-a", serverAddr(), "read", path});
m_thread = QThread::create([this]() { run(); });
connect(m_thread, &QThread::finished, this, &NativeStreamer::onThreadFinished);
m_thread->start();
} }
void NativeStreamer::stop() void NativeStreamer::stop()
{ {
m_stopping.store(true); m_stopping = true;
m_running.store(false);
m_restartTimer->stop(); m_restartTimer->stop();
// Close the fid to unblock any pending read if (m_process) {
int fid = m_fid.exchange(-1); disconnect(m_process, nullptr, this, nullptr);
if (fid >= 0) { if (m_process->state() != QProcess::NotRunning) {
m_client->closeFid(fid); m_process->kill();
} m_process->waitForFinished(1000);
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_process->deleteLater();
m_thread = nullptr; m_process = nullptr;
} }
} }
bool NativeStreamer::isRunning() const bool NativeStreamer::isRunning() const
{ {
return m_running.load(); return m_process && m_process->state() == QProcess::Running;
} }
void NativeStreamer::setGuard(std::function<bool()> guard) void NativeStreamer::setGuard(std::function<bool()> guard)
@ -86,18 +81,24 @@ void NativeStreamer::setGuard(std::function<bool()> guard)
m_guard = guard; m_guard = guard;
} }
void NativeStreamer::onThreadFinished() void NativeStreamer::onReadyRead()
{ {
// Thread finished naturally (not via stop()) - clean up and maybe restart if (m_process) {
if (m_thread) { QByteArray data = m_process->readAllStandardOutput();
m_thread->deleteLater(); if (!data.isEmpty()) {
m_thread = nullptr; emit dataReady(data);
}
} }
}
// Close any remaining fid
int fid = m_fid.exchange(-1); void NativeStreamer::onProcessFinished(int exitCode, QProcess::ExitStatus status)
if (fid >= 0) { {
m_client->closeFid(fid); Q_UNUSED(exitCode)
Q_UNUSED(status)
if (m_process) {
m_process->deleteLater();
m_process = nullptr;
} }
scheduleRestart(); scheduleRestart();
@ -105,7 +106,7 @@ void NativeStreamer::onThreadFinished()
void NativeStreamer::scheduleRestart() void NativeStreamer::scheduleRestart()
{ {
if (m_stopping.load()) { if (m_stopping) {
return; return;
} }
@ -128,35 +129,3 @@ void NativeStreamer::scheduleRestart()
emit finished(); 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);
}

View File

@ -2,19 +2,19 @@
#define NATIVESTREAMER_H #define NATIVESTREAMER_H
#include <QObject> #include <QObject>
#include <QThread> #include <QProcess>
#include <QString> #include <QString>
#include <QTimer> #include <QTimer>
#include <atomic>
#include <functional> #include <functional>
class Lib9pClient;
/** /**
* NativeStreamer - Streaming file reader using native 9P library. * NativeStreamer - Streaming file reader using ollie-9p subprocess.
* *
* Reads from a 9P file in a worker thread, emitting data as it arrives. * Uses ollie-9p subprocess for streaming reads (blocking wait files like
* Replaces QProcess-based streaming with direct library calls. * chat, statewait, eventwait). The native library is used for request/response
* operations, but subprocess is needed for streaming because:
* - Blocking reads in Go can't be interrupted by closing the fid
* - Subprocess can be killed cleanly
* *
* Restart policies: * Restart policies:
* Oneshot — runs once, emits finished on EOF/error * Oneshot — runs once, emits finished on EOF/error
@ -27,7 +27,7 @@ class NativeStreamer : public QObject
public: public:
enum RestartPolicy { Oneshot, Looping, Guarded }; enum RestartPolicy { Oneshot, Looping, Guarded };
explicit NativeStreamer(Lib9pClient *client, RestartPolicy policy, QObject *parent = nullptr); explicit NativeStreamer(RestartPolicy policy, QObject *parent = nullptr);
~NativeStreamer() override; ~NativeStreamer() override;
void start(const QString &path); void start(const QString &path);
@ -38,25 +38,24 @@ public:
signals: signals:
void dataReady(const QByteArray &data); void dataReady(const QByteArray &data);
void finished(); // Only emitted when not restarting (Oneshot, or Guarded with false guard) void finished(); // Only emitted when not restarting
void errorOccurred(const QString &error); void errorOccurred(const QString &error);
private slots: private slots:
void onThreadFinished(); void onReadyRead();
void onProcessFinished(int exitCode, QProcess::ExitStatus status);
private: private:
void run();
void scheduleRestart(); void scheduleRestart();
static QString ollie9pBin();
static QString serverAddr();
Lib9pClient *m_client;
QString m_path; QString m_path;
QThread *m_thread = nullptr; QProcess *m_process = nullptr;
std::atomic<bool> m_running{false};
std::atomic<bool> m_stopping{false};
std::atomic<int> m_fid{-1};
RestartPolicy m_policy; RestartPolicy m_policy;
std::function<bool()> m_guard; std::function<bool()> m_guard;
QTimer *m_restartTimer = nullptr; QTimer *m_restartTimer = nullptr;
bool m_stopping = false;
}; };
#endif // NATIVESTREAMER_H #endif // NATIVESTREAMER_H

View File

@ -15,7 +15,7 @@ Ollie9pClient::Ollie9pClient(QObject *parent)
m_9p->connectDefault(); m_9p->connectDefault();
// Chat stream — one-shot, reads agent chat log. Starts/stops with agent switch. // Chat stream — one-shot, reads agent chat log. Starts/stops with agent switch.
m_chat = new NativeStreamer(m_9p, NativeStreamer::Oneshot, this); m_chat = new NativeStreamer(NativeStreamer::Oneshot, this);
connect(m_chat, &NativeStreamer::dataReady, this, [this](const QByteArray &data) { connect(m_chat, &NativeStreamer::dataReady, this, [this](const QByteArray &data) {
if (!data.isEmpty()) if (!data.isEmpty())
emit chatReceived(QString::fromUtf8(data)); emit chatReceived(QString::fromUtf8(data));
@ -23,7 +23,7 @@ Ollie9pClient::Ollie9pClient(QObject *parent)
// State stream — guarded, reads agent statewait. Auto-restarts until // State stream — guarded, reads agent statewait. Auto-restarts until
// session/agent is deselected or daemon disconnects (guard returns false). // session/agent is deselected or daemon disconnects (guard returns false).
m_state = new NativeStreamer(m_9p, NativeStreamer::Guarded, this); m_state = new NativeStreamer(NativeStreamer::Guarded, this);
m_state->setGuard([this]() { m_state->setGuard([this]() {
return m_daemonConnected && !m_activeSessionId.isEmpty() && !m_agentId.isEmpty(); return m_daemonConnected && !m_activeSessionId.isEmpty() && !m_agentId.isEmpty();
}); });
@ -36,7 +36,7 @@ Ollie9pClient::Ollie9pClient(QObject *parent)
}); });
// Event stream — looping, reads server eventwait for session/agent events. // Event stream — looping, reads server eventwait for session/agent events.
m_event = new NativeStreamer(m_9p, NativeStreamer::Looping, this); m_event = new NativeStreamer(NativeStreamer::Looping, this);
connect(m_event, &NativeStreamer::dataReady, this, [this](const QByteArray &data) { connect(m_event, &NativeStreamer::dataReady, this, [this](const QByteArray &data) {
// Events may come in batches separated by newlines // Events may come in batches separated by newlines
QString text = QString::fromUtf8(data); QString text = QString::fromUtf8(data);