diff --git a/CMakeLists.txt b/CMakeLists.txt index b01d226..061f318 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -91,10 +91,9 @@ if(OLLIE_KF5) gui/main.cpp gui/ollie9pclient.cpp gui/lib9pclient.cpp - gui/ninepconnection.cpp + gui/nativestreamer.cpp gui/chatblockmodel.cpp gui/thememanager.cpp - gui/streamfsm.cpp gui/clipboardhelper.cpp ${GUI_KF5_QRC} ) @@ -161,10 +160,9 @@ else() gui/main.cpp gui/ollie9pclient.cpp gui/lib9pclient.cpp - gui/ninepconnection.cpp + gui/nativestreamer.cpp gui/chatblockmodel.cpp gui/thememanager.cpp - gui/streamfsm.cpp gui/clipboardhelper.cpp ) target_include_directories(ollie-gui PRIVATE diff --git a/gui/lib9pclient.cpp b/gui/lib9pclient.cpp index 924102c..77292bc 100644 --- a/gui/lib9pclient.cpp +++ b/gui/lib9pclient.cpp @@ -243,3 +243,36 @@ bool Lib9pClient::rename(const QString &oldPath, const QString &newName) m_lastError.clear(); return true; } + +int Lib9pClient::open(const QString &path) +{ + if (m_handle <= 0) { + m_lastError = QStringLiteral("not connected"); + return -1; + } + + QByteArray pathUtf8 = path.toUtf8(); + int fid = ollie9p_open(m_handle, pathUtf8.data()); + if (fid < 0) { + updateLastError(); + return -1; + } + m_lastError.clear(); + return fid; +} + +int Lib9pClient::readFid(int fid, char *buf, int bufLen) +{ + int n = ollie9p_read_fid(fid, buf, bufLen); + if (n < 0) { + updateLastError(); + return -1; + } + m_lastError.clear(); + return n; +} + +void Lib9pClient::closeFid(int fid) +{ + ollie9p_close_fid(fid); +} diff --git a/gui/lib9pclient.h b/gui/lib9pclient.h index 4ec48cd..93050bf 100644 --- a/gui/lib9pclient.h +++ b/gui/lib9pclient.h @@ -45,6 +45,11 @@ public: bool mkdir(const QString &path); bool rename(const QString &oldPath, const QString &newName); + // Streaming API - for continuous reads (chat, statewait, etc.) + int open(const QString &path); // Returns fid handle, -1 on error + int readFid(int fid, char *buf, int bufLen); // Returns bytes read, 0 on EOF, -1 on error + void closeFid(int fid); + // Error handling QString lastError() const { return m_lastError; } diff --git a/gui/nativestreamer.cpp b/gui/nativestreamer.cpp new file mode 100644 index 0000000..89d99a1 --- /dev/null +++ b/gui/nativestreamer.cpp @@ -0,0 +1,157 @@ +#include "nativestreamer.h" +#include "lib9pclient.h" +#include + +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) { + m_thread->quit(); + m_thread->wait(1000); + if (m_thread->isRunning()) { + 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 guard) +{ + m_guard = guard; +} + +void NativeStreamer::onThreadFinished() +{ + if (m_thread) { + delete m_thread; + 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); +} diff --git a/gui/nativestreamer.h b/gui/nativestreamer.h new file mode 100644 index 0000000..b57bad2 --- /dev/null +++ b/gui/nativestreamer.h @@ -0,0 +1,62 @@ +#ifndef NATIVESTREAMER_H +#define NATIVESTREAMER_H + +#include +#include +#include +#include +#include +#include + +class Lib9pClient; + +/** + * NativeStreamer - Streaming file reader using native 9P library. + * + * Reads from a 9P file in a worker thread, emitting data as it arrives. + * Replaces QProcess-based streaming with direct library calls. + * + * Restart policies: + * Oneshot — runs once, emits finished on EOF/error + * Looping — always restarts after EOF/error + * Guarded — restarts only if guard() returns true + */ +class NativeStreamer : public QObject +{ + Q_OBJECT +public: + enum RestartPolicy { Oneshot, Looping, Guarded }; + + explicit NativeStreamer(Lib9pClient *client, RestartPolicy policy, QObject *parent = nullptr); + ~NativeStreamer() override; + + void start(const QString &path); + void stop(); + bool isRunning() const; + + void setGuard(std::function guard); + +signals: + void dataReady(const QByteArray &data); + void finished(); // Only emitted when not restarting (Oneshot, or Guarded with false guard) + void errorOccurred(const QString &error); + +private slots: + void onThreadFinished(); + +private: + void run(); + void scheduleRestart(); + + 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}; + RestartPolicy m_policy; + std::function m_guard; + QTimer *m_restartTimer = nullptr; +}; + +#endif // NATIVESTREAMER_H diff --git a/gui/ollie9pclient.cpp b/gui/ollie9pclient.cpp index 29ffc19..bf4d751 100644 --- a/gui/ollie9pclient.cpp +++ b/gui/ollie9pclient.cpp @@ -7,52 +7,27 @@ #include #include -static QString ninepBin() -{ - // plan9port's 9p command - QString bin = QStandardPaths::findExecutable("9p"); - if (!bin.isEmpty()) return bin; - QString plan9 = qEnvironmentVariable("PLAN9"); - if (!plan9.isEmpty() && QFileInfo::exists(plan9 + "/bin/9p")) - return plan9 + "/bin/9p"; - return "9p"; -} - -static QString 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"; -} - -// Returns the path to ollie-9p binary, or empty string if not found. -static QString ollie9pBin() -{ - return QStandardPaths::findExecutable("ollie-9p"); -} - Ollie9pClient::Ollie9pClient(QObject *parent) : QObject(parent) { + // Initialize native 9P client + m_9p = new Lib9pClient(this); + m_9p->connectDefault(); + // Chat stream — one-shot, reads agent chat log. Starts/stops with agent switch. - m_chat = new StreamFsm(StreamFsm::Oneshot, this); - connect(m_chat, &StreamFsm::readyRead, this, [this](const QByteArray &data) { + m_chat = new NativeStreamer(m_9p, NativeStreamer::Oneshot, this); + connect(m_chat, &NativeStreamer::dataReady, this, [this](const QByteArray &data) { if (!data.isEmpty()) emit chatReceived(QString::fromUtf8(data)); }); // State stream — guarded, reads agent statewait. Auto-restarts until // session/agent is deselected or daemon disconnects (guard returns false). - m_state = new StreamFsm(StreamFsm::Guarded, this); + m_state = new NativeStreamer(m_9p, NativeStreamer::Guarded, this); m_state->setGuard([this]() { return m_daemonConnected && !m_activeSessionId.isEmpty() && !m_agentId.isEmpty(); }); - connect(m_state, &StreamFsm::readyRead, this, [this](const QByteArray &data) { + connect(m_state, &NativeStreamer::dataReady, this, [this](const QByteArray &data) { QString state = QString::fromUtf8(data).trimmed(); if (!state.isEmpty()) { m_activeState = state; @@ -60,40 +35,36 @@ Ollie9pClient::Ollie9pClient(QObject *parent) } }); - // Event stream — reads eventwait for server events. - // Note: NinePConnection handles retry/reconnect but we use heartbeat for disconnect detection. - m_daemon = new NinePConnection(this); - connect(m_daemon, &NinePConnection::connected, this, [this]() { - // Initialize native 9P client for fast operations - if (!m_9p) { - m_9p = new Lib9pClient(this); + // Event stream — looping, reads server eventwait for session/agent events. + m_event = new NativeStreamer(m_9p, 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); + for (const QString &line : text.split('\n', Qt::SkipEmptyParts)) { + handleEvent(line.trimmed()); } - if (!m_9p->isConnected()) { - m_9p->connectDefault(); - } - // Heartbeat will set daemonConnected on first successful ping - heartbeat(); - }); - connect(m_daemon, &NinePConnection::readyRead, this, [this](const QByteArray &data) { - handleEvent(QString::fromUtf8(data).trimmed()); }); - // Heartbeat timer — detect daemon death by periodic reads + // Heartbeat timer — detects daemon death m_heartbeatTimer = new QTimer(this); m_heartbeatTimer->setInterval(2500); // 2.5 seconds connect(m_heartbeatTimer, &QTimer::timeout, this, &Ollie9pClient::heartbeat); m_heartbeatTimer->start(); - refreshSessions(); - ensureRootDataLoaded(); - m_daemon->start(ollie9pBin(), {"-a", serverAddr(), "read", "--open-marker", "eventwait"}); + // Initial connection check and data load + heartbeat(); + if (m_daemonConnected) { + refreshSessions(); + ensureRootDataLoaded(); + startEventStream(); + } } Ollie9pClient::~Ollie9pClient() { if (m_heartbeatTimer) m_heartbeatTimer->stop(); stopAgentStreams(); - if (m_daemon) m_daemon->stop(); + if (m_event) m_event->stop(); } void Ollie9pClient::setActiveSessionId(const QString &id) @@ -134,9 +105,6 @@ void Ollie9pClient::setDaemonConnected(bool connected) void Ollie9pClient::heartbeat() { // Try to read a cheap file to verify daemon is alive - if (!m_9p) { - m_9p = new Lib9pClient(this); - } if (!m_9p->isConnected()) { m_9p->connectDefault(); } @@ -145,17 +113,19 @@ void Ollie9pClient::heartbeat() bool alive = !result.isEmpty(); if (alive && !m_daemonConnected) { - // Daemon came up — refresh state + // Daemon came up — refresh state and start event stream setDaemonConnected(true); refreshSessions(); ensureRootDataLoaded(); + startEventStream(); if (!m_activeSessionId.isEmpty() && !m_agentId.isEmpty()) { startActiveAgentStreams(); } } else if (!alive && m_daemonConnected) { - // Daemon died — clear state + // Daemon died — clear state and stop streams setDaemonConnected(false); stopAgentStreams(); + if (m_event) m_event->stop(); m_sessions.clear(); m_activeSessionId.clear(); m_agentId.clear(); @@ -164,12 +134,17 @@ void Ollie9pClient::heartbeat() emit activeSessionIdChanged(); emit activeAgentIdChanged(); emit activeStateChanged(); - if (m_9p) { - m_9p->disconnect(); - } + m_9p->disconnect(); } } +void Ollie9pClient::startEventStream() +{ + if (!m_daemonConnected) return; + if (m_event->isRunning()) return; + m_event->start(QStringLiteral("eventwait")); +} + QString Ollie9pClient::sessionConnectionColor(const QString &sessionId) const { for (const QVariant &v : m_sessions) { @@ -283,14 +258,9 @@ QString Ollie9pClient::agentState(const QString &sessionId, const QString &agent void Ollie9pClient::refreshSessions() { if (!m_daemonConnected) return; + if (!m_9p || !m_9p->isConnected()) return; - // Use native 9P client if available (fast path) - QByteArray out; - if (m_9p && m_9p->isConnected()) { - out = m_9p->read(QStringLiteral("session/idx")); - } else { - out = run9p({"read", "session/idx"}); - } + QByteArray out = m_9p->read(QStringLiteral("session/idx")); QString raw = QString::fromUtf8(out); m_sessions.clear(); @@ -391,21 +361,17 @@ void Ollie9pClient::refreshSessions() QString Ollie9pClient::readLog() { if (m_activeSessionId.isEmpty()) return {}; - if (m_9p && m_9p->isConnected()) { - return QString::fromUtf8(m_9p->read(agentPath() + "/log")); - } - return QString::fromUtf8(run9p({"read", agentPath() + "/log"})); + if (!m_9p || !m_9p->isConnected()) return {}; + return QString::fromUtf8(m_9p->read(agentPath() + "/log")); } QString Ollie9pClient::readLogForSession(const QString &sessionId, const QString &agentId) { if (sessionId.isEmpty() || agentId.isEmpty()) return {}; + if (!m_9p || !m_9p->isConnected()) return {}; // Use immutable IDs — the 9P namespace resolves them via aliases. QString path = "session/" + sessionId + "/agent/" + agentId + "/log"; - if (m_9p && m_9p->isConnected()) { - return QString::fromUtf8(m_9p->read(path)); - } - return QString::fromUtf8(run9p({"read", path})); + return QString::fromUtf8(m_9p->read(path)); } bool Ollie9pClient::submit(const QString &prompt) @@ -487,25 +453,17 @@ bool Ollie9pClient::isSessionPaused(const QString &sessionId) const QString Ollie9pClient::getConfig() { if (m_activeSessionId.isEmpty()) return {}; - if (m_9p && m_9p->isConnected()) { - return QString::fromUtf8(m_9p->read(agentPath() + "/cfg")); - } - return QString::fromUtf8(run9p({"read", agentPath() + "/cfg"})); + if (!m_9p || !m_9p->isConnected()) return {}; + return QString::fromUtf8(m_9p->read(agentPath() + "/cfg")); } QStringList Ollie9pClient::getAgents(const QString &sessionId) { if (sessionId.isEmpty()) return {}; + if (!m_9p || !m_9p->isConnected()) return {}; // Use immutable ID — the 9P namespace resolves it via alias. QString path = "session/" + sessionId + "/agent"; - QStringList all; - if (m_9p && m_9p->isConnected()) { - all = m_9p->ls(path); - } else { - QString raw = QString::fromUtf8(run9p({"ls", path})).trimmed(); - if (!raw.isEmpty()) - all = raw.split('\n', Qt::SkipEmptyParts); - } + QStringList all = m_9p->ls(path); QStringList agents; for (const QString &a : all) { QString name = a.trimmed(); @@ -548,12 +506,8 @@ void Ollie9pClient::switchAgent(const QString &sessionId, const QString &agentId void Ollie9pClient::loadRootBackends() { if (m_rootBackendsLoaded) return; - QByteArray out; - if (m_9p && m_9p->isConnected()) { - out = m_9p->read(QStringLiteral("backends")); - } else { - out = run9p({"read", "backends"}); - } + if (!m_9p || !m_9p->isConnected()) return; + QByteArray out = m_9p->read(QStringLiteral("backends")); QString raw = QString::fromUtf8(out).trimmed(); if (!raw.isEmpty()) { m_availableBackends = raw.split('\n', Qt::SkipEmptyParts); @@ -566,12 +520,8 @@ void Ollie9pClient::loadRootBackends() void Ollie9pClient::loadRootAgents() { if (m_rootAgentsLoaded) return; - QByteArray out; - if (m_9p && m_9p->isConnected()) { - out = m_9p->read(QStringLiteral("agents")); - } else { - out = run9p({"read", "agents"}); - } + if (!m_9p || !m_9p->isConnected()) return; + QByteArray out = m_9p->read(QStringLiteral("agents")); QString raw = QString::fromUtf8(out).trimmed(); if (!raw.isEmpty()) { m_availableAgents = raw.split('\n', Qt::SkipEmptyParts); @@ -692,8 +642,8 @@ void Ollie9pClient::ensureRootDataLoaded() { if (!m_rootBackendsLoaded) loadRootBackends(); if (!m_rootAgentsLoaded) loadRootAgents(); - if (!m_rootModelsLoaded) { - QByteArray out = run9p({"read", "models"}); + if (!m_rootModelsLoaded && m_9p && m_9p->isConnected()) { + QByteArray out = m_9p->read(QStringLiteral("models")); QString raw = QString::fromUtf8(out); if (!raw.isEmpty()) { for (const QString &line : raw.split('\n')) { @@ -721,10 +671,8 @@ void Ollie9pClient::startActiveAgentStreams() if (!m_daemonConnected) return; if (m_activeSessionId.isEmpty() || m_agentId.isEmpty()) return; - const QString srv = serverAddr(); - const QString bin = ninepBin(); - m_chat->start(bin, {"-a", srv, "read", agentPath() + "/chat"}); - m_state->start(bin, {"-a", srv, "read", agentPath() + "/statewait"}); + m_chat->start(agentPath() + "/chat"); + m_state->start(agentPath() + "/statewait"); } void Ollie9pClient::stopAgentStreams() @@ -747,22 +695,3 @@ QString Ollie9pClient::agentPath() const // Use immutable IDs — the 9P namespace resolves them via aliases. return "session/" + m_activeSessionId + "/agent/" + m_agentId; } - -QByteArray Ollie9pClient::run9p(const QStringList &args) -{ - QProcess proc; - QStringList fullArgs = {"-a", serverAddr()}; - fullArgs.append(args); - proc.start(ninepBin(), fullArgs); - - // Drain output while waiting. The log file can be up to 64KB; if the - // parent doesn't read, the child blocks on a full pipe and never exits, - // causing "QProcess: Destroyed while process is still running". - QByteArray out; - while (proc.state() != QProcess::NotRunning) { - if (proc.waitForReadyRead(100)) - out += proc.readAllStandardOutput(); - } - out += proc.readAllStandardOutput(); - return out; -} diff --git a/gui/ollie9pclient.h b/gui/ollie9pclient.h index 6cf4c51..9dc113e 100644 --- a/gui/ollie9pclient.h +++ b/gui/ollie9pclient.h @@ -10,12 +10,11 @@ #include #endif -#include "streamfsm.h" -#include "ninepconnection.h" #include "lib9pclient.h" +#include "nativestreamer.h" -// 9P-based client for ollie. Uses ollie-9p subprocess for transport. -// Streaming files (chat, statewait) use StreamFsm instances for lifecycle. +// 9P-based client for ollie. Uses native libollie9p for all operations. +// Streaming files (chat, statewait, eventwait) use NativeStreamer threads. class Ollie9pClient : public QObject { Q_OBJECT @@ -106,10 +105,10 @@ private: void stopAgentStreams(); void stopStreams(); QString agentPath() const; - QByteArray run9p(const QStringList &args); // Fallback for streaming (uses subprocess) void ensureRootDataLoaded(); void setDaemonConnected(bool connected); void startActiveAgentStreams(); + void startEventStream(); void handleEvent(const QString &eventLine); void heartbeat(); @@ -133,15 +132,14 @@ private: QVariantMap m_rootModels; QString m_currentBackend; bool m_daemonConnected = false; - NinePConnection *m_daemon = nullptr; // Agent state cache — updated via delta events QHash m_agentStateValues; - // Streaming FSMs — each wraps a QProcess lifecycle - StreamFsm *m_chat = nullptr; // Oneshot — reads agent chat log - StreamFsm *m_state = nullptr; // Guarded — reads agent statewait - StreamFsm *m_event = nullptr; // Looping — reads server eventwait + // Native streaming readers (replace StreamFsm/NinePConnection) + NativeStreamer *m_chat = nullptr; // Oneshot — reads agent chat log + NativeStreamer *m_state = nullptr; // Guarded — reads agent statewait + NativeStreamer *m_event = nullptr; // Looping — reads server eventwait }; #endif // OLLIE9PCLIENT_H \ No newline at end of file