refactor: promote remote to session level, remove agent connectivity, formalize NinePConnection state machine

- Add createSession(name, remote) with New Session dialog
- Remove agent connectivity code (redundant with session connectivity)
- Refactor NinePConnection as formal state machine with explicit events
- Fix silent restart to not cycle through Connecting state
This commit is contained in:
Levi Neely 2026-08-04 15:22:08 +02:00
parent ec89eaae18
commit 265a1fdda9
5 changed files with 242 additions and 205 deletions

View File

@ -580,6 +580,81 @@ ApplicationWindow {
}
}
// --- New Session Dialog ---
Dialog {
id: newSessionDialog
title: "New Session"
anchors.centerIn: parent
width: 380
modal: true
contentItem: ColumnLayout {
spacing: 16
Layout.margins: 16
Label {
id: newSessionErrorLabel
Layout.fillWidth: true
visible: text !== ""
color: "red"
wrapMode: Text.Wrap
}
// Name
RowLayout {
Layout.fillWidth: true
Label { text: "Name:"; Layout.preferredWidth: 80 }
TextField {
id: sessionNameField
Layout.fillWidth: true
placeholderText: "(auto-generated)"
}
}
// Remote
RowLayout {
Layout.fillWidth: true
Label { text: "Remote:"; Layout.preferredWidth: 80 }
TextField {
id: sessionRemoteField
Layout.fillWidth: true
placeholderText: "user@host (optional)"
}
}
}
footer: RowLayout {
Layout.fillWidth: true
Item { Layout.fillWidth: true }
Button {
text: "OK"
onClicked: {
newSessionErrorLabel.text = ""
var name = sessionNameField.text.trim()
var remote = sessionRemoteField.text.trim()
var err = ollie.createSession(name, remote)
if (err === "") {
ollie.refreshSessions()
newSessionDialog.close()
} else {
newSessionErrorLabel.text = err
}
}
}
Button {
text: "Cancel"
onClicked: newSessionDialog.close()
}
}
onOpened: {
newSessionErrorLabel.text = ""
sessionNameField.text = ""
sessionRemoteField.text = ""
}
}
// --- New Agent Dialog ---
Dialog {
id: newAgentDialog
@ -771,10 +846,7 @@ ApplicationWindow {
Label { text: "Sessions"; Layout.fillWidth: true; font.bold: true; leftPadding: 8 }
ToolButton {
text: "\u2795"
onClicked: {
if (ollie.createQuickSession())
ollie.refreshSessions()
}
onClicked: newSessionDialog.open()
ToolTip.text: "New Session"
ToolTip.visible: hovered
}

View File

@ -1,10 +1,27 @@
#include "ninepconnection.h"
/*
* TRANSITION TABLE
* ================
*
* State | Event | Guard | Next State | Actions
* -------------|--------------|------------------------|--------------|------------------
* Disconnected | EvStart | | Connecting | spawnProcess
* Connecting | EvOpenMarker | | Connected | emit connected
* Connecting | EvExitFail | canRetry | RetryWait | scheduleRetry
* Connecting | EvExitFail | !canRetry | Disconnected | emit disconnected
* Connected | EvExitOk | | Connected | scheduleRetry (silent)
* Connected | EvExitFail | canRetry | RetryWait | scheduleRetry
* Connected | EvExitFail | !canRetry | Disconnected | emit disconnected
* RetryWait | EvRetryTimer | | Connecting | spawnProcess
* Any | EvStop | | Disconnected | killProcess, emit disconnected
*/
NinePConnection::NinePConnection(QObject *parent)
: QObject(parent)
{
m_retryTimer.setSingleShot(true);
connect(&m_retryTimer, &QTimer::timeout, this, &NinePConnection::retry);
connect(&m_retryTimer, &QTimer::timeout, this, [this]() { dispatch(EvRetryTimer); });
}
NinePConnection::~NinePConnection()
@ -12,24 +29,113 @@ NinePConnection::~NinePConnection()
stop();
}
void NinePConnection::start(const QString &program, const QStringList &args, int retryMs, bool probe, int maxRetries)
void NinePConnection::start(const QString &program, const QStringList &args, int retryMs, int maxRetries)
{
m_program = program;
m_args = args;
m_retryMs = retryMs;
m_probe = probe;
m_maxRetries = maxRetries;
m_retryCount = 0;
m_stopRequested = false;
m_openMarkerSeen = false;
m_retryTimer.stop();
if (!m_process)
startProcess();
dispatch(EvStart);
}
void NinePConnection::stop()
{
m_stopRequested = true;
dispatch(EvStop);
}
void NinePConnection::dispatch(Event ev)
{
State next = m_state;
switch (m_state) {
case Disconnected:
if (ev == EvStart) {
next = Connecting;
spawnProcess();
}
break;
case Connecting:
if (ev == EvStop) {
next = Disconnected;
killProcess();
} else if (ev == EvOpenMarker) {
next = Connected;
} else if (ev == EvExitFail) {
bool canRetry = m_retryMs > 0 && (m_maxRetries == 0 || m_retryCount < m_maxRetries);
if (canRetry) {
next = RetryWait;
++m_retryCount;
scheduleRetry();
} else {
next = Disconnected;
}
}
break;
case Connected:
if (ev == EvStop) {
next = Disconnected;
killProcess();
} else if (ev == EvExitOk) {
// Silent restart - stay Connected
m_retryCount = 0;
scheduleRetry();
} else if (ev == EvExitFail) {
bool canRetry = m_retryMs > 0 && (m_maxRetries == 0 || m_retryCount < m_maxRetries);
if (canRetry) {
next = RetryWait;
++m_retryCount;
scheduleRetry();
} else {
next = Disconnected;
}
}
break;
case RetryWait:
if (ev == EvStop) {
next = Disconnected;
m_retryTimer.stop();
} else if (ev == EvRetryTimer) {
next = Connecting;
spawnProcess();
}
break;
}
if (next != m_state)
enter(next);
}
void NinePConnection::enter(State state)
{
State prev = m_state;
m_state = state;
emit stateChanged(state);
if (state == Connected && prev != Connected)
emit connected();
else if (state == Disconnected && prev != Disconnected)
emit disconnected();
}
void NinePConnection::spawnProcess()
{
if (m_process) return;
m_stderr.clear();
m_openMarkerSeen = false;
m_process = new QProcess(this);
connect(m_process, &QProcess::readyReadStandardOutput, this, &NinePConnection::onStdout);
connect(m_process, &QProcess::readyReadStandardError, this, &NinePConnection::onStderr);
connect(m_process, QOverload<int, QProcess::ExitStatus>::of(&QProcess::finished),
this, &NinePConnection::onFinished);
m_process->start(m_program, m_args);
}
void NinePConnection::killProcess()
{
m_retryTimer.stop();
if (m_process) {
m_process->disconnect(this);
@ -38,32 +144,14 @@ void NinePConnection::stop()
m_process->deleteLater();
m_process = nullptr;
}
setState(Disconnected);
}
void NinePConnection::startProcess()
void NinePConnection::scheduleRetry()
{
if (m_stopRequested || m_process)
return;
m_stderr.clear();
m_openMarkerSeen = false;
m_process = new QProcess(this);
connect(m_process, &QProcess::started, this, &NinePConnection::onStarted);
connect(m_process, &QProcess::readyReadStandardOutput, this, &NinePConnection::onStdout);
connect(m_process, &QProcess::readyReadStandardError, this, &NinePConnection::onStderr);
connect(m_process, QOverload<int, QProcess::ExitStatus>::of(&QProcess::finished),
this, &NinePConnection::onFinished);
// A finite probe refreshes an already connected agent. Keep the
// established state while the replacement probe is opening so a
// successful periodic probe does not flash disconnected in the UI.
if (!(m_probe && m_state == Connected))
setState(Connecting);
m_process->start(m_program, m_args);
}
void NinePConnection::onStarted()
{
// Process started - wait for OLLIE_9P_OPEN marker on stderr
if (m_retryMs > 0)
m_retryTimer.start(m_retryMs);
else
spawnProcess();
}
void NinePConnection::onStdout()
@ -77,13 +165,13 @@ void NinePConnection::onStderr()
{
m_stderr += m_process->readAllStandardError();
while (true) {
const int end = m_stderr.indexOf('\n');
int end = m_stderr.indexOf('\n');
if (end < 0) break;
const QByteArray line = m_stderr.left(end).trimmed();
QByteArray line = m_stderr.left(end).trimmed();
m_stderr.remove(0, end + 1);
if (line == "OLLIE_9P_OPEN") {
m_openMarkerSeen = true;
setState(Connected);
dispatch(EvOpenMarker);
}
}
}
@ -93,42 +181,7 @@ void NinePConnection::onFinished(int exitCode, QProcess::ExitStatus)
if (!m_process) return;
m_process->deleteLater();
m_process = nullptr;
// Successful probe that was connected: schedule re-probe (reset retry count).
if (m_probe && m_state == Connected && m_openMarkerSeen && exitCode == 0 && !m_stopRequested) {
m_retryCount = 0;
if (m_retryMs > 0)
m_retryTimer.start(m_retryMs);
else
startProcess(); // immediate re-probe
return;
}
// No retry configured, or stop requested, or max retries exceeded.
if (m_stopRequested || m_retryMs <= 0 || (m_maxRetries > 0 && m_retryCount >= m_maxRetries)) {
setState(Disconnected);
return;
}
++m_retryCount;
setState(RetryWait);
m_retryTimer.start(m_retryMs);
}
void NinePConnection::retry()
{
if (!m_stopRequested)
startProcess();
}
void NinePConnection::setState(State state)
{
if (m_state == state) return;
const State previous = m_state;
m_state = state;
emit stateChanged(state);
if (state == Connected) {
emit connected();
if (m_probe && m_process)
m_process->closeWriteChannel();
}
else if (state != Connected && previous != Disconnected)
emit disconnected();
bool success = (exitCode == 0 && m_openMarkerSeen);
dispatch(success ? EvExitOk : EvExitFail);
}

