diff --git a/gui/nativestreamer.cpp b/gui/nativestreamer.cpp index c04a42d..00ba7ee 100644 --- a/gui/nativestreamer.cpp +++ b/gui/nativestreamer.cpp @@ -1,17 +1,33 @@ #include "nativestreamer.h" -#include "lib9pclient.h" +#include #include -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) - , 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()) { + if (!m_stopping && !m_path.isEmpty()) { start(m_path); } }); @@ -24,61 +40,40 @@ NativeStreamer::~NativeStreamer() void NativeStreamer::start(const QString &path) { - if (m_running.load()) { + if (m_process) { stop(); } m_path = path; - m_running.store(true); - m_stopping.store(false); + m_stopping = 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); + m_process = new QProcess(this); + connect(m_process, &QProcess::readyReadStandardOutput, this, &NativeStreamer::onReadyRead); + connect(m_process, QOverload::of(&QProcess::finished), + this, &NativeStreamer::onProcessFinished); - // Start worker thread - m_thread = QThread::create([this]() { run(); }); - connect(m_thread, &QThread::finished, this, &NativeStreamer::onThreadFinished); - m_thread->start(); + m_process->start(ollie9pBin(), {"-a", serverAddr(), "read", path}); } void NativeStreamer::stop() { - m_stopping.store(true); - m_running.store(false); + m_stopping = true; 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(); + if (m_process) { + disconnect(m_process, nullptr, this, nullptr); + if (m_process->state() != QProcess::NotRunning) { + m_process->kill(); + m_process->waitForFinished(1000); } - delete m_thread; - m_thread = nullptr; + m_process->deleteLater(); + m_process = nullptr; } } bool NativeStreamer::isRunning() const { - return m_running.load(); + return m_process && m_process->state() == QProcess::Running; } void NativeStreamer::setGuard(std::function guard) @@ -86,18 +81,24 @@ void NativeStreamer::setGuard(std::function guard) m_guard = guard; } -void NativeStreamer::onThreadFinished() +void NativeStreamer::onReadyRead() { - // Thread finished naturally (not via stop()) - clean up and maybe restart - if (m_thread) { - m_thread->deleteLater(); - m_thread = nullptr; + if (m_process) { + QByteArray data = m_process->readAllStandardOutput(); + if (!data.isEmpty()) { + emit dataReady(data); + } } - - // Close any remaining fid - int fid = m_fid.exchange(-1); - if (fid >= 0) { - m_client->closeFid(fid); +} + +void NativeStreamer::onProcessFinished(int exitCode, QProcess::ExitStatus status) +{ + Q_UNUSED(exitCode) + Q_UNUSED(status) + + if (m_process) { + m_process->deleteLater(); + m_process = nullptr; } scheduleRestart(); @@ -105,7 +106,7 @@ void NativeStreamer::onThreadFinished() void NativeStreamer::scheduleRestart() { - if (m_stopping.load()) { + if (m_stopping) { return; } @@ -128,35 +129,3 @@ void NativeStreamer::scheduleRestart() 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); -} diff --git a/gui/nativestreamer.h b/gui/nativestreamer.h index b57bad2..3e63056 100644 --- a/gui/nativestreamer.h +++ b/gui/nativestreamer.h @@ -2,19 +2,19 @@ #define NATIVESTREAMER_H #include -#include +#include #include #include -#include #include -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. - * Replaces QProcess-based streaming with direct library calls. + * Uses ollie-9p subprocess for streaming reads (blocking wait files like + * 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: * Oneshot — runs once, emits finished on EOF/error @@ -27,7 +27,7 @@ class NativeStreamer : public QObject public: enum RestartPolicy { Oneshot, Looping, Guarded }; - explicit NativeStreamer(Lib9pClient *client, RestartPolicy policy, QObject *parent = nullptr); + explicit NativeStreamer(RestartPolicy policy, QObject *parent = nullptr); ~NativeStreamer() override; void start(const QString &path); @@ -38,25 +38,24 @@ public: signals: 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); private slots: - void onThreadFinished(); + void onReadyRead(); + void onProcessFinished(int exitCode, QProcess::ExitStatus status); private: - void run(); void scheduleRestart(); + static QString ollie9pBin(); + static QString serverAddr(); - Lib9pClient *m_client; QString m_path; - QThread *m_thread = nullptr; - std::atomic m_running{false}; - std::atomic m_stopping{false}; - std::atomic m_fid{-1}; + QProcess *m_process = nullptr; RestartPolicy m_policy; std::function m_guard; QTimer *m_restartTimer = nullptr; + bool m_stopping = false; }; #endif // NATIVESTREAMER_H diff --git a/gui/ollie9pclient.cpp b/gui/ollie9pclient.cpp index bf4d751..5cfe0fd 100644 --- a/gui/ollie9pclient.cpp +++ b/gui/ollie9pclient.cpp @@ -15,7 +15,7 @@ Ollie9pClient::Ollie9pClient(QObject *parent) m_9p->connectDefault(); // 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) { if (!data.isEmpty()) emit chatReceived(QString::fromUtf8(data)); @@ -23,7 +23,7 @@ Ollie9pClient::Ollie9pClient(QObject *parent) // State stream — guarded, reads agent statewait. Auto-restarts until // 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]() { 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. - 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) { // Events may come in batches separated by newlines QString text = QString::fromUtf8(data);