500 lines
16 KiB
C++
500 lines
16 KiB
C++
#include "ollie9pclient.h"
|
|
|
|
#include <QDir>
|
|
#include <QSet>
|
|
#include <QStandardPaths>
|
|
#include <QDebug>
|
|
|
|
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)
|
|
{
|
|
// 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) {
|
|
if (!data.isEmpty())
|
|
emit chatReceived(QString::fromUtf8(data));
|
|
});
|
|
connect(m_chat, &StreamFsm::died, this, &Ollie9pClient::streamingDone);
|
|
|
|
// State stream — guarded, reads agent statewait. Auto-restarts until
|
|
// session/agent is deselected (guard returns false).
|
|
m_state = new StreamFsm(StreamFsm::Guarded, this);
|
|
m_state->setGuard([this]() {
|
|
return !m_activeSessionId.isEmpty() && !m_agentId.isEmpty();
|
|
});
|
|
connect(m_state, &StreamFsm::readyRead, this, [this](const QByteArray &data) {
|
|
QString state = QString::fromUtf8(data).trimmed();
|
|
if (!state.isEmpty()) {
|
|
m_activeState = state;
|
|
emit activeStateChanged();
|
|
}
|
|
});
|
|
|
|
// Event stream — looping, always restarts. Listens for server events
|
|
// (session create/kill/rename) and refreshes the session list.
|
|
m_event = new StreamFsm(StreamFsm::Looping, this);
|
|
connect(m_event, &StreamFsm::readyRead, this, [this](const QByteArray &data) {
|
|
QString raw = QString::fromUtf8(data);
|
|
if (!raw.isEmpty())
|
|
refreshSessions();
|
|
});
|
|
|
|
refreshSessions();
|
|
ensureRootDataLoaded();
|
|
m_event->start(ollie9pBin(), {"-a", serverAddr(), "read", "eventwait"});
|
|
}
|
|
|
|
Ollie9pClient::~Ollie9pClient()
|
|
{
|
|
stopStreams();
|
|
}
|
|
|
|
void Ollie9pClient::setActiveSessionId(const QString &id)
|
|
{
|
|
qDebug() << "setActiveSessionId" << id << "prev session" << m_activeSessionId << "prev agent" << m_agentId;
|
|
// Session selection must be side-effect free for per-agent streams.
|
|
// Even re-selecting the same session should clear agent selection and stop
|
|
// statewait/chat so the UI can stay in a pure "session selected" state.
|
|
stopStreams();
|
|
m_activeSessionId = id;
|
|
m_agentId.clear();
|
|
m_activeState = "idle";
|
|
|
|
if (id.isEmpty()) {
|
|
emit activeSessionIdChanged();
|
|
emit activeStateChanged();
|
|
return;
|
|
}
|
|
|
|
emit activeSessionIdChanged();
|
|
emit activeAgentIdChanged();
|
|
emit activeStateChanged();
|
|
}
|
|
|
|
void Ollie9pClient::refreshSessions()
|
|
{
|
|
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["id"] = 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] : "";
|
|
session["agent"] = parts.size() > 5 ? parts[5] : "";
|
|
m_sessions.append(session);
|
|
}
|
|
}
|
|
emit sessionsChanged();
|
|
|
|
// Auto-select first session if none active
|
|
if (m_activeSessionId.isEmpty() && !m_sessions.isEmpty()) {
|
|
setActiveSessionId(m_sessions.first().toMap()["id"].toString());
|
|
}
|
|
}
|
|
|
|
QString Ollie9pClient::readLog()
|
|
{
|
|
if (m_activeSessionId.isEmpty()) return {};
|
|
QByteArray out = run9p({"read", agentPath() + "/log"});
|
|
return QString::fromUtf8(out);
|
|
}
|
|
|
|
QString Ollie9pClient::readLogForSession(const QString &sessionId, const QString &agentId)
|
|
{
|
|
if (sessionId.isEmpty() || agentId.isEmpty()) return {};
|
|
QByteArray out = run9p({"read", "session/" + sessionId + "/agent/" + agentId + "/log"});
|
|
return QString::fromUtf8(out);
|
|
}
|
|
|
|
bool Ollie9pClient::submit(const QString &prompt)
|
|
{
|
|
if (m_activeSessionId.isEmpty() || m_agentId.isEmpty() || prompt.trimmed().isEmpty()) return false;
|
|
QProcess proc;
|
|
proc.start(ninepBin(), {"-a", serverAddr(), "write", agentPath() + "/prompt"});
|
|
proc.waitForStarted(3000);
|
|
proc.write(prompt.toUtf8());
|
|
proc.closeWriteChannel();
|
|
proc.waitForFinished(5000);
|
|
return proc.exitCode() == 0;
|
|
}
|
|
|
|
bool Ollie9pClient::interrupt()
|
|
{
|
|
if (m_activeSessionId.isEmpty() || m_agentId.isEmpty()) return false;
|
|
QProcess proc;
|
|
proc.start(ninepBin(), {"-a", serverAddr(), "write", agentPath() + "/ctl"});
|
|
proc.waitForStarted(3000);
|
|
proc.write("stop");
|
|
proc.closeWriteChannel();
|
|
proc.waitForFinished(3000);
|
|
return proc.exitCode() == 0;
|
|
}
|
|
|
|
bool Ollie9pClient::kill()
|
|
{
|
|
if (m_activeSessionId.isEmpty() || m_agentId.isEmpty()) return false;
|
|
QProcess proc;
|
|
proc.start(ninepBin(), {"-a", serverAddr(), "write", agentPath() + "/ctl"});
|
|
proc.waitForStarted(3000);
|
|
proc.write("kill");
|
|
proc.closeWriteChannel();
|
|
proc.waitForFinished(3000);
|
|
return proc.exitCode() == 0;
|
|
}
|
|
|
|
bool Ollie9pClient::killSession(const QString &sessionId)
|
|
{
|
|
if (sessionId.isEmpty()) return false;
|
|
|
|
// If killing the active session, shut down its streams first
|
|
if (sessionId == m_activeSessionId) {
|
|
stopStreams();
|
|
}
|
|
|
|
QString path = "session/" + sessionId + "/ctl";
|
|
QProcess proc;
|
|
proc.start(ninepBin(), {"-a", serverAddr(), "write", path});
|
|
proc.waitForStarted(3000);
|
|
proc.write("kill");
|
|
proc.closeWriteChannel();
|
|
proc.waitForFinished(3000);
|
|
return proc.exitCode() == 0;
|
|
}
|
|
|
|
QString Ollie9pClient::getConfig()
|
|
{
|
|
if (m_activeSessionId.isEmpty()) return {};
|
|
QByteArray out = run9p({"read", agentPath() + "/cfg"});
|
|
return QString::fromUtf8(out);
|
|
}
|
|
|
|
QStringList Ollie9pClient::getAgents(const QString &sessionId)
|
|
{
|
|
if (sessionId.isEmpty()) return {};
|
|
QByteArray out = run9p({"ls", "session/" + sessionId + "/agent"});
|
|
QString raw = QString::fromUtf8(out).trimmed();
|
|
if (raw.isEmpty()) return {};
|
|
QStringList all = raw.split('\n', Qt::SkipEmptyParts);
|
|
QStringList agents;
|
|
for (const QString &a : all) {
|
|
QString name = a.trimmed();
|
|
// "new" is the agent creation file, not an actual agent
|
|
if (name != "new")
|
|
agents.append(name);
|
|
}
|
|
return agents;
|
|
}
|
|
|
|
void Ollie9pClient::setActiveAgentId(const QString &agentId)
|
|
{
|
|
if (m_agentId == agentId) return;
|
|
stopStreams();
|
|
m_agentId = agentId;
|
|
m_activeState = "idle";
|
|
emit activeAgentIdChanged();
|
|
emit activeStateChanged();
|
|
}
|
|
|
|
void Ollie9pClient::switchAgent(const QString &sessionId, const QString &agentId)
|
|
{
|
|
qDebug() << "switchAgent" << sessionId << agentId << "prev session" << m_activeSessionId << "prev agent" << m_agentId;
|
|
if (m_activeSessionId == sessionId && m_agentId == agentId) return;
|
|
|
|
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;
|
|
|
|
emit activeSessionIdChanged();
|
|
emit activeAgentIdChanged();
|
|
emit activeStateChanged();
|
|
|
|
// switchAgent is the ONLY path that starts per-agent streams.
|
|
QString srv = serverAddr();
|
|
QString bin = ninepBin();
|
|
m_chat->start(bin, {"-a", srv, "read", agentPath() + "/chat"});
|
|
m_state->start(bin, {"-a", srv, "read", agentPath() + "/statewait"});
|
|
}
|
|
|
|
void Ollie9pClient::loadRootBackends()
|
|
{
|
|
if (m_rootBackendsLoaded) return;
|
|
QByteArray out = run9p({"read", "backends"});
|
|
QString raw = QString::fromUtf8(out).trimmed();
|
|
if (!raw.isEmpty()) {
|
|
m_availableBackends = raw.split('\n', Qt::SkipEmptyParts);
|
|
m_rootBackendsLoaded = true;
|
|
emit availableBackendsChanged();
|
|
}
|
|
emit rootBackendsLoadedChanged();
|
|
}
|
|
|
|
void Ollie9pClient::loadRootAgents()
|
|
{
|
|
if (m_rootAgentsLoaded) return;
|
|
QByteArray out = run9p({"read", "agents"});
|
|
QString raw = QString::fromUtf8(out).trimmed();
|
|
if (!raw.isEmpty()) {
|
|
m_availableAgents = raw.split('\n', Qt::SkipEmptyParts);
|
|
m_rootAgentsLoaded = true;
|
|
emit availableAgentsChanged();
|
|
} else {
|
|
m_availableAgents = {"default"};
|
|
m_rootAgentsLoaded = true;
|
|
emit availableAgentsChanged();
|
|
}
|
|
emit rootAgentsLoadedChanged();
|
|
}
|
|
|
|
QStringList Ollie9pClient::getAvailableModels(const QString &backend) const
|
|
{
|
|
if (!m_rootModelsLoaded) return {};
|
|
auto it = m_rootModels.find(backend);
|
|
if (it == m_rootModels.end()) return {};
|
|
return it.value().toStringList();
|
|
}
|
|
|
|
bool Ollie9pClient::createSession(const QString &cwd, const QString &name, const QString &backend, const QString &model, const QString &agent, const QString &remote, const QString &agentAlias)
|
|
{
|
|
if (cwd.isEmpty()) return false;
|
|
|
|
// Step 1: Create the session (name only) using rdwr (write then read back result)
|
|
QString sessName = name;
|
|
{
|
|
QProcess proc;
|
|
QStringList args = {"-a", serverAddr(), "rdwr", "session/new"};
|
|
proc.start(ninepBin(), args);
|
|
proc.waitForStarted(3000);
|
|
QString input = "name=" + sessName + "\n";
|
|
proc.write(input.toUtf8());
|
|
proc.closeWriteChannel();
|
|
proc.waitForFinished(5000);
|
|
if (proc.exitCode() != 0) {
|
|
QString err = QString::fromUtf8(proc.readAllStandardError()).trimmed();
|
|
if (err.isEmpty())
|
|
err = QString::fromUtf8(proc.readAllStandardOutput()).trimmed();
|
|
qDebug() << "createSession (step 1) failed:" << err;
|
|
return false;
|
|
}
|
|
// Read the returned session name (auto-generated if blank)
|
|
QByteArray out = proc.readAllStandardOutput().trimmed();
|
|
if (!out.isEmpty())
|
|
sessName = QString::fromUtf8(out);
|
|
}
|
|
|
|
// Step 2: Create the agent within the session using rdwr
|
|
{
|
|
QStringList agentArgs;
|
|
agentArgs << "cwd=" + cwd;
|
|
if (!backend.isEmpty()) agentArgs << "backend=" + backend;
|
|
if (!model.isEmpty()) agentArgs << "model=" + model;
|
|
if (!agent.isEmpty()) agentArgs << "agent=" + agent;
|
|
if (!remote.isEmpty()) agentArgs << "remote=" + remote;
|
|
if (!agentAlias.isEmpty()) agentArgs << "name=" + agentAlias;
|
|
|
|
QProcess proc;
|
|
QStringList args = {"-a", serverAddr(), "rdwr", "session/" + sessName + "/agent/new"};
|
|
proc.start(ninepBin(), args);
|
|
proc.waitForStarted(3000);
|
|
QString input = agentArgs.join(" ") + "\n";
|
|
proc.write(input.toUtf8());
|
|
proc.closeWriteChannel();
|
|
proc.waitForFinished(5000);
|
|
|
|
if (proc.exitCode() != 0) {
|
|
QString err = QString::fromUtf8(proc.readAllStandardError()).trimmed();
|
|
if (err.isEmpty())
|
|
err = QString::fromUtf8(proc.readAllStandardOutput()).trimmed();
|
|
qDebug() << "createSession (step 2) failed:" << err;
|
|
// Clean up the session
|
|
QProcess killProc;
|
|
killProc.start(ninepBin(), {"-a", serverAddr(), "write", "session/" + sessName + "/ctl"});
|
|
killProc.waitForStarted(3000);
|
|
killProc.write("kill");
|
|
killProc.closeWriteChannel();
|
|
killProc.waitForFinished(3000);
|
|
return false;
|
|
}
|
|
}
|
|
|
|
return true;
|
|
}
|
|
|
|
bool Ollie9pClient::renameSession(const QString &sessionId, const QString &newName)
|
|
{
|
|
if (sessionId.isEmpty() || newName.isEmpty()) return false;
|
|
|
|
QString bin = ollie9pBin();
|
|
if (bin.isEmpty()) {
|
|
qDebug() << "renameSession: ollie-9p not found, cannot rename";
|
|
return false;
|
|
}
|
|
|
|
QProcess proc;
|
|
proc.setProgram(bin);
|
|
proc.setArguments({"-a", serverAddr(), "mv", "session/" + sessionId, "session/" + newName});
|
|
proc.start();
|
|
proc.waitForFinished(5000);
|
|
|
|
if (proc.exitCode() != 0) {
|
|
qDebug() << "renameSession failed:" << QString::fromUtf8(proc.readAllStandardError()).trimmed();
|
|
return false;
|
|
}
|
|
|
|
// Update active session ID if we renamed the active session
|
|
if (m_activeSessionId == sessionId) {
|
|
m_activeSessionId = newName;
|
|
emit activeSessionIdChanged();
|
|
}
|
|
|
|
refreshSessions();
|
|
return true;
|
|
}
|
|
|
|
bool Ollie9pClient::renameAgent(const QString &sessionId, const QString &agentId, const QString &newName)
|
|
{
|
|
if (sessionId.isEmpty() || agentId.isEmpty() || newName.isEmpty()) return false;
|
|
|
|
QString bin = ollie9pBin();
|
|
if (bin.isEmpty()) {
|
|
qDebug() << "renameAgent: ollie-9p not found, cannot rename";
|
|
return false;
|
|
}
|
|
|
|
// Write name=<newName> to the agent's ctl file — Plan 9 style.
|
|
QProcess proc;
|
|
proc.setProgram(bin);
|
|
proc.setArguments({"-a", serverAddr(), "rdwr",
|
|
"session/" + sessionId + "/agent/" + agentId + "/ctl"});
|
|
proc.start();
|
|
proc.waitForStarted(3000);
|
|
QByteArray input = QString("name=" + newName + "\n").toUtf8();
|
|
proc.write(input);
|
|
proc.closeWriteChannel();
|
|
proc.waitForFinished(5000);
|
|
|
|
if (proc.exitCode() != 0) {
|
|
qDebug() << "renameAgent failed:" << QString::fromUtf8(proc.readAllStandardError()).trimmed();
|
|
return false;
|
|
}
|
|
|
|
if (m_activeSessionId == sessionId) {
|
|
m_agentId = newName;
|
|
emit activeAgentIdChanged();
|
|
}
|
|
|
|
refreshSessions();
|
|
return true;
|
|
}
|
|
|
|
// --- Root data loading (lazy) ---
|
|
|
|
void Ollie9pClient::ensureRootDataLoaded()
|
|
{
|
|
if (!m_rootBackendsLoaded) loadRootBackends();
|
|
if (!m_rootAgentsLoaded) loadRootAgents();
|
|
if (!m_rootModelsLoaded) {
|
|
QByteArray out = run9p({"read", "models"});
|
|
QString raw = QString::fromUtf8(out);
|
|
if (!raw.isEmpty()) {
|
|
for (const QString &line : raw.split('\n')) {
|
|
int tab = line.indexOf('\t');
|
|
if (tab < 0) continue;
|
|
QString be = line.left(tab).trimmed();
|
|
QString mo = line.mid(tab + 1).trimmed();
|
|
if (!be.isEmpty() && !mo.isEmpty()) {
|
|
QStringList list = m_rootModels[be].toStringList();
|
|
list.append(mo);
|
|
m_rootModels[be] = QVariant(list);
|
|
}
|
|
}
|
|
m_rootModelsLoaded = true;
|
|
emit rootModelsLoadedChanged();
|
|
}
|
|
}
|
|
refreshSessions();
|
|
}
|
|
|
|
// --- Streaming ---
|
|
|
|
void Ollie9pClient::stopStreams()
|
|
{
|
|
if (m_chat) m_chat->stop();
|
|
if (m_state) m_state->stop();
|
|
if (m_event) m_event->stop();
|
|
}
|
|
|
|
// --- Helpers ---
|
|
|
|
QString Ollie9pClient::agentPath() const
|
|
{
|
|
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;
|
|
}
|