View File

@ -6,6 +6,11 @@
#include <QTimer>
#include <QStringList>
/**
* NinePConnection - Resilient 9P subprocess connection.
*
* State machine with explicit events and transition table.
*/
class NinePConnection : public QObject
{
Q_OBJECT
@ -13,10 +18,13 @@ public:
enum State { Disconnected, Connecting, Connected, RetryWait };
Q_ENUM(State)
enum Event { EvStart, EvStop, EvOpenMarker, EvExitOk, EvExitFail, EvRetryTimer };
Q_ENUM(Event)
explicit NinePConnection(QObject *parent = nullptr);
~NinePConnection() override;
void start(const QString &program, const QStringList &args, int retryMs = 4000, bool probe = false, int maxRetries = 0);
void start(const QString &program, const QStringList &args, int retryMs = 4000, int maxRetries = 0);
void stop();
State state() const { return m_state; }
@ -27,15 +35,16 @@ signals:
void stateChanged(State state);
private slots:
void onStarted();
void onStdout();
void onStderr();
void onFinished(int exitCode, QProcess::ExitStatus status);
void retry();
private:
void setState(State state);
void startProcess();
void dispatch(Event ev);
void enter(State state);
void spawnProcess();
void killProcess();
void scheduleRetry();
QProcess *m_process = nullptr;
QTimer m_retryTimer;
@ -43,12 +52,10 @@ private:
QStringList m_args;
QByteArray m_stderr;
State m_state = Disconnected;
bool m_stopRequested = false;
int m_retryMs = 4000;
bool m_probe = false;
bool m_openMarkerSeen = false;
int m_maxRetries = 0; // 0 = unlimited
int m_maxRetries = 0;
int m_retryCount = 0;
bool m_openMarkerSeen = false;
};
#endif // NINEPCONNECTION_H

View File

@ -46,18 +46,6 @@ Ollie9pClient::Ollie9pClient(QObject *parent)
emit chatReceived(QString::fromUtf8(data));
});
connect(m_chat, &StreamFsm::died, this, [this]() {
if (!m_activeSessionId.isEmpty()) {
for (const QVariant &value : std::as_const(m_sessions)) {
const QVariantMap session = value.toMap();
if (session.value("id").toString() != m_activeSessionId) continue;
for (const QVariant &agentValue : session.value("agents").toList()) {
const QString key = agentKey(m_activeSessionId,
agentValue.toMap().value("id").toString());
setAgentConnected(key, false);
}
break;
}
}
emit streamingDone();
});
@ -90,14 +78,17 @@ Ollie9pClient::Ollie9pClient(QObject *parent)
}
refreshSessions();
ensureRootDataLoaded();
startAgentConnections();
// Restart streams for the active agent if we have one
if (!m_activeSessionId.isEmpty() && !m_agentId.isEmpty()) {
startActiveAgentStreams();
emit sessionConnected();
}
});
connect(m_daemon, &NinePConnection::disconnected, this, [this]() {
setDaemonConnected(false);
if (m_9p) {
m_9p->disconnect();
}
stopAgentConnections();
});
connect(m_daemon, &NinePConnection::readyRead, this, [this](const QByteArray &data) {
handleEvent(QString::fromUtf8(data).trimmed());
@ -112,7 +103,6 @@ Ollie9pClient::~Ollie9pClient()
{
stopAgentStreams();
if (m_daemon) m_daemon->stop();
stopAgentConnections();
}
void Ollie9pClient::setActiveSessionId(const QString &id)
@ -146,19 +136,10 @@ void Ollie9pClient::setDaemonConnected(bool connected)
{
if (m_daemonConnected == connected) return;
m_daemonConnected = connected;
if (!connected) {
for (auto it = m_agentConnected.begin(); it != m_agentConnected.end(); ++it)
it.value() = false;
}
emit daemonConnectedChanged();
emit sessionsChanged();
}
bool Ollie9pClient::agentConnected(const QString &sessionId, const QString &agentId) const
{
return m_daemonConnected && m_agentConnected.value(agentKey(sessionId, agentId), false);
}
QString Ollie9pClient::sessionConnectionColor(const QString &sessionId) const
{
for (const QVariant &v : m_sessions) {
@ -179,83 +160,6 @@ QString Ollie9pClient::sessionConnectionColor(const QString &sessionId) const
return "gray";
}
void Ollie9pClient::setAgentConnected(const QString &key, bool connected)
{
const bool cached = m_agentConnected.value(key, false);
m_agentConnected[key] = connected;
bool modelChanged = false;
bool foundInModel = false;
// Update the in-memory model without starting another round of probes.
for (QVariant &value : m_sessions) {
QVariantMap session = value.toMap();
QVariantList agents = session.value("agents").toList();
for (QVariant &agentValue : agents) {
QVariantMap agent = agentValue.toMap();
if (agentKey(session["id"].toString(), agent["id"].toString()) == key) {
foundInModel = true;
if (agent.value("connected").toBool() != connected)
modelChanged = true;
agent["connected"] = connected;
}
agentValue = agent;
}
session["agents"] = agents;
value = session;
}
// Ignore stale disconnects for agents no longer in the model (killed).
if (!foundInModel && !connected)
return;
if (cached != connected || modelChanged)
emit sessionsChanged();
}
void Ollie9pClient::startAgentConnections()
{
// Agent connection state now comes from session/idx (refreshSessions)
// and eventwait events. No per-agent processes needed.
// Just start streams for the active agent if we have one.
if (!m_activeSessionId.isEmpty() && !m_agentId.isEmpty()) {
startActiveAgentStreams();
}
}
// reconcileAgentConnections updates agent connected state from m_sessions.
// Connection state is derived from session.connected in session/idx.
void Ollie9pClient::reconcileAgentConnections()
{
if (!m_daemonConnected) return;
// Update connected state for all agents based on session data
for (const QVariant &value : std::as_const(m_sessions)) {
const QVariantMap session = value.toMap();
const QString sid = session.value("id").toString();
const bool isConnected = session.value("connected").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);
bool wasConnected = m_agentConnected.value(key, false);
m_agentConnected[key] = isConnected;
// Emit signal if this is the active agent and connection state changed
if (key == agentKey(m_activeSessionId, m_agentId)) {
if (!wasConnected && isConnected) {
emit sessionConnected();
startActiveAgentStreams();
}
}
}
}
}
void Ollie9pClient::stopAgentConnections()
{
// No per-agent connections to stop anymore
m_agentConnected.clear();
}
// handleEvent processes a single event from eventwait.
// Event format: "topic payload" where topic is "session.{sid}.{action}" or
// "session.{sid}.agent.{aid}.{action}"
@ -364,8 +268,6 @@ void Ollie9pClient::refreshSessions()
// Track if active session/agent still exist during parsing
bool activeSessionFound = false;
bool activeAgentFound = false;
bool activeAgentConnected = false;
bool activeAgentWasConnected = m_agentConnected.value(agentKey(m_activeSessionId, m_agentId), false);
if (!raw.isEmpty()) {
for (const QString &line : raw.split('\n', Qt::SkipEmptyParts)) {
@ -406,14 +308,9 @@ void Ollie9pClient::refreshSessions()
av["connected"] = isConnected;
agents.append(av);
// Update agent connected cache inline
const QString key = agentKey(sessionId, agentId);
m_agentConnected[key] = isConnected;
// Check if this is the active agent
if (isActiveSession && agentId == m_agentId) {
activeAgentFound = true;
activeAgentConnected = isConnected;
}
}
}
@ -423,12 +320,6 @@ void Ollie9pClient::refreshSessions()
}
}
// Handle active agent connection state change
if (activeAgentFound && !activeAgentWasConnected && activeAgentConnected) {
emit sessionConnected();
startActiveAgentStreams();
}
emit sessionsChanged();
// Handle active session/agent removal
@ -700,6 +591,25 @@ bool Ollie9pClient::createQuickSession()
return proc.exitCode() == 0;
}
QString Ollie9pClient::createSession(const QString &name, const QString &remote)
{
QStringList args;
if (!name.isEmpty()) args << "name=" + name;
if (!remote.isEmpty()) args << "remote=" + remote;
QProcess proc;
proc.start(ninepBin(), {"-a", serverAddr(), "rdwr", "session/new"});
proc.waitForStarted(3000);
proc.write((args.join(" ") + "\n").toUtf8());
proc.closeWriteChannel();
proc.waitForFinished(5000);
if (proc.exitCode() != 0) {
QString err = QString::fromUtf8(proc.readAllStandardError()).trimmed();
return err.isEmpty() ? QStringLiteral("failed to create session") : err;
}
return QString(); // success
}
QString Ollie9pClient::createAgent(const QString &sessionId, const QString &cwd, const QString &backend, const QString &model, const QString &agent, const QString &remote, const QString &agentAlias)
{
if (sessionId.isEmpty() || cwd.isEmpty()) return QStringLiteral("session ID and directory are required");

View File

@ -52,7 +52,6 @@ public:
bool rootModelsLoaded() const { return m_rootModelsLoaded; }
QString currentBackend() const { return m_currentBackend; }
bool daemonConnected() const { return m_daemonConnected; }
Q_INVOKABLE bool agentConnected(const QString &sessionId, const QString &agentId) const;
Q_INVOKABLE QString agentState(const QString &sessionId, const QString &agentId) const;
Q_INVOKABLE QString sessionConnectionColor(const QString &sessionId) const;
Q_INVOKABLE void setCurrentBackend(const QString &backend) {
@ -74,7 +73,6 @@ public:
Q_INVOKABLE QString getConfig();
Q_INVOKABLE QStringList getAgents(const QString &sessionId);
void setAgentConnected(const QString &key, bool connected);
QString agentKey(const QString &sessionId, const QString &agentId) const;
Q_INVOKABLE void setActiveAgentId(const QString &agentId);
Q_INVOKABLE void switchAgent(const QString &sessionId, const QString &agentId);
@ -82,6 +80,7 @@ public:
Q_INVOKABLE void loadRootAgents();
Q_INVOKABLE QStringList getAvailableModels(const QString &backend) const;
Q_INVOKABLE bool createQuickSession();
Q_INVOKABLE QString createSession(const QString &name, const QString &remote);
Q_INVOKABLE QString createAgent(const QString &sessionId, const QString &cwd, const QString &backend, const QString &model, const QString &agent, const QString &remote, const QString &agentAlias);
Q_INVOKABLE bool killAgent(const QString &sessionId, const QString &agentId);
Q_INVOKABLE bool renameSession(const QString &sessionId, const QString &newName);
@ -111,10 +110,7 @@ private:
QByteArray run9p(const QStringList &args); // Fallback for streaming (uses subprocess)
void ensureRootDataLoaded();
void setDaemonConnected(bool connected);
void startAgentConnections();
void startActiveAgentStreams();
void reconcileAgentConnections();
void stopAgentConnections();
void handleEvent(const QString &eventLine);
// Native 9P client (persistent connection, no subprocess overhead)
@ -134,7 +130,6 @@ private:
QVariantMap m_rootModels;
QString m_currentBackend;
bool m_daemonConnected = false;
QHash<QString, bool> m_agentConnected; // Derived from session.connected in session/idx
NinePConnection *m_daemon = nullptr;
// Agent state cache — updated via delta events