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.
This commit is contained in:
parent
eedee30332
commit
8e8466d4ef
|
|
@ -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<QString> 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<QString> 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();
|
||||
|
|
|
|||
|
|
@ -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<QString> m_seenAgentConnections;
|
||||
NinePConnection *m_daemon = nullptr;
|
||||
|
||||
// Per-agent state tracking
|
||||
QHash<QString, StreamFsm *> m_agentStateStreams;
|
||||
// Agent state cache — updated via delta events
|
||||
QHash<QString, QString> m_agentStateValues;
|
||||
|
||||
// Streaming FSMs — each wraps a QProcess lifecycle
|
||||
|
|
|
|||
Reference in New Issue