From 8e8466d4efcd8198d15ebb43a6052c8f6166c122 Mon Sep 17 00:00:00 2001 From: Levi Neely Date: Mon, 3 Aug 2026 19:17:40 +0200 Subject: [PATCH] optimize session refresh with pubsub delta events - Parse new session/idx format with all data in one read: name, id, paused, connected, cwd, backend, model, agents - Handle delta events from eventwait instead of full refresh: - session.{sid}.agent.{aid}.state -> update cache, emit signal - Structural changes (new/kill/pause/resume) -> full refresh - Remove per-agent statewait streams (eliminated N processes) - Single eventwait stream handles all events via pubsub This eliminates O(sessions * agents) synchronous reads per refresh, replacing them with one read + incremental delta updates. --- gui/ollie9pclient.cpp | 174 ++++++++++++++++++------------------------ gui/ollie9pclient.h | 6 +- 2 files changed, 75 insertions(+), 105 deletions(-) diff --git a/gui/ollie9pclient.cpp b/gui/ollie9pclient.cpp index 536df36..b615bf6 100644 --- a/gui/ollie9pclient.cpp +++ b/gui/ollie9pclient.cpp @@ -84,15 +84,13 @@ Ollie9pClient::Ollie9pClient(QObject *parent) refreshSessions(); ensureRootDataLoaded(); startAgentConnections(); - reconcileAgentStateStreams(); }); connect(m_daemon, &NinePConnection::disconnected, this, [this]() { setDaemonConnected(false); stopAgentConnections(); - stopAgentStateStreams(); }); - connect(m_daemon, &NinePConnection::readyRead, this, [this](const QByteArray &) { - refreshSessions(); + connect(m_daemon, &NinePConnection::readyRead, this, [this](const QByteArray &data) { + handleEvent(QString::fromUtf8(data).trimmed()); }); refreshSessions(); @@ -288,63 +286,52 @@ void Ollie9pClient::stopAgentConnections() m_agentConnections.clear(); } -void Ollie9pClient::reconcileAgentStateStreams() +// handleEvent processes a single event from eventwait. +// Event format: "topic [payload]" +// Topics: session.{sid}.{action}, session.{sid}.agent.{aid}.{action} +void Ollie9pClient::handleEvent(const QString &eventLine) { - if (!m_daemonConnected) return; - QSet wanted; - for (const QVariant &value : std::as_const(m_sessions)) { - const QVariantMap session = value.toMap(); - const QString sid = session.value("id").toString(); - const bool paused = session.value("paused").toBool(); - for (const QVariant &agentValue : session.value("agents").toList()) { - const QVariantMap agent = agentValue.toMap(); - const QString aid = agent.value("id").toString(); - const QString key = agentKey(sid, aid); - wanted.insert(key); - if (m_agentStateStreams.contains(key)) continue; - if (paused) continue; // Don't start streams for paused sessions - auto *stream = new StreamFsm(StreamFsm::Looping, this); - m_agentStateStreams.insert(key, stream); - m_agentStateValues.insert(key, "idle"); - connect(stream, &StreamFsm::readyRead, this, [this, key, sid, aid](const QByteArray &data) { - QString state = QString::fromUtf8(data).trimmed(); - if (!state.isEmpty() && m_agentStateValues.value(key) != state) { - m_agentStateValues[key] = state; - emit agentStateChanged(sid, aid, state); + if (eventLine.isEmpty()) return; + + // Split into topic and payload + int spaceIdx = eventLine.indexOf(' '); + QString topic = spaceIdx > 0 ? eventLine.left(spaceIdx) : eventLine; + QString payload = spaceIdx > 0 ? eventLine.mid(spaceIdx + 1) : QString(); + + QStringList parts = topic.split('.'); + // Expected: session.{sid}.* or session.{sid}.agent.{aid}.* + + if (parts.size() >= 3 && parts[0] == "session") { + QString sessionId = parts[1]; + QString action = parts[2]; + + if (action == "agent" && parts.size() >= 5) { + // Agent event: session.{sid}.agent.{aid}.{action} + QString agentId = parts[3]; + QString agentAction = parts[4]; + QString key = agentKey(sessionId, agentId); + + if (agentAction == "state") { + // Update cached state + if (!payload.isEmpty() && m_agentStateValues.value(key) != payload) { + m_agentStateValues[key] = payload; + emit agentStateChanged(sessionId, agentId, payload); // Also update activeState if this is the active agent - if (sid == m_activeSessionId && aid == m_agentId) { - m_activeState = state; + if (sessionId == m_activeSessionId && agentId == m_agentId) { + m_activeState = payload; emit activeStateChanged(); } } - }); - stream->start(ollie9pBin(), {"-a", serverAddr(), "read", - "session/" + sid + "/agent/" + aid + "/statewait"}); + } else if (agentAction == "new" || agentAction == "kill") { + // Structural change — need full refresh + refreshSessions(); + } + } else if (action == "pause" || action == "resume" || + action == "new" || action == "kill" || action == "rename") { + // Session-level structural change — need full refresh + refreshSessions(); } } - // Remove streams for agents that no longer exist - for (auto it = m_agentStateStreams.begin(); it != m_agentStateStreams.end();) { - if (wanted.contains(it.key())) { - ++it; - continue; - } - it.value()->disconnect(this); - it.value()->stop(); - it.value()->deleteLater(); - m_agentStateValues.remove(it.key()); - it = m_agentStateStreams.erase(it); - } -} - -void Ollie9pClient::stopAgentStateStreams() -{ - for (StreamFsm *stream : std::as_const(m_agentStateStreams)) { - stream->disconnect(this); - stream->stop(); - } - qDeleteAll(m_agentStateStreams); - m_agentStateStreams.clear(); - m_agentStateValues.clear(); } QString Ollie9pClient::agentState(const QString &sessionId, const QString &agentId) const @@ -356,57 +343,46 @@ void Ollie9pClient::refreshSessions() { if (!m_daemonConnected) return; QByteArray out = run9p({"read", "session/idx"}); - // Don't trim — trailing tabs are significant fields QString raw = QString::fromUtf8(out); m_sessions.clear(); - QSet seen; if (!raw.isEmpty()) { for (const QString &line : raw.split('\n', Qt::SkipEmptyParts)) { QStringList parts = line.split('\t'); - // Skip lines where the first field (session name) is empty - if (parts.isEmpty() || parts[0].isEmpty()) continue; - QString sid = parts[0]; - if (seen.contains(sid)) continue; // already added (multi-agent) - seen.insert(sid); - QVariantMap session; - // session/idx format: name\tstate\tcwd\tbackend\tmodel\tagentName\tid - // The immutable ID is in field 6 (0-indexed), avoiding a round-trip. - QString sessionId = parts.size() > 6 ? parts[6] : ""; - if (sessionId.isEmpty()) sessionId = sid; - session["id"] = sessionId; - session["name"] = sid; - session["state"] = parts.size() > 1 ? parts[1] : ""; - session["cwd"] = parts.size() > 2 ? parts[2] : ""; - session["backend"] = parts.size() > 3 ? parts[3] : ""; - session["model"] = parts.size() > 4 ? parts[4] : ""; - // Check if session is paused - QString pausedStr = QString::fromUtf8(run9p({"read", "session/" + sessionId + "/paused"})).trimmed(); - session["paused"] = (pausedStr == "true"); - // Check RPC connection status - QString connectedStr = QString::fromUtf8(run9p({"read", "session/" + sessionId + "/connected"})).trimmed(); - session["connected"] = (connectedStr == "true"); - QVariantList agents; - for (const QString &name : QString::fromUtf8(run9p({"ls", "session/" + sessionId + "/agent"})).trimmed().split('\n', Qt::SkipEmptyParts)) { - if (name == "new") continue; - const QString path = "session/" + sessionId + "/agent/" + name; - AgentRecord a; - a.id = QString::fromUtf8(run9p({"read", path + "/id"})).trimmed(); - a.name = name; - if (a.id.isEmpty()) a.id = name; - QString agentState = QString::fromUtf8(run9p({"read", path + "/state"})).trimmed(); - if (agentState.isEmpty()) agentState = "idle"; - QVariantMap av; av["id"] = a.id; av["name"] = a.name; - av["state"] = agentState; - av["connected"] = m_daemonConnected && m_agentConnected.value(agentKey(sessionId, a.id), false); + // New format: name\tid\tpaused\tconnected\tcwd\tbackend\tmodel\tagents + // Where agents is semicolon-separated: name:id:state;name:id:state + if (parts.size() < 2) continue; + + QVariantMap session; + session["name"] = parts.value(0); + session["id"] = parts.value(1); + session["paused"] = (parts.value(2) == "true"); + session["connected"] = (parts.value(3) == "true"); + session["cwd"] = parts.value(4); + session["backend"] = parts.value(5); + session["model"] = parts.value(6); + + // Parse agents from semicolon-separated list + QVariantList agents; + QString agentStr = parts.value(7); + if (!agentStr.isEmpty()) { + for (const QString &agentEntry : agentStr.split(';', Qt::SkipEmptyParts)) { + QStringList agentParts = agentEntry.split(':'); + if (agentParts.size() >= 3) { + QVariantMap av; + av["name"] = agentParts.value(0); + av["id"] = agentParts.value(1); + av["state"] = agentParts.value(2); + av["connected"] = session["connected"]; // Agent inherits session connectivity agents.append(av); } - session["agents"] = agents; - m_sessions.append(session); + } + } + session["agents"] = agents; + m_sessions.append(session); } } reconcileAgentConnections(); - reconcileAgentStateStreams(); emit sessionsChanged(); // Check if the currently active session/agent still exist @@ -625,13 +601,9 @@ void Ollie9pClient::switchAgent(const QString &sessionId, const QString &agentId stopStreams(); m_activeSessionId = sessionId; m_agentId = agentId; - m_activeState = "idle"; - - // Read initial state immediately - QByteArray stateOut = run9p({"read", agentPath() + "/state"}); - QString stateStr = QString::fromUtf8(stateOut).trimmed(); - if (!stateStr.isEmpty()) - m_activeState = stateStr; + + // Use cached state from live tracking if available, otherwise default to idle + m_activeState = m_agentStateValues.value(agentKey(sessionId, agentId), "idle"); emit activeSessionIdChanged(); emit activeAgentIdChanged(); diff --git a/gui/ollie9pclient.h b/gui/ollie9pclient.h index cc9eea4..4175dcc 100644 --- a/gui/ollie9pclient.h +++ b/gui/ollie9pclient.h @@ -120,9 +120,8 @@ private: void startAgentConnections(); void startActiveAgentStreams(); void reconcileAgentConnections(); - void reconcileAgentStateStreams(); void stopAgentConnections(); - void stopAgentStateStreams(); + void handleEvent(const QString &eventLine); QVariantList m_sessions; QString m_activeSessionId; @@ -143,8 +142,7 @@ private: QSet m_seenAgentConnections; NinePConnection *m_daemon = nullptr; - // Per-agent state tracking - QHash m_agentStateStreams; + // Agent state cache — updated via delta events QHash m_agentStateValues; // Streaming FSMs — each wraps a QProcess lifecycle