ollie/kde/gui/plumber.cpp

286 lines
8.9 KiB
C++

/*
* SPDX-License-Identifier: GPL-3.0-or-later
*/
#include "plumber.h"
#include "ninep.h"
#include "plumbmsg.h"
#include <QByteArray>
#include <QDir>
#include <QFileInfo>
#include <QThread>
#include <unistd.h>
#include <cstdarg>
#include <cstdio>
// Canonicalise the X11 display the way plan9 getns() does: "xxx:0.0" -> "xxx:0"
// and '/' -> '_'. Only the trailing ".0" is stripped.
static QString canonicalDisplay(QString disp)
{
const int colon = disp.lastIndexOf(QLatin1Char(':'));
if (colon >= 0) {
int p = colon + 1;
while (p < disp.size() && disp[p].isDigit())
p++;
if (disp.mid(p) == QLatin1String(".0"))
disp.truncate(p);
}
disp.replace(QLatin1Char('/'), QLatin1Char('_'));
return disp;
}
QString Plumber::socketPath()
{
const QByteArray ns = qgetenv("NAMESPACE");
QString nsDir;
if (!ns.isEmpty()) {
nsDir = QString::fromLocal8Bit(ns);
} else {
const QByteArray disp = qgetenv("DISPLAY");
if (disp.isEmpty())
return QString(); // no namespace resolvable
QString user = QString::fromLocal8Bit(qgetenv("USER"));
if (user.isEmpty())
user = QString::fromLocal8Bit(qgetenv("LOGNAME"));
nsDir = QStringLiteral("/tmp/ns.%1.%2")
.arg(user, canonicalDisplay(QString::fromLocal8Bit(disp)));
}
return nsDir + QStringLiteral("/plumb");
}
// Current user name for the 9P attach
static QString currentUser()
{
QString user = QString::fromLocal8Bit(qgetenv("USER"));
if (user.isEmpty())
user = QString::fromLocal8Bit(qgetenv("LOGNAME"));
if (user.isEmpty())
user = QStringLiteral("none");
return user;
}
// ---------------------------------------------------------------------------
// PlumbReader: owns a NineP session on a named port, loops on blocking reads.
// Lives on its own QThread. Emits parsed messages via queued signals so the
// Plumber (GUI thread) can act on them.
// ---------------------------------------------------------------------------
class PlumbReader : public QThread
{
Q_OBJECT
public:
explicit PlumbReader(const QString &socketPath, const QString &portName, QObject *parent = nullptr)
: QThread(parent)
, m_socketPath(socketPath)
, m_portName(portName)
{
}
void stop()
{
m_stop.storeRelaxed(1);
}
Q_SIGNALS:
void message(const QString &data, const QString &addr, const QString &wdir);
void failed(const QString &message);
protected:
void run() override
{
static constexpr int PollMs = 200;
static constexpr int ReconnectMs = 2000;
// Log to file for debugging
FILE *logf = fopen("/tmp/ollie-plumber.log", "a");
auto log = [logf](const char *fmt, ...) {
if (!logf) return;
va_list ap;
va_start(ap, fmt);
vfprintf(logf, fmt, ap);
fprintf(logf, "\n");
fflush(logf);
va_end(ap);
};
log("PlumbReader starting for port %s on %s", qPrintable(m_portName), qPrintable(m_socketPath));
while (!m_stop.loadRelaxed()) {
NineP nine;
log("Connecting to %s...", qPrintable(m_socketPath));
if (!nine.connectAndAttach(m_socketPath, currentUser())) {
log("Connect failed: %s", qPrintable(nine.errorString()));
QThread::msleep(ReconnectMs);
continue;
}
log("Connected, walking to %s...", qPrintable(m_portName));
static constexpr uint32_t PortFid = 1;
if (!nine.walk(PortFid, m_portName)) {
log("Walk failed: %s", qPrintable(nine.errorString()));
QThread::msleep(ReconnectMs);
continue;
}
log("Opening %s for read...", qPrintable(m_portName));
if (!nine.open(PortFid, NineP::OREAD)) {
log("Open failed: %s", qPrintable(nine.errorString()));
QThread::msleep(ReconnectMs);
continue;
}
log("Port %s connected and open!", qPrintable(m_portName));
bool pending = false;
while (!m_stop.loadRelaxed()) {
if (!pending) {
if (!nine.beginRead(PortFid, 0, nine.msize())) {
log("beginRead failed: %s", qPrintable(nine.errorString()));
break; // reconnect
}
pending = true;
}
QByteArray buf;
bool timedOut = false;
const int n = nine.recvReadReply(&buf, PollMs, &timedOut);
if (timedOut)
continue;
if (n < 0) {
log("Read error: %s", qPrintable(nine.errorString()));
break; // reconnect
}
pending = false;
if (n == 0)
continue;
log("Received message: %d bytes", buf.size());
PlumbMsg m;
if (PlumbMsg::unpack(buf, &m)) {
log("Emitting message: %s", m.data.constData());
Q_EMIT message(QString::fromUtf8(m.data), m.lookup(QStringLiteral("addr")),
m.wdir);
}
}
// Connection lost, wait before reconnecting
if (!m_stop.loadRelaxed()) {
log("Disconnected, reconnecting in %dms...", ReconnectMs);
QThread::msleep(ReconnectMs);
}
}
log("PlumbReader exiting");
if (logf) fclose(logf);
}
private:
QString m_socketPath;
QString m_portName;
QAtomicInt m_stop{0};
};
// ---------------------------------------------------------------------------
// Plumber
// ---------------------------------------------------------------------------
Plumber::Plumber(QObject *parent)
: QObject(parent)
{
}
Plumber::~Plumber()
{
auto stopReader = [](PlumbReader *&reader) {
if (!reader)
return;
reader->stop();
if (reader->wait(3000)) {
delete reader;
} else {
QObject::connect(reader, &QThread::finished, reader, &QObject::deleteLater);
}
reader = nullptr;
};
stopReader(m_reader);
stopReader(m_ollieReader);
}
bool Plumber::available() const
{
const QString path = socketPath();
if (path.isEmpty())
return false;
return QFileInfo::exists(path);
}
bool Plumber::send(const QString &data, const QString &wdir)
{
return sendTo(data, QString(), wdir);
}
bool Plumber::sendTo(const QString &data, const QString &dst, const QString &wdir)
{
if (data.isEmpty())
return false;
const QString path = socketPath();
if (path.isEmpty()) {
m_error = QStringLiteral("no plumb namespace ($NAMESPACE/$DISPLAY unset)");
Q_EMIT plumbError(m_error);
return false;
}
NineP nine;
if (!nine.connectAndAttach(path, currentUser())) {
m_error = nine.errorString();
Q_EMIT plumbError(m_error);
return false;
}
static constexpr uint32_t SendFid = 1;
if (!nine.walk(SendFid, QStringLiteral("send")) || !nine.open(SendFid, NineP::OWRITE)) {
m_error = nine.errorString();
Q_EMIT plumbError(m_error);
return false;
}
PlumbMsg m;
m.src = QStringLiteral("ollie");
m.dst = dst;
m.wdir = wdir.isEmpty() ? QDir::currentPath() : wdir;
m.type = QStringLiteral("text");
m.data = data.toUtf8();
const QByteArray packed = m.pack();
const int w = nine.write(SendFid, 0, packed);
if (w != packed.size()) {
m_error = nine.errorString().isEmpty()
? QStringLiteral("short plumb write (%1/%2)").arg(w).arg(packed.size())
: nine.errorString();
Q_EMIT plumbError(m_error);
return false;
}
m_error.clear();
return true;
}
void Plumber::startReader()
{
const QString path = socketPath();
if (path.isEmpty()) {
Q_EMIT readerError(QStringLiteral("no plumb namespace ($NAMESPACE/$DISPLAY unset)"));
return;
}
// Start ollie port reader (for ollie:// URLs)
if (!m_ollieReader) {
qDebug("plumber: starting ollie port reader on %s", qPrintable(path));
m_ollieReader = new PlumbReader(path, QStringLiteral("ollie"));
connect(m_ollieReader, &PlumbReader::message, this, [this](const QString &data, const QString &, const QString &) {
qDebug("plumber: received ollie message: %s", qPrintable(data));
Q_EMIT ollieMessage(data);
}, Qt::QueuedConnection);
connect(m_ollieReader, &PlumbReader::failed, this, [](const QString &msg) {
qWarning("plumber: ollie port FAILED: %s", qPrintable(msg));
}, Qt::QueuedConnection);
m_ollieReader->start();
}
}
#include "plumber.moc